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

