NiFi同步Kafka Topic消息不全求助:源Topic 2200万仅同步1200万
Kafka主题压缩特性导致源Topic实际保留的消息数少于2200万
由于源Topic x1的cleanup.policy设置为compact,Kafka会定期对主题进行压缩,仅保留每个key对应的最新版本消息。如果x1中存在大量重复key的消息,旧版本的消息会被清理,实际保留的唯一key的最新消息总数可能就是1200万条,这与同步到x2的数量一致。此时你看到的"2200万条消息"可能是累计生产的总条数,而非当前主题实际存储的消息数。消费者组已有提交的偏移量,导致未从最早位置开始消费
ConsumeKafka_2_6的Offset Reset设置为earliest仅在消费者组没有已提交的偏移量时生效。如果该NiFi流之前已经运行过,并且消费者组已经提交过x1的偏移量,那么本次消费会从上次提交的偏移量位置开始,而非从头消费。这会导致仅同步后续的消息,而非全部2200万条。PublishKafka_2_6组件存在消息发送失败,且未被监控到
部分消息可能在发送到x2时失败(例如网络波动、broker临时不可用、消息格式问题等),而NiFi的PublishKafka组件可能将失败的消息路由到了failure关系,但你未对该关系进行处理或监控。需要检查PublishKafka组件的连接,确认是否有消息流入失败队列,以及失败原因。消费者未订阅源Topic的所有分区
如果x1包含多个分区,但ConsumeKafka_2_6的配置中未正确订阅所有分区(例如手动指定了部分分区,或者消费者组的分区分配策略导致部分分区未被分配),则会导致仅消费部分分区的消息,最终同步到x2的数量不足。Kafka日志清理导致部分早期消息已被删除
即使cleanup.policy是compact,Kafka仍会根据retention.ms或retention.bytes配置清理超出保留期限或大小的日志段。如果x1中部分早期消息的日志段已被清理(即使未被压缩),消费者也无法读取这些消息,导致同步数量减少。
内容的提问来源于stack exchange,提问作者EdiM

