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

AWS MSK Connect MirrorMaker2刷新及偏移量提交失败求助

解决MSK Connect(2.7.1)MirrorMaker2刷新超时与偏移量提交失败问题

问题详情

在MSK Connect 2.7.1环境下运行MirrorMaker2源连接器时,触发以下错误:

[Worker-0d8c5a576b5ef6e99] [2023-12-22 16:01:00,771] ERROR [msk-dev-conv-mm2-sourceconnector|task-1|offsets] WorkerSourceTask{id=msk-dev-conv-mm2-sourceconnector-1} Failed to flush, timed out while waiting for producer to flush outstanding 3 messages (org.apache.kafka.connect.runtime.WorkerSourceTask:509)
[Worker-0d8c5a576b5ef6e99] [2023-12-22 16:01:00,771] ERROR [msk-dev-conv-mm2-sourceconnector|task-1|offsets] WorkerSourceTask{id=msk-dev-conv-mm2-sourceconnector-1} Failed to commit offsets (org.apache.kafka.connect.runtime.SourceTaskOffsetCommitter:116)

集群配置

allow.everyone.if.no.acl.found = false
auto.create.topics.enable = true
delete.topic.enable = true
log.cleaner.delete.retention.ms = 86400000
log.cleanup.policy = compact
log.retention.hours = -1
message.max.bytes = 5242940
min.insync.replicas = 2
unclean.leader.election.enable = false

MirrorMaker2源连接器配置

connector.class=org.apache.kafka.connect.mirror.MirrorSourceConnector
errors.log.include.messages=false
replication.factor=3
source.cluster.ssl.truststore.location=${s3import:ca-central-1:msk-manual-msk-configurations/test.truststore.jks}
target.cluster.sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required awsDebugCreds=true;
sync.topic.acls.enabled=false
tasks.max=16
source.cluster.alias=
sync.topic.configs.interval.seconds=20
target.cluster.security.protocol=SASL_SSL
replication.policy.separator=
value.converter=org.apache.kafka.connect.converters.ByteArrayConverter
errors.log.enable=true
key.converter=org.apache.kafka.connect.converters.ByteArrayConverter
refresh.groups.interval.seconds=20
refresh.topics.interval.seconds=20
offset-syncs.topic.replication.factor=3
ssl.protocol=TLS
consumer.group.id=mm2-dev
target.cluster.sasl.mechanism=AWS_MSK_IAM
topics=axis.*|ausmle-test
target.cluster.sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler
producer.enable.idempotence=true
source.cluster.sasl.jaas.config=org.apache.kafka.common.security.scram.ScramLoginModule required username='123' password='123';
source.cluster.bootstrap.servers=10.0.0.1:19093,10.0.0.1:19094,10.0.0.1:19095
source.cluster.sasl.mechanism=SCRAM-SHA-512
target.cluster.alias=
target.cluster.bootstrap.servers=b-1.msk.123.c3.kafka.ca-central-1.amazonaws.com:9098,b-2.msk.123.c3.kafka.ca-central-1.amazonaws.com:9098,b-3.msk.123.c3.kafka.ca-central-1.amazonaws.com:9098
source.cluster.ssl.truststore.password=123
sync.topic.configs.enabled=true
source.cluster.security.protocol=SASL_SSL
source.cluster.ssl.endpoint.identification.algorithm=

核心原因分析

该错误本质是MirrorMaker2内置生产者向目标集群发送消息时,flush操作超时,进而触发偏移量提交失败。常见触发因素包括:目标集群负载过高、生产者超时配置过短、网络延迟大、权限不足、任务负载不均等。

解决方案

1. 调整生产者超时配置

在连接器配置中追加以下参数,延长生产者的超时窗口:

# 延长请求超时时间
producer.request.timeout.ms=60000
# 延长生产者阻塞最大时间
producer.max.block.ms=120000
# 确保该值大于request.timeout.ms + linger.ms(默认linger.ms为0)
producer.delivery.timeout.ms=130000

2. 检查目标集群健康状态

  • 查看目标MSK集群的主题ISR(同步副本)状态,确保所有副本均在ISR列表内,无严重滞后
  • 监控集群CPU、内存、磁盘IO使用率,若负载过高,优先扩容节点或调整主题分区数
  • 集群设置min.insync.replicas=2,生产者需要至少2个副本确认消息,若副本同步延迟,会直接导致确认超时

3. 优化连接器任务分配

当前tasks.max=16,需匹配源集群总分区数:

  • 若源集群总分区数远小于16,减少tasks.max至接近分区数的整数,避免资源浪费
  • 若分区数远大于16,增加tasks.max或拆分同步主题,确保单任务负载均衡

4. 验证目标集群IAM权限

目标集群使用AWS_MSK_IAM认证,需确保MSK Connect执行角色具备:

  • kafka:DescribeCluster
  • kafka:WriteData(针对同步的所有主题)
  • kafka:CreateTopic(若依赖自动创建目标主题)
  • 生产环境移除awsDebugCreds=true,避免调试模式影响性能

5. 排查网络连通性

  • 确认MSK Connect与目标集群之间的网络延迟<100ms,无丢包
  • 跨VPC场景下,检查VPC peering/中转网关的带宽是否满足同步需求

6. 调整Connect Worker参数

修改MSK Connect Worker的全局配置:

# 延长偏移量提交超时时间
offset.commit.timeout.ms=10000
# 延长任务flush超时时间
flush.timeout.ms=10000

验证步骤

修改配置后重启连接器,观察日志是否消除超时错误,同时监控目标集群的消息同步进度,确保数据正常复制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 13:55:29