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

MirrorMaker2同步本地Kafka到Azure EventHubs报错排查求助

MirrorMaker2同步Azure EventHubs报错排查求助

问题背景

已搭建MirrorMaker2(MM2)实现本地Kafka集群到Azure EventHubs的单向数据同步,此前Kafka到Kafka同步无异常。初始运行正常,但数周后出现报错,疑似与开启某大型主题同步相关。

问题现象

日志频繁出现以下错误,重启服务、重新启用主题同步均无效:

2023-04-12 11:47:32,787 INFO org.apache.kafka.connect.runtime.WorkerSourceTask: WorkerSourceTask{id=MirrorSourceConnector-102} Committing offsets
2023-04-12 11:47:32,787 INFO org.apache.kafka.connect.runtime.WorkerSourceTask: WorkerSourceTask{id=MirrorSourceConnector-102} flushing 0 outstanding messages for offset commit
2023-04-12 11:47:32,789 ERROR org.apache.kafka.clients.producer.internals.Sender: [Producer clientId=producer-731] Uncaught error in kafka producer I/O thread:
java.lang.IllegalStateException: There are no in-flight requests for node 0
        at org.apache.kafka.clients.InFlightRequests.requestQueue(InFlightRequests.java:62)
        at org.apache.kafka.clients.InFlightRequests.completeNext(InFlightRequests.java:70)
        at org.apache.kafka.clients.NetworkClient.handleCompletedReceives(NetworkClient.java:838)
        at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:558)
        at org.apache.kafka.clients.producer.internals.Sender.runOnce(Sender.java:324)
        at org.apache.kafka.clients.producer.internals.Sender.run(Sender.java:239)
        at java.lang.Thread.run(Thread.java:750)

旧日志中还发现关联错误:

Failed to flush, timed out while waiting for producer to flush outstanding 
53465 messages
.....
org.apache.kafka.common.KafkaException: Producer is closed forcefully

环境配置

沿用了此前Kafka到Kafka同步的部分非必要配置:
MM2配置截图

已尝试操作

  • 重启MM2(基于Streams Replication Manager)
  • 移除所有白名单主题后重新添加,问题未解决
  • 检查Kafka Broker及Connect日志,未发现异常
  • 确认Azure侧无报错,但大型主题同步延迟严重
  • 定位到报错生产者关联同步主题及内部offset-sync-topic,该主题存在大量inactive producers;曾将其min ISR改为1测试(怀疑网络延迟),无效后改回2

排查建议方向

生产者参数调优

  • 增大request.timeout.ms:跨公网同步到Azure需更长超时时间,建议从默认30000ms调整至60000-120000ms
  • 修改max.in.flight.requests.per.connection:从默认5改为1,避免乱序请求引发的处理异常
  • 提升buffer.memory与batch.size:缓解大流量场景下的消息堆积,降低因缓冲区溢出导致的生产者强制关闭风险

Offset同步主题优化

  • 调整offset-sync-topic的分区数:若大型主题分区数过多,需增加offset-sync-topic分区数以匹配负载
  • 清理无效生产者元数据:手动清理offset-sync-topic中旧的无效生产者记录(操作前需备份数据)

网络与Azure资源排查

  • 检测链路稳定性:使用ping/traceroute工具排查本地到Azure EventHubs的网络延迟与丢包率
  • 确认EventHubs资源配额:检查吞吐量单位(TU)是否足够支撑大型主题的同步流量,是否触发了限流

MM2配置精简与针对性优化

  • 移除Kafka到Kafka同步的冗余配置:如replica.fetch.*类参数,此类参数对EventHubs无效,可能引发配置冲突
  • 为大型主题单独配置同步参数:通过topics.<大型主题名>.producer.*前缀覆盖全局生产者配置,针对性提升大流量场景下的同步性能

监控与日志增强

  • 开启生产者DEBUG日志:添加log4j.logger.org.apache.kafka.clients.producer=DEBUG,追踪请求发送与响应的完整流程
  • 增加关键指标监控:监控MM2任务的outstandingRecords指标、EventHubs的入站流量及分区负载,定位瓶颈点

内容的提问来源于stack exchange,提问作者aireducs

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 07:52:49