PySpark结构化流连接Kafka时allowAutoTopicCreation相关报错
问题解决:Spark Structured Streaming对接Kafka的UnsupportedVersionException
问题根源
报错提示MetadataRequest versions older than 4 don't support the allowAutoTopicCreation field,本质是Spark Kafka连接器版本与Spark主版本不匹配,且连接器版本过高,与Kafka 2.6.0集群的协议不兼容:
- 你使用的Spark版本是3.2.4,但提交命令中指定的
spark-sql-kafka-0-10包版本是3.4.1,跨版本的连接器会引入不兼容的Kafka客户端逻辑 - 高版本的Kafka客户端会向集群发送包含
allowAutoTopicCreation字段的元数据请求,而Kafka 2.6.0集群不支持该字段的协议版本
另外,kafka-python能正常消费是因为它的客户端版本(2.0.2)与Kafka 2.6.0兼容,和Spark使用的Java Kafka客户端不是同一套实现。
解决方案
1. 匹配Spark与Kafka连接器版本
Spark的Kafka连接器版本必须与Spark主版本完全一致,修改spark-submit命令中的包版本:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.4 pyspark_structured_streaming.py
2. 禁用自动主题创建
即使确认topic存在,客户端仍可能触发自动创建检查,添加配置禁用该功能,避免发送不兼容的请求:
修改流数据读取部分的代码,增加kafka.allow.auto.create.topics选项:
stream_data = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "host1:port1,host2:port2") \ .option("subscribe", "topic") \ .option("kafka.allow.auto.create.topics", "false") # 新增该配置 .load()
完整修改后的代码
from pyspark.sql import SparkSession from pyspark.sql.functions import from_json from pyspark.sql.types import StructType, StringType spark = SparkSession \ .builder \ .appName("PysparkTesting") \ .getOrCreate() spark.sparkContext.setLogLevel('WARN') schema = StructType().add("changeType", StringType()) stream_data = spark \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", "host1:port1,host2:port2") \ .option("subscribe", "topic") \ .option("kafka.allow.auto.create.topics", "false") \ .load() parsed_data = stream_data.selectExpr("CAST(value as STRING)") \ .select(from_json("value", schema).alias("data")) \ .select("data.changeType") query = parsed_data.writeStream \ .outputMode("append") \ .format("console") \ .start() query.awaitTermination()
验证说明
执行修改后的spark-submit命令,即可避免协议版本不兼容的报错,正常解析并输出changeType字段到控制台。
内容的提问来源于stack exchange,提问作者John F
相关产品推荐
相关产品推荐

