多Spark Structured Streaming作业消费同一Kafka Topic异常问题
问题描述
我有两个独立的Python脚本(job1.py和job2.py),使用Spark Structured Streaming消费Kafka主题test1的数据。两个脚本配置了相同的Kafka消费组consumer-group-1,预期同一时间仅一个作业消费数据,但同时运行两个脚本并向test1发送消息时,两个作业均处理了相同数据。
脚本代码
job1.py
from pyspark.sql import SparkSession from pyspark.sql.functions import col, length import signal import sys checkpoint_dir = "/tmp/checkpoints" kafka_bootstrap_servers = "localhost:9092" spark = SparkSession.builder \ .appName("KafkaConsumer1") \ .getOrCreate() spark.conf.set("spark.sql.streaming.stateStore.stateSchemaCheck", "true") spark.sparkContext.setLogLevel("WARN") shutdown_requested = False def shutdown_handler(signum, frame): global shutdown_requested print("Graceful shutdown initiated...") shutdown_requested = True query.stop() signal.signal(signal.SIGINT, shutdown_handler) signal.signal(signal.SIGTERM, shutdown_handler) df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", kafka_bootstrap_servers) \ .option("subscribe", "test1") \ .option("startingOffsets", "latest") \ .option("kafka.group.id", "consumer-group-1") \ .load() df = df.selectExpr("CAST(value AS STRING) as message") df = df.withColumn("char_count", length(col("message"))) query = df.writeStream \ .outputMode("append") \ .format("console") \ .option("checkpointLocation", f"{checkpoint_dir}/wordcount_dta") \ .start() try: query.awaitTermination() except Exception as e: print(f"Exception encountered: {e}")
job2.py
from pyspark.sql import SparkSession from pyspark.sql.functions import col, length import signal import sys checkpoint_dir = "/tmp/checkpoints" kafka_bootstrap_servers = "localhost:9092" spark = SparkSession.builder \ .appName("KafkaConsumer2") \ .getOrCreate() spark.conf.set("spark.sql.streaming.stateStore.stateSchemaCheck", "true") spark.sparkContext.setLogLevel("WARN") shutdown_requested = False def shutdown_handler(signum, frame): global shutdown_requested print("Graceful shutdown initiated...") shutdown_requested = True query.stop() signal.signal(signal.SIGINT, shutdown_handler) signal.signal(signal.SIGTERM, shutdown_handler) df = spark.readStream \ .format("kafka") \ .option("kafka.bootstrap.servers", kafka_bootstrap_servers) \ .option("subscribe", "test1") \ .option("startingOffsets", "latest") \ .option("kafka.group.id", "consumer-group-1") \ .load() df = df.selectExpr("CAST(value AS STRING) as message") df = df.withColumn("char_count", length(col("message"))) query = df.writeStream \ .outputMode("append") \ .format("console") \ .option("checkpointLocation", f"{checkpoint_dir}/wordcount") \ .start() try: query.awaitTermination() except Exception as e: print(f"Exception encountered: {e}")
复现步骤
- 同时启动job1.py和job2.py
- 使用
bin/kafka-console-producer.sh向Kafka主题test1发送消息
预期与实际行为
- 预期:同一消费组下仅一个Spark作业(job1.py或job2.py)消费消息
- 实际:两个Spark作业同时处理相同消息,不符合消费组预期
疑问
- 为何同一消费组下的两个作业会重复消费Kafka数据?
- 如何确保同一时间仅一个作业消费消息,实现作业无缝更新(启动新查询作业、停止旧作业时无中断或重复)?
解答
一、重复消费的原因
核心问题是Spark Structured Streaming的Kafka消费逻辑和原生Kafka消费组机制不兼容:
- 你设置的
kafka.group.id仅用于向Kafka的__consumer_offsets主题提交偏移量,Spark不会用它实现作业间的负载均衡。 - 两个作业使用了不同的检查点目录(
/tmp/checkpoints/wordcount_dta和/tmp/checkpoints/wordcount),每个作业独立维护自己的消费状态和偏移量,互相不感知对方的存在。 - Spark流查询是独立的执行单元,即使配置相同消费组,也不会实现类似原生Kafka的消费组协同。
二、解决方案
1. 确保同一时间仅一个作业消费
让多个作业共享同一检查点目录,利用Spark流查询的单点锁机制:
- 将两个作业的
checkpointLocation修改为同一个路径,比如统一设置为/tmp/checkpoints/kafka_consumer。当第二个作业启动时,若第一个作业仍在运行,它会因无法获取检查点的独占锁而启动失败,从而保证同一时间只有一个作业消费。 - 注意:如果作业的业务逻辑(如输出模式、状态Schema)不一致,共享检查点会导致报错,因此更新作业时需确保逻辑兼容。
2. 实现作业无缝更新(无中断无重复)
遵循以下步骤即可实现新旧作业的平滑切换:
步骤1:准备新作业
新作业必须使用和旧作业完全相同的检查点目录,确保它能读取旧作业的最后消费偏移量。同时保证业务逻辑与旧作业兼容。
步骤2:启动新作业
先启动新作业,它会处于等待状态(或报错,取决于Spark版本),直到旧作业释放检查点锁。
步骤3:优雅停止旧作业
向旧作业发送SIGINT或SIGTERM信号,触发优雅停止逻辑,确保旧作业将最后处理的偏移量提交到检查点。
步骤4:新作业接管消费
旧作业停止后,新作业会自动获取检查点锁,从旧作业停止的偏移量继续消费,实现无缝衔接。
修改后的示例代码(统一检查点)
以job2.py为例,修改检查点路径:
query = df.writeStream \ .outputMode("append") \ .format("console") \ .option("checkpointLocation", f"{checkpoint_dir}/kafka_consumer") # 统一检查点路径 .start()
job1.py做同样修改即可。
额外注意事项
- 不要依赖
kafka.group.id实现Spark作业的负载均衡,Spark Structured Streaming不支持同一消费组下多作业分摊消费(这是原生Kafka消费者的特性)。 - 检查点目录需使用所有作业都能访问的共享存储(如HDFS、S3),本地目录仅适用于同一机器上的作业。
- 若必须修改状态逻辑(如聚合方式),则无法共享旧检查点,此时需重新初始化,或指定
startingOffsets从最新偏移量开始,这种情况可能存在短暂的数据重复或丢失风险,建议先停止旧作业,处理完未消费数据后再启动新作业。
内容的提问来源于stack exchange,提问作者Nagaraj
相关产品推荐
相关产品推荐

