Spark Structured Streaming写入Kafka报错,寻求排查方案
Spark Structured Streaming写入Kafka报错:
scala.collection.Seq org.apache.spark.sql.types.StructType.toAttributes() 开发Spark Structured Streaming应用时,从Kafka主题读取JSON数据、解析转换后写入另一Kafka主题,启动流查询时触发错误,同时尝试写入本地JSON文件也失败。
环境配置
- Master、Worker及PySpark Jupyter环境运行在独立Docker容器中
- Spark版本:3.4.0
- Scala版本:2.12
- Kafka版本:Confluent 7.6.1(
confluentinc/cp-kafka:latest) - 已验证Schema与Kafka中JSON数据结构匹配,输出主题已存在
完整代码
import logging from pyspark.sql import SparkSession from pyspark.sql.types import StructType,StructField,FloatType,IntegerType,StringType from pyspark.sql.functions import from_json,col, to_json, struct import uuid import os import time SPARK_VERSION = '3.4.0' SCALA_VERSION = '2.12' SPARK_MASTER = "spark://spark-master:7077" KAFKA_BOOTSTRAP_SERVERS = "kafka:9092" KAFKA_INPUT_TOPIC = "event-window" KAFKA_OUTPUT_TOPIC = "event-window-output" # Docker环境建议使用共享卷路径,避免节点间目录隔离问题 CHECKPOINT_LOCATION = f"/shared/tmp/event-window-output/{uuid.uuid4()}" packages = [ f'org.apache.spark:spark-sql-kafka-0-10_{SCALA_VERSION}:{SPARK_VERSION}', 'org.apache.kafka:kafka-clients:3.3.2' ] logging.basicConfig(level=logging.INFO, format='%(asctime)s:%(funcName)s:%(levelname)s:%(message)s') logger = logging.getLogger("spark_structured_streaming") try: spark_session = SparkSession.builder \ .master(SPARK_MASTER) \ .appName("SparkStructuredStreaming") \ .config("spark.jars.packages", ",".join(packages)) \ .getOrCreate() spark_session.sparkContext.setLogLevel("DEBUG") logging.info('Spark session created successfully') except Exception: logging.error("Couldn't create the spark session") try: df = spark_session \ .readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) \ .option("subscribe", KAFKA_INPUT_TOPIC) \ .option("startingOffsets", "latest") \ .load() logging.info("Initial dataframe created successfully") except Exception as e: logging.warning(f"Initial dataframe couldn't be created due to exception: {e}") df = df.selectExpr("CAST(value as STRING)", "timestamp") df.printSchema() schema = StructType([ StructField("user", StringType()), StructField("int_value", IntegerType()), ]) df = df.select( from_json(col("value"), schema).alias("info"), "timestamp" ) df = df.select("info.*", "timestamp") df.printSchema() # 修正Kafka Sink要求的数据结构 kafka_output_df = df.select( to_json(struct("user", "int_value", "timestamp")).alias("value") ) result = ( kafka_output_df.writeStream .outputMode("append") .format("kafka") .option("topic", KAFKA_OUTPUT_TOPIC) .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) .option("checkpointLocation", CHECKPOINT_LOCATION) .start() .awaitTermination() )
完整错误栈
--------------------------------------------------------------------------- StreamingQueryException Traceback (most recent call last) Cell In[1], line 72 61 df = df.select("info.*", "timestamp") 62 df.printSchema() 64 result = ( 65 df.writeStream.trigger(processingTime="10 seconds") 66 .outputMode("append") 67 .format("kafka") 68 .option("topic", KAFKA_OUTPUT_TOPIC) 69 .option("kafka.bootstrap.servers", KAFKA_BOOTSTRAP_SERVERS) 70 .option("checkpointLocation", CHECKPOINT_LOCATION) 71 .start() ---> 72 .awaitTermination() 73 ) File /usr/local/spark/python/pyspark/sql/streaming/query.py:221, in StreamingQuery.awaitTermination(self, timeout) 219 return self._jsq.awaitTermination(int(timeout * 1000)) 220 else: ---> 221 return self._jsq.awaitTermination() File /usr/local/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py:1322, in JavaMember.__call__(self, *args) 1316 command = proto.CALL_COMMAND_NAME +\ 1317 self.command_header +\ 1318 args_command +\ 1319 proto.END_COMMAND_PART 1321 answer = self.gateway_client.send_command(command) ---> 1322 return_value = get_return_value( 1323 answer, self.gateway_client, self.target_id, self.name) 1325 for temp_arg in temp_args: 1326 if hasattr(temp_arg, "_detach"): File /usr/local/spark/python/pyspark/errors/exceptions/captured.py:185, in capture_sql_exception.<locals>.deco(*a, **kw) 181 converted = convert_exception(e.java_exception) 182 if not isinstance(converted, UnknownException): 183 # Hide where the exception came from that shows a non-Pythonic 184 # JVM exception message. ---> 185 raise converted from None 186 else: 187 raise StreamingQueryException: [STREAM_FAILED] Query [id = 848c3d42-7794-4836-9818-cb6cdc3c361c, runId = 6934b4a7-874c-4deb-a93d-2f26273d5e93] terminated with exception: 'scala.collection.Seq org.apache.spark.sql.types.StructType.toAttributes()'
问题分析与解决建议
核心原因
- Kafka Sink数据结构不匹配:Spark Kafka Sink要求输出DataFrame必须包含
value(必填)和key(可选)字段,且类型为StringType或BinaryType。原代码输出的DataFrame无符合要求的value字段,导致Sink处理Schema时触发方法签名错误。 - Checkpoint目录访问限制:Docker容器间本地目录隔离,原代码使用的
./tmp仅在Jupyter容器内存在,Worker节点无法访问,导致流查询无法持久化状态。 - 潜在依赖冲突:若Spark集群与Jupyter环境依赖版本不一致,会引发Scala方法签名不匹配(如
StructType.toAttributes()返回类型定义差异)。
具体解决方案
修正输出DataFrame结构
使用to_json和struct函数将字段序列化为JSON字符串作为Kafka的value字段:from pyspark.sql.functions import to_json, struct kafka_output_df = df.select( to_json(struct("user", "int_value", "timestamp")).alias("value") )如需指定
key,可添加col("user").alias("key")到select中。修复Checkpoint目录
为所有Spark节点配置共享卷(如将主机目录挂载到容器/shared路径),修改Checkpoint路径为共享卷内路径,确保所有节点可读写。验证依赖兼容性
- 确保
spark-sql-kafka-0-10版本与Spark版本(3.4.0)、Scala版本(2.12)完全匹配。 - 移除手动指定的
kafka-clients依赖:该依赖会由spark-sql-kafka-0-10自动引入,手动指定易引发版本冲突。
- 确保
测试基础流功能
先将数据写入控制台验证逻辑是否正常,排除Kafka Sink干扰:df.writeStream.outputMode("append").format("console").start().awaitTermination()
内容的提问来源于stack exchange,提问作者Camilo Piñón
相关产品推荐
相关产品推荐

