Apache Spark 3.5批处理模式下Kafka偏移量重复读取问题
问题描述
基于Spark官方Kafka集成指南编写了Kafka数据源的批处理查询,希望每日定时运行以处理上次运行后新增的记录,但测试时发现每次批处理都会读取所有历史记录,而非仅新增数据。代码示例如下:
builder = (pyspark.sql.SparkSession.builder.appName("MyApp") .master("local[*]") .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") .config("spark.sql.execution.arrow.pyspark.enabled", "true") .config("spark.hadoop.fs.s3a.access.key", s3a_access_key) .config("spark.hadoop.fs.s3a.secret.key", s3a_secret_key) .config("spark.hadoop.fs.s3a.endpoint", s3a_host_port) .config("spark.hadoop.fs.s3a.path.style.access", "true") .config("spark.hadoop.fs.s3a.connection.ssl.enabled", "false") .config("spark.databricks.delta.retentionDurationCheck.enabled", "false") .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") .config("spark.driver.extraJavaOptions", "-Dlog4j.configuration=file:/data/custom-log4j.properties") ) my_packages = [ # "io.delta:delta-spark_2.12:3.0.0", -> no need, since configure_spark_with_delta_pip below adds it "org.apache.hadoop:hadoop-aws:3.3.4", "org.apache.hadoop:hadoop-client-runtime:3.3.4", "org.apache.hadoop:hadoop-client-api:3.3.4", "io.delta:delta-contribs_2.12:3.0.0", "io.delta:delta-hive_2.12:3.0.0", "com.amazonaws:aws-java-sdk-bundle:1.12.603", "org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0", ] # Create a Spark instance with the builder # As a result, you now can read and write Delta tables spark = configure_spark_with_delta_pip(builder, extra_packages=my_packages).getOrCreate() kdf = (spark .read .format("kafka") .option("kafka.bootstrap.servers", kafka_bootstrap_servers) .option("kafka.security.protocol", kafka_security_protocol) .option("kafka.sasl.mechanism", "SCRAM-SHA-256") .option("kafka.sasl.jaas.config", f"org.apache.kafka.common.security.scram.ScramLoginModule required username=\"{kafka_username}\" password=\"{kafka_password}\";") .option("includeHeaders", "true") .option("subscribe", "filebeat") .option("checkpointLocation", "s3a://checkpointlocation/") .load()) kdf = kdf.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)", "headers", "CAST(topic AS STRING)", "CAST(partition AS STRING)", "CAST(offset AS STRING)") out = kdf... (out.select(["message", "partition", "offset"]) .show( truncate=False, n=MAX_JAVA_INT )) spark.stop()
输出显示每次运行都会处理相同的偏移量,请问需要修改哪些配置或代码才能实现每次仅处理新增的Kafka记录?
解决方案
Spark批处理API spark.read.format("kafka")不会自动维护偏移量,你设置的checkpointLocation仅对结构化流(readStream)生效,因此每次批处理都会从头读取全量数据。要实现增量读取,需手动管理每个Kafka分区的偏移量,具体步骤如下:
1. 存储已处理的偏移量
利用你已集成的Delta表来持久化每个Kafka分区的最后处理偏移量(也可使用数据库等外部存储)。首次运行时初始化偏移量表:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, LongType # 定义偏移量表结构 offset_schema = StructType([ StructField("topic", StringType(), nullable=False), StructField("partition", IntegerType(), nullable=False), StructField("offset", LongType(), nullable=False) ]) # 初始化:-2代表从最早偏移量开始(首次运行) initial_offsets = spark.createDataFrame( [("filebeat", 0, -2), ("filebeat", 1, -2)], # 替换为你的实际分区列表 schema=offset_schema ) initial_offsets.write.format("delta").mode("overwrite").save("s3a://your-bucket/offset-table")
2. 读取上次偏移量,构建增量读取条件
从偏移量表中获取已处理的最大偏移量,设置Kafka的startingOffsets参数以读取新增数据:
import json from pyspark.sql import functions as F # 读取已处理的偏移量 offset_df = spark.read.format("delta").load("s3a://your-bucket/offset-table") # 转换为Kafka startingOffsets所需的JSON格式 starting_offsets = {} for row in offset_df.collect(): topic = row["topic"] partition = row["partition"] # 从已处理偏移量的下一个位置开始读取 start_offset = row["offset"] + 1 if row["offset"] != -2 else "earliest" if topic not in starting_offsets: starting_offsets[topic] = {} starting_offsets[topic][str(partition)] = start_offset starting_offsets_json = json.dumps(starting_offsets) # 增量读取Kafka数据 kdf = (spark .read .format("kafka") .option("kafka.bootstrap.servers", kafka_bootstrap_servers) .option("kafka.security.protocol", kafka_security_protocol) .option("kafka.sasl.mechanism", "SCRAM-SHA-256") .option("kafka.sasl.jaas.config", f"org.apache.kafka.common.security.scram.ScramLoginModule required username=\"{kafka_username}\" password=\"{kafka_password}\";") .option("includeHeaders", "true") .option("subscribe", "filebeat") .option("startingOffsets", starting_offsets_json) .option("endingOffsets", "latest") # 读取到当前最新偏移量 .load()) # 转换字段类型(注意offset转为LongType以便后续计算) kdf = kdf.selectExpr( "CAST(key AS STRING)", "CAST(value AS STRING)", "headers", "CAST(topic AS STRING)", "CAST(partition AS INTEGER)", "CAST(offset AS BIGINT)" )
3. 更新偏移量表
处理完数据后,计算本次处理的每个分区最大偏移量,原子更新到偏移量表:
# 计算本次处理的每个分区最大偏移量 current_max_offsets = kdf.groupBy("topic", "partition").agg(F.max("offset").alias("offset")) # 原子覆盖旧的偏移量数据 current_max_offsets.write.format("delta").mode("overwrite").save("s3a://your-bucket/offset-table")
关键注意事项
- 偏移量更新必须是原子操作:Delta表的写入是原子性的,能避免重复处理或漏处理数据。
- 首次运行初始化:如果偏移量表不存在,需先初始化,或默认从
"earliest"或"latest"开始读取。 - 区分批处理与流处理配置:
checkpointLocation仅适用于结构化流,批处理中无需设置。
内容的提问来源于stack exchange,提问作者Ant Kutschera
相关产品推荐
相关产品推荐

