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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 04:22:16