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

AWS MSK Connect配置MongoDB Kafka源连接器连接超时问题排查

AWS MSK Connect MongoDB源连接器连接超时问题排查

环境与配置情况

  • 基于公有子网搭建的AWS MSK集群,已开启公网访问,集群运行正常
  • 配置MSK Connect MongoDB Kafka源连接器,采用IAM认证,已创建对应IAM角色与权限策略
  • 连接器配置:
connector.class=com.mongodb.kafka.connect.MongoSourceConnector
tasks.max=1
connection.uri=mongodb+srv://XXXXX:XXXXXXX@XXXXXXX.qricz.mongodb.net/user?authSource=admin&replicaSet=atlas-bftscy-shard-0&readPreference=primary
database=user
topic.prefix=mongo
poll.max.batch.size=1000
poll.await.time.ms=5000
pipeline=[]
batch.size=0
publish.full.document.only=true
mongo.errors.tolerance=all
mongo.errors.deadletterqueue.topic.name=data-ingestion-dlq-beta
offset.partition.name=mongodb-source-connector-beta.1
startup.mode=copy_existing
output.format.value=json
output.format.key=json
key.converter.schemas.enable=false
value.converter.schemas.enable=false
key.converter=org.apache.kafka.connect.storage.StringConverter
value.converter=org.apache.kafka.connect.storage.StringConverter
change.stream.full.document=updateLookup

报错信息

[Worker-0aef29e0af9181e21] [2024-01-31 20:47:37,864] INFO Exception in monitor thread while connecting to server XXXXXXX.qricz.mongodb.net:27017 (org.mongodb.driver.cluster:76)
[Worker-0aef29e0af9181e21] com.mongodb.MongoSocketOpenException: Exception opening socket
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.SocketStream.open(SocketStream.java:70)
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.InternalStreamConnection.open(InternalStreamConnection.java:180)
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.DefaultServerMonitor$ServerMonitorRunnable.lookupServerDescription(DefaultServerMonitor.java:193)
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.DefaultServerMonitor$ServerMonitorRunnable.run(DefaultServerMonitor.java:157)
[Worker-0aef29e0af9181e21]  at java.base/java.lang.Thread.run(Thread.java:829)
[Worker-0aef29e0af9181e21] Caused by: java.net.SocketTimeoutException: connect timed out
[Worker-0aef29e0af9181e21]  at java.base/java.net.PlainSocketImpl.socketConnect(Native Method)
[Worker-0aef29e0af9181e21]  at java.base/java.net.AbstractPlainSocketImpl.doConnect(AbstractPlainSocketImpl.java:412)
[Worker-0aef29e0af9181e21]  at java.base/java.net.AbstractPlainSocketImpl.connectToAddress(AbstractPlainSocketImpl.java:255)
[Worker-0aef29e0af9181e21]  at java.base/java.net.AbstractPlainSocketImpl.connect(AbstractPlainSocketImpl.java:237)
[Worker-0aef29e0af9181e21]  at java.base/java.net.SocksSocketImpl.connect(SocksSocketImpl.java:392)
[Worker-0aef29e0af9181e21]  at java.base/java.net.Socket.connect(Socket.java:609)
[Worker-0aef29e0af9181e21]  at java.base/sun.security.ssl.SSLSocketImpl.connect(SSLSocketImpl.java:305)
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.SocketStreamHelper.initialize(SocketStreamHelper.java:107)
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.SocketStream.initializeSocket(SocketStream.java:79)
[Worker-0aef29e0af9181e21]  at com.mongodb.internal.connection.SocketStream.open(SocketStream.java:65)
[Worker-0aef29e0af9181e21]  ... 4 more

排查与解决方案

1. 验证网络连通性

  • 安全组出站规则:检查MSK Connect工作节点所在安全组的出站规则,确保允许向0.0.0.0/0(或MongoDB Atlas的IP段)发起TCP 27017端口的请求
  • 子网路由与NACL:确认MSK Connect所在公有子网的路由表默认路由指向互联网网关,且网络ACL(NACL)允许出站TCP 27017流量,同时允许入站的响应流量(TCP 1024-65535)
  • MongoDB Atlas白名单:登录MongoDB Atlas控制台,确认IP白名单已添加MSK Connect所在VPC的公有IP段,或临时设置为0.0.0.0/0(仅测试用,生产环境需严格限制IP)

2. 修正MongoDB连接URI

  • DNS解析验证:在MSK Connect所在VPC的EC2实例中执行nslookup XXXXXXX.qricz.mongodb.net,确认SRV记录能正常解析到MongoDB节点IP
  • URI参数补充:在connection.uri中添加超时参数,延长连接等待时间:
    mongodb+srv://XXXXX:XXXXXXX@XXXXXXX.qricz.mongodb.net/user?authSource=admin&replicaSet=atlas-bftscy-shard-0&readPreference=primary&connectTimeoutMS=30000&socketTimeoutMS=30000
    
  • 凭证校验:确认URI中的用户名、密码无特殊字符未转义问题,可通过MongoDB Compass测试连接字符串是否有效

3. 检查MSK Connect执行角色权限

  • 确认执行角色已附加AmazonMSKConnectFullAccess或自定义权限策略,包含对MSK集群、DLQ主题的访问权限

4. 连接器配置优化

  • 若初始同步数据量较大,可适当降低poll.max.batch.size,避免因MongoDB响应缓慢导致超时
  • 确认startup.mode=copy_existing符合业务需求,若仅需同步增量变更,可改为startup.mode=latest

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 23:17:17