Kafka实时数据无法推送至S3 Bucket:序列化错误求助
问题描述
我是该领域新手,正在尝试将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()
额外检查项
- AWS权限验证:确保配置的AWS密钥拥有S3读写权限(
s3:PutObject、s3:GetObject、s3:ListBucket等); - 依赖兼容性:Spark 3.5.0与hadoop-aws 3.3.1版本兼容,确认所有依赖包版本匹配;
- Kafka连通性:验证Spark集群能正常访问Kafka Broker(
broker:29092)。
内容的提问来源于stack exchange,提问作者Subhamoy Paul

