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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 20:02:01