You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Azure Databricks将PySpark DataFrame写入Kafka(Avro)时的报错问题

解决PySpark中to_avro()参数数量错误(Azure Schema Registry集成)

你遇到的错误完全是因为导入了错误的to_avro函数——你用的是PySpark原生的pyspark.sql.avro.functions.to_avro,这个函数确实只接受1-2个参数,而要和Azure Schema Registry集成自动注册Schema,必须使用Azure专门提供的Avro函数。

解决步骤:

  1. 安装依赖包
    在Databricks notebook中先执行以下命令安装适配DBR 11.3 LTS的Azure Schema Registry Spark集成包:

    %pip install azure-schema-registry-spark-avro==1.0.0
    
  2. 替换导入路径
    把原来的Avro函数导入替换为Azure专用版本:

    from pyspark.sql.functions import col, lit
    # 替换原生avro函数为Azure Schema Registry集成版
    from azure.spark.sql.avro.functions import from_avro, to_avro
    
  3. 修正后的完整代码
    你的原有调用逻辑是正确的,替换导入后即可正常运行:

    from pyspark.sql.functions import col, lit
    from azure.spark.sql.avro.functions import from_avro, to_avro
    
    schema_registry_addr = "https://someeventhub.servicebus.windows.net:8081"
    
    df \
      .select(
        to_avro(col("key"), lit("t-key"), schema_registry_addr).alias("key"),
        to_avro(col("value"), lit("t-value"), schema_registry_addr).alias("value")) \
      .writeStream \
      .format("kafka") \
      .option("kafka.sasl.mechanism", "PLAIN") \
      .option("kafka.security.protocol", "SASL_SSL") \
      .option("kafka.sasl.jaas.config", EH_SASL) \
      .option("kafka.bootstrap.servers", BOOTSTRAP_SERVERS) \
      .option("topic", TOPIC) \
      .option("checkpointLocation", "/FileStore/checkpoint") \
      .start()
    

额外注意:

  • 包版本要匹配DBR版本:DBR 11.3基于Scala 2.12,1.0.0版本的azure-schema-registry-spark-avro是适配的。
  • 确认Databricks集群有访问Azure Schema Registry地址的网络权限,且SASL配置中的账号密钥拥有Schema的读写权限。

内容的提问来源于stack exchange,提问作者Tomáš Sedloň

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.31 04:45:44