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函数。
解决步骤:
安装依赖包
在Databricks notebook中先执行以下命令安装适配DBR 11.3 LTS的Azure Schema Registry Spark集成包:%pip install azure-schema-registry-spark-avro==1.0.0替换导入路径
把原来的Avro函数导入替换为Azure专用版本:from pyspark.sql.functions import col, lit # 替换原生avro函数为Azure Schema Registry集成版 from azure.spark.sql.avro.functions import from_avro, to_avro修正后的完整代码
你的原有调用逻辑是正确的,替换导入后即可正常运行: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ň
相关产品推荐
相关产品推荐

