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

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()'

问题分析与解决建议

核心原因

  1. Kafka Sink数据结构不匹配:Spark Kafka Sink要求输出DataFrame必须包含value(必填)和key(可选)字段,且类型为StringType或BinaryType。原代码输出的DataFrame无符合要求的value字段,导致Sink处理Schema时触发方法签名错误。
  2. Checkpoint目录访问限制:Docker容器间本地目录隔离,原代码使用的./tmp仅在Jupyter容器内存在,Worker节点无法访问,导致流查询无法持久化状态。
  3. 潜在依赖冲突:若Spark集群与Jupyter环境依赖版本不一致,会引发Scala方法签名不匹配(如StructType.toAttributes()返回类型定义差异)。

具体解决方案

  1. 修正输出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中。

  2. 修复Checkpoint目录
    为所有Spark节点配置共享卷(如将主机目录挂载到容器/shared路径),修改Checkpoint路径为共享卷内路径,确保所有节点可读写。

  3. 验证依赖兼容性

    • 确保spark-sql-kafka-0-10版本与Spark版本(3.4.0)、Scala版本(2.12)完全匹配。
    • 移除手动指定的kafka-clients依赖:该依赖会由spark-sql-kafka-0-10自动引入,手动指定易引发版本冲突。
  4. 测试基础流功能
    先将数据写入控制台验证逻辑是否正常,排除Kafka Sink干扰:

    df.writeStream.outputMode("append").format("console").start().awaitTermination()
    

内容的提问来源于stack exchange,提问作者Camilo Piñón

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 06:56:19