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

多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作业同时处理相同消息,不符合消费组预期

疑问

  1. 为何同一消费组下的两个作业会重复消费Kafka数据?
  2. 如何确保同一时间仅一个作业消费消息,实现作业无缝更新(启动新查询作业、停止旧作业时无中断或重复)?

解答

一、重复消费的原因

核心问题是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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 03:43:18