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:DescribeClusterkafka: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
相关产品推荐
相关产品推荐

