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

Kafka MongoDB Sink连接器无法消费事件 求助排查原因

MongoDB Kafka Sink连接器消费停滞问题分析

问题现象

  • 特定Event Topic的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)


原因分析

  1. 单条消息处理超时导致消费者被踢出组
    调整max.batch.size=1仍无效,是因为max.poll.interval.ms(默认300000ms/5分钟)限制的是从拉取消息到提交偏移的总时长,而非单批次的消息数量。如果分区1中的单条消息处理时间超过这个阈值(比如MongoDB写入耗时过久),消费者仍会因未及时发送心跳被踢出消费组,进而触发偏移提交失败。

  2. 分区1的消息存在特殊处理瓶颈
    该分区积压的核心原因是消息处理成本远高于其他分区,可能的场景包括:

    • 消息payload异常庞大,导致序列化/反序列化耗时过长
    • 使用ReplaceOneBusinessKeyStrategy时,该分区消息对应的业务键(date,clientId,orderId)重复率极高,引发大量MongoDB写冲突或锁等待
    • MongoDB针对Event集合的业务键缺少组合索引,导致ReplaceOne操作触发全表扫描,大幅增加写入耗时
    • MongoDB正在对该集合进行索引维护、分片迁移等后台操作,导致写入性能骤降
  3. 控制台消费者的局限性
    控制台消费者仅做消息拉取和打印,无需执行复杂的MongoDB写入逻辑,因此不会触发超时;但当分区积压量极大时,其自身的拉取和处理能力有限,无法快速清理积压。

  4. 新连接器无效的原因
    问题根源在于Event Topic的分区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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 09:52:55