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

Kafka实时数据无法推送至S3 Bucket:序列化错误求助

Spark流处理写入S3报错:Task not serializable

问题描述

我是该领域新手,正在尝试将Kafka实时流数据推送到AWS S3 Bucket。Kafka Broker已生成数据,消费者代码能在S3创建checkpoint文件夹,但无法推送实时数据,报错信息如下:

org.apache.spark.serializer.SerializationDebugger$.improveException(SerializationDebugger.scala:41)
at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:49)
at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:115)
at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:441)
... 46 more
Traceback (most recent call last):
File "/opt/bitnami/spark/jobs/spark-realtime.py", line 131, in
main()
File "/opt/bitnami/spark/jobs/spark-realtime.py", line 126, in main
query5.awaitTermination()
File "/opt/bitnami/spark/python/lib/pyspark.zip/pyspark/sql/streaming/query.py", line 221, in awaitTermination
File "/opt/bitnami/spark/python/lib/py4j-0.10.9.7-src.zip/py4j/java_gateway.py", line 1322, in call
File "/opt/bitnami/spark/python/lib/pyspark.zip/pyspark/errors/exceptions/captured.py", line 185, in deco
pyspark.errors.exceptions.captured.StreamingQueryException: [STREAM_FAILED] Query [id = a1964eae-4521-408a-a0c4-43546366574d, runId = 60a6c320-4a1c-42de-8663-07b8ab5bf91c] terminated with exception: Task not serializable

原消费者代码

from pyspark.sql import SparkSession
from config import configuration
from pyspark.sql.types import StructType,StructField,StringType,DoubleType,IntegerType,TimestampType
from pyspark.sql.functions import from_json,col


def main():

    from pyspark.conf import SparkConf


    spark=SparkSession.builder.appName("RealtimeStreaming")\
    .config("spark.jars.packages", 
                "org.apache.spark:spark-sql-kafka-0-10_2.13:3.5.0,"
                "org.apache.hadoop:hadoop-aws:3.3.1,"
                "com.amazonaws:aws-java-sdk:1.11.469" )\
    .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")\
    .config("spark.hadoop.fs.s3a.access.key", configuration.get("AWS_ACCESS_KEY"))\
    .config("spark.hadoop.fs.s3a.secret.key", configuration.get("AWS_SECRET_KEY"))\
    .config('spark.hadoop.fs.s3a.aws.credentials.provider', 'org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider')\
    .getOrCreate()

    #print(spark.sparkContext.getConf().getAll())

    spark.sparkContext.setLogLevel('WARN')

    vehicleSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("location", StringType(), True),
                StructField("speed", DoubleType(), True),
                StructField("direction", StringType(), True),
                StructField("make", StringType(), True),
                StructField("model", StringType(), True),
                StructField("year", IntegerType(), True),
                StructField("fuelType", StringType(), True)                             
                ])

    gpsSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("speed", DoubleType(), True),
                StructField("direction", StringType(), True),
                StructField("vehicle_type", StringType(), True)
                ])

    trafficSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("cameraId", StringType(), True),
                StructField("location", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("snapshot", StringType(), True)
                ])

    weatherSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("location", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("temperature", DoubleType(), True),
                StructField("weatherCondition", StringType(), True),
                StructField("precipitation", DoubleType(), True),
                StructField("windspeed", DoubleType(), True),
                StructField("humidity", IntegerType(), True),
                StructField("airQualityIndex", DoubleType(), True)
                ])

    emergencySchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("incidentId", StringType(), True),
                StructField("type", StringType(), True),
                StructField("location", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("status", StringType(), True),
                StructField("description", StringType(), True)
                ])
    

    def read_kafka_topic(topic, schema):
            try:
                return (spark.readStream
                    .format('kafka')
                    .option('kafka.bootstrap.servers', 'broker:29092')
                    .option('subscribe', topic)
                    .option('startingOffsets','earliest')
                    .load()
                    .selectExpr('CAST(value AS STRING)')
                    .select(from_json(col('value'), schema).alias('data'))
                    .select('data.*')
                    .withWatermark('timestamp', '2 minutes')
                    )
            except Exception as e:
                print(f"Error: {e}")
                raise

    def streamWriter(DataFrame, checkpointFolder, output):
        return(DataFrame.writeStream
               .format('parquet')
               .option('checkpointlocation', checkpointFolder)
               .option('path', output)
               .outputMode('append')
               .start()
               )

    vehicleDF=read_kafka_topic('vehicle_data', vehicleSchema).alias('vehicle')
    gpsDF=read_kafka_topic('gps_data', gpsSchema).alias('gps')
    trafficDF=read_kafka_topic('traffic_data', trafficSchema).alias('traffic')
    weatherDF=read_kafka_topic('weather_data', weatherSchema).alias('weather')
    emergencyDF=read_kafka_topic('emergency_data', emergencySchema).alias('emergency')

    query1=streamWriter(vehicleDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/vehicle_data',
                 output='s3a://realtime-datastream-bucket/data/vehicle_data')
    query2=streamWriter(gpsDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/gps_data',
                 output='s3a://realtime-datastream-bucket/data/gps_data')
    query3=streamWriter(trafficDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/traffic_data',
                 output='s3a://realtime-datastream-bucket/data/traffic_data')
    query4=streamWriter(weatherDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/weather_data',
                 output='s3a://realtime-datastream-bucket/data/weather_data')
    query5=streamWriter(emergencyDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/emergency_data',
                 output='s3a://realtime-datastream-bucket/data/emergency_data')
    
    query5.awaitTermination()



if __name__=="__main__":
    main()

问题分析

