ksqlDB重分区流对应Topic首偏移量接近末偏移量原因咨询
ksqlDB重分区Topic首偏移量接近尾偏移量的核心成因
该现象属于ksqlDB默认配置下的正常表现,不存在数据异常丢失,具体触发逻辑如下:
- 重分区持久查询默认从最新位点启动消费。你通过
CREATE STREAM AS SELECT(简称CSAS)语句创建重分区流时,若未显式指定消费起始位点,ksqlDB默认会从语句提交时刻源Topicdatachanges.user的最新偏移量开始处理数据,不会回溯读取查询启动前源Topic内已存在的历史存量数据。重分区Topic仅会写入查询启动后处理的增量数据,不会包含全量历史变更。 - ksqlDB自动创建的内部重分区Topic默认启用短周期日志清理策略。这类由ksqlDB自主生命周期管理的重分区Topic,不属于面向业务长期留存数据的输出Topic,默认配置的消息留存周期极短(主流版本默认留存时长为6小时),超过留存窗口的旧消息段会被Kafka Broker主动删除,分区的日志起始偏移量(即你观测到的first offset)会随着旧消息清理持续向前推进。
- Kafka的first offset统计的是当前分区实际留存的最早消息偏移量,而非Topic创建以来写入的第一条消息的偏移量。当旧消息被持续清理、新消息写入速率平稳时,就会出现当前留存最早消息的偏移量与最新写入消息的偏移量(last offset)数值接近的情况。
调整方案
如果需要重分区流覆盖源Topic全量历史数据,可在创建CSAS语句时显式指定从最早位点启动消费,示例代码如下:
CREATE OR REPLACE STREAM REPARTITIONED_DATACHANGES_USER WITH (KAFKA_TOPIC='repartitioned.datachanges.user', VALUE_FORMAT='AVRO', KEY_FORMAT='AVRO') AS SELECT u.ROWKEY->ID AS USER_ID, u.BEFORE, U.AFTER FROM DATACHANGES_USER u PARTITION by u.ROWKEY->ID EMIT CHANGES STARTING FROM 'earliest';
如果需要延长重分区Topic的数据留存时长,可在流定义的WITH子句中显式指定留存参数覆盖默认值,例如添加RETENTION_MS = 604800000即可将留存周期设置为7天。
内容的提问来源于stack exchange,提问作者Ciro di Marzo
相关产品推荐
相关产品推荐

