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
相关产品推荐
相关产品推荐

