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

Spark Structured Streaming按时间戳读Kafka因分区无数据报错如何解决

问题:Spark Structured Streaming按时间戳读取Kafka时因空分区报错

报错信息:

Caused by: java.lang.AssertionError: assertion failed: No offset matched from request of topic-partition topicA-0 and timestamp 1686877634000.

场景说明:Kafka主题topicA的partition-1有数据,但partition-0无数据,使用时间戳指定偏移量范围读取时,作业直接失败,即便设置了failOnDataLoss=false也无法规避错误。

使用的代码:

df_stream = (
    spark.read.format("kafka")
    .option("kafka.bootstrap.servers", kafka_servers)
    .option("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    .option("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    .option("subscribe", kafka_topic)
    .option("kafka.security.protocol", "SASL_PLAINTEXT")
    .option("kafka.group.id", kafka_group_id)
    .option("startingOffsetsByTimestamp", """{"topicA":{"0": 1686877634000, "1": 1686877634000}}""")
    .option("endingOffsetsByTimestamp", """{"topicA":{"0": 1686881234000, "1": 1686881234000}}""")
    .option("failOnDataLoss", False) 
    .options(**options)
    .load()
)

目前已尝试全量读取后按Kafka时间戳过滤,但效率偏低,以下是更优的解决方案:


方案1:给空分区指定earliest/latest替代时间戳

直接修改startingOffsetsByTimestamp和endingOffsetsByTimestamp的配置,对空分区用earliest(起始)或latest(结束)代替具体时间戳。Spark会自动跳过空分区的时间戳匹配,不会触发断言错误,且不会读取到空分区的数据。

修改后的代码示例:

df_stream = (
    spark.read.format("kafka")
    .option("kafka.bootstrap.servers", kafka_servers)
    .option("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    .option("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    .option("subscribe", kafka_topic)
    .option("kafka.security.protocol", "SASL_PLAINTEXT")
    .option("kafka.group.id", kafka_group_id)
    # 空分区0用earliest代替时间戳,有数据的分区1保留时间戳
    .option("startingOffsetsByTimestamp", """{"topicA":{"0": "earliest", "1": 1686877634000}}""")
    .option("endingOffsetsByTimestamp", """{"topicA":{"0": "latest", "1": 1686881234000}}""")
    .option("failOnDataLoss", False) 
    .options(**options)
    .load()
)

方案2:动态生成仅包含非空分区的时间戳配置

通过Kafka AdminClient提前查询主题的分区状态,筛选出有数据的分区,只给这些分区设置时间戳,空分区不加入配置。这样Spark只会处理有数据的分区,完全规避空分区的问题。

示例代码(用kafka-python实现):

from kafka import KafkaAdminClient
import json

# 初始化AdminClient
admin_client = KafkaAdminClient(
    bootstrap_servers=kafka_servers,
    security_protocol="SASL_PLAINTEXT"
    # 若需SASL认证,添加对应配置:
    # sasl_mechanism="PLAIN",
    # sasl_plain_username="your_username",
    # sasl_plain_password="your_password"
)

# 获取指定主题的分区信息
topic_info = admin_client.describe_topics([kafka_topic])[0]
valid_partitions = {}

for p in topic_info["partitions"]:
    partition_id = p["partition"]
    # 获取分区的最早/最晚偏移量,判断是否有数据
    earliest_offset, latest_offset = admin_client.get_watermark_offsets(kafka_topic, partition_id)
    if earliest_offset < latest_offset:
        valid_partitions[str(partition_id)] = 1686877634000  # 起始时间戳

# 构建起始和结束时间戳配置
starting_config = {kafka_topic: valid_partitions}
ending_config = {kafka_topic: {k:1686881234000 for k in valid_partitions.keys()}}

# 转换为JSON字符串
starting_offsets_ts = json.dumps(starting_config)
ending_offsets_ts = json.dumps(ending_config)

# 构建流读取DF
df_stream = (
    spark.read.format("kafka")
    .option("kafka.bootstrap.servers", kafka_servers)
    .option("key.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    .option("value.serializer", "org.apache.kafka.common.serialization.StringSerializer")
    .option("subscribe", kafka_topic)
    .option("kafka.security.protocol", "SASL_PLAINTEXT")
    .option("kafka.group.id", kafka_group_id)
    .option("startingOffsetsByTimestamp", starting_offsets_ts)
    .option("endingOffsetsByTimestamp", ending_offsets_ts)
    .option("failOnDataLoss", False) 
    .options(**options)
    .load()
)

方案3:升级Spark版本

部分旧版本的Spark Kafka连接器存在空分区下时间戳匹配的断言错误,升级到Spark 3.3.x及以上版本,官方可能已经修复了该问题,允许空分区跳过时间戳匹配而不触发报错。


内容的提问来源于stack exchange,提问作者Anna

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 04:42:50