报错核心是Task not serializable,根源在于:
main函数内部定义的嵌套函数read_kafka_topic和streamWriter会隐式捕获main函数中的spark(SparkSession实例),而SparkSession是不可序列化的对象。当Spark将任务分发到Executor节点时,需要序列化这些函数及其依赖,因此触发序列化失败。

另外原代码中存在一个拼写错误:option('checkpointlocation', checkpointFolder)应为checkpointLocation(驼峰式命名),这也可能导致checkpoint逻辑异常。

解决方案

将嵌套函数移到main函数外部,避免捕获不可序列化的SparkSession上下文,同时修正参数拼写错误:

from pyspark.sql import SparkSession
from config import configuration
from pyspark.sql.types import StructType,StructField,StringType,DoubleType,IntegerType,TimestampType
from pyspark.sql.functions import from_json,col

# 移到main外部,显式传入spark参数
def read_kafka_topic(spark, topic, schema):
    try:
        return (spark.readStream
            .format('kafka')
            .option('kafka.bootstrap.servers', 'broker:29092')
            .option('subscribe', topic)
            .option('startingOffsets','earliest')
            .load()
            .selectExpr('CAST(value AS STRING)')
            .select(from_json(col('value'), schema).alias('data'))
            .select('data.*')
            .withWatermark('timestamp', '2 minutes')
            )
    except Exception as e:
        print(f"Error: {e}")
        raise

def streamWriter(DataFrame, checkpointFolder, output):
    return(DataFrame.writeStream
           .format('parquet')
           .option('checkpointLocation', checkpointFolder)  # 修正拼写错误
           .option('path', output)
           .outputMode('append')
           .start()
           )

def main():
    from pyspark.conf import SparkConf

    spark=SparkSession.builder.appName("RealtimeStreaming")\
    .config("spark.jars.packages", 
                "org.apache.spark:spark-sql-kafka-0-10_2.13:3.5.0,"
                "org.apache.hadoop:hadoop-aws:3.3.1,"
                "com.amazonaws:aws-java-sdk:1.11.469" )\
    .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")\
    .config("spark.hadoop.fs.s3a.access.key", configuration.get("AWS_ACCESS_KEY"))\
    .config("spark.hadoop.fs.s3a.secret.key", configuration.get("AWS_SECRET_KEY"))\
    .config('spark.hadoop.fs.s3a.aws.credentials.provider', 'org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider')\
    .getOrCreate()

    spark.sparkContext.setLogLevel('WARN')

    # 定义各类Schema(与原代码一致)
    vehicleSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("location", StringType(), True),
                StructField("speed", DoubleType(), True),
                StructField("direction", StringType(), True),
                StructField("make", StringType(), True),
                StructField("model", StringType(), True),
                StructField("year", IntegerType(), True),
                StructField("fuelType", StringType(), True)                             
                ])

    gpsSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("speed", DoubleType(), True),
                StructField("direction", StringType(), True),
                StructField("vehicle_type", StringType(), True)
                ])

    trafficSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("cameraId", StringType(), True),
                StructField("location", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("snapshot", StringType(), True)
                ])

    weatherSchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("location", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("temperature", DoubleType(), True),
                StructField("weatherCondition", StringType(), True),
                StructField("precipitation", DoubleType(), True),
                StructField("windspeed", DoubleType(), True),
                StructField("humidity", IntegerType(), True),
                StructField("airQualityIndex", DoubleType(), True)
                ])

    emergencySchema= StructType([
                StructField("id", StringType(), True),
                StructField("deviceId", StringType(), True),
                StructField("incidentId", StringType(), True),
                StructField("type", StringType(), True),
                StructField("location", StringType(), True),
                StructField("timestamp", TimestampType(), True),
                StructField("status", StringType(), True),
                StructField("description", StringType(), True)
                ])
    
    # 调用函数时传入spark实例
    vehicleDF=read_kafka_topic(spark, 'vehicle_data', vehicleSchema).alias('vehicle')
    gpsDF=read_kafka_topic(spark, 'gps_data', gpsSchema).alias('gps')
    trafficDF=read_kafka_topic(spark, 'traffic_data', trafficSchema).alias('traffic')
    weatherDF=read_kafka_topic(spark, 'weather_data', weatherSchema).alias('weather')
    emergencyDF=read_kafka_topic(spark, 'emergency_data', emergencySchema).alias('emergency')

    query1=streamWriter(vehicleDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/vehicle_data',
                 output='s3a://realtime-datastream-bucket/data/vehicle_data')
    query2=streamWriter(gpsDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/gps_data',
                 output='s3a://realtime-datastream-bucket/data/gps_data')
    query3=streamWriter(trafficDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/traffic_data',
                 output='s3a://realtime-datastream-bucket/data/traffic_data')
    query4=streamWriter(weatherDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/weather_data',
                 output='s3a://realtime-datastream-bucket/data/weather_data')
    query5=streamWriter(emergencyDF, checkpointFolder='s3a://realtime-datastream-bucket/checkpoints/emergency_data',
                 output='s3a://realtime-datastream-bucket/data/emergency_data')
    
    query5.awaitTermination()

if __name__=="__main__":
    main()

额外检查项

  1. AWS权限验证:确保配置的AWS密钥拥有S3读写权限(s3:PutObject、s3:GetObject、s3:ListBucket等);
  2. 依赖兼容性:Spark 3.5.0与hadoop-aws 3.3.1版本兼容,确认所有依赖包版本匹配;
  3. Kafka连通性:验证Spark集群能正常访问Kafka Broker(broker:29092)。

内容的提问来源于stack exchange,提问作者Subhamoy Paul

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 03:07:02