Kafka MongoDB Sink连接器无法消费事件 求助排查原因
问题现象
- 特定
EventTopic的MongoDB Sink连接器停止消费,其中分区1积压大量事件 - 其他同类型连接器(处理不同Topic)运行正常
- 为该Topic创建新的同类型连接器仍出现相同问题
- 控制台消费者可消费部分事件,但无法处理过载的分区1
- MongoDB运行正常,其他事件可正常写入
连接器配置
{ "name": "event-mongodb-sink", "config": { "connector.class": "com.mongodb.kafka.connect.MongoSinkConnector", "database":"Event", "collection":"Event", "topics":"Event", "connection.uri":"mongodb://uri", "mongo.errors.tolerance": "all", "mongo.errors.log.enable": "true", "errors.log.include.messages": "true", "writemodel.strategy":"com.mongodb.kafka.connect.sink.writemodel.strategy.ReplaceOneBusinessKeyStrategy", "document.id.strategy": "com.mongodb.kafka.connect.sink.processor.id.strategy.PartialValueStrategy", "document.id.strategy.overwrite.existing": "true", "document.id.strategy.partial.value.projection.type": "allowlist", "document.id.strategy.partial.value.projection.list": "date,clientId,orderId", "tasks.max": 2, "max.batch.size": 1000, "bulk.write.ordered": false, "errors.log.include.messages": true } }
关键错误日志
偏移提交失败错误
[2023-04-12 16:57:28,752] ERROR WorkerSinkTask{id=event-mongodb-sink-2-0} Commit of offsets threw an unexpected exception for sequence number 5: {Event-7=OffsetAndMetadata{offset=1215069, leaderEpoch=null, metadata=''}, Event-6=OffsetAndMetadata{offset=1217175, leaderEpoch=null, metadata=''}, Event-5=OffsetAndMetadata{offset=1213520, leaderEpoch=null, metadata=''}, Event-4=OffsetAndMetadata{offset=1217328, leaderEpoch=null, metadata=''}, Event-3=OffsetAndMetadata{offset=1216765, leaderEpoch=null, metadata=''}, Event-2=OffsetAndMetadata{offset=1216501, leaderEpoch=null, metadata=''}, Event-1=OffsetAndMetadata{offset=101580287, leaderEpoch=null, metadata=''}, Event-0=OffsetAndMetadata{offset=1214247, leaderEpoch=null, metadata=''}} (org.apache.kafka.connect.runtime.WorkerSinkTask)
org.apache.kafka.clients.consumer.CommitFailedException: Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.sendOffsetCommitRequest(ConsumerCoordinator.java:1163)
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.doCommitOffsetsAsync(ConsumerCoordinator.java:981)
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.commitOffsetsAsync(ConsumerCoordinator.java:948)
消费者组踢出日志
[Consumer clientId=event-mongodb-sink-2-0, groupId=event-mongodb-sink-2] Giving away all assigned partitions as lost since generation has been reset,indicating that consumer is no longer part of the group (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator)
[2023-04-12 16:57:28,753] INFO [Consumer clientId=event-mongodb-sink-2-0, groupId=event-mongodb-sink-2] Lost previously assigned partitions Event-7, Event-6, TradeEvent-5, Event-4, Event-3, Event-2, Event-1, Event-0 (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator)
[2023-04-12 16:57:28,753] INFO [Consumer clientId=event-mongodb-sink-2-0, groupId=event-mongodb-sink-2] (Re-)joining group (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator)
[2023-04-12 16:57:28,753] INFO [Consumer clientId=event-mongodb-sink-2-0, groupId=event-mongodb-sink-2] Request joining group due to: need to re-join with the given member-id (org.apache.kafka.clients.consumer.internals.ConsumerCoordinator)
原因分析
单条消息处理超时导致消费者被踢出组
调整max.batch.size=1仍无效,是因为max.poll.interval.ms(默认300000ms/5分钟)限制的是从拉取消息到提交偏移的总时长,而非单批次的消息数量。如果分区1中的单条消息处理时间超过这个阈值(比如MongoDB写入耗时过久),消费者仍会因未及时发送心跳被踢出消费组,进而触发偏移提交失败。分区1的消息存在特殊处理瓶颈
该分区积压的核心原因是消息处理成本远高于其他分区,可能的场景包括:- 消息payload异常庞大,导致序列化/反序列化耗时过长
- 使用
ReplaceOneBusinessKeyStrategy时,该分区消息对应的业务键(date,clientId,orderId)重复率极高,引发大量MongoDB写冲突或锁等待 - MongoDB针对
Event集合的业务键缺少组合索引,导致ReplaceOne操作触发全表扫描,大幅增加写入耗时 - MongoDB正在对该集合进行索引维护、分片迁移等后台操作,导致写入性能骤降
控制台消费者的局限性
控制台消费者仅做消息拉取和打印,无需执行复杂的MongoDB写入逻辑,因此不会触发超时;但当分区积压量极大时,其自身的拉取和处理能力有限,无法快速清理积压。新连接器无效的原因
问题根源在于EventTopic的分区1数据本身,而非连接器实例,因此新连接器消费同一分区时仍会遇到相同的处理超时问题。
排查与解决方向
- 抽样检查分区1的消息:使用
kafka-console-consumer.sh指定分区1,查看消息是否存在超大payload、异常字段等问题 - 优化MongoDB写入性能:
- 检查
Event集合是否存在date,clientId,orderId组合索引,若缺失则创建,提升ReplaceOne操作的效率 - 查看MongoDB慢查询日志,定位耗时过长的写入操作,排查锁等待、资源瓶颈
- 检查
- 调整Kafka消费者超时参数:
- 增大
max.poll.interval.ms(例如设置为900000即15分钟),给单条消息足够的处理时间 - 配套调整
session.timeout.ms(建议30000ms)和heartbeat.interval.ms(建议10000ms),保持心跳间隔为会话超时的1/3左右,避免误判消费者离线
- 增大
- 临时分流处理:使用Kafka工具将分区1的消息分流到临时Topic,分批处理或过滤异常消息,缓解积压
- 检查MongoDB连接池:确认连接器的MongoDB连接池配置是否足够,避免因连接耗尽导致写入等待时间过长
内容的提问来源于stack exchange,提问作者Omegaspard

