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

修改AWS MSK集群Advertised端口后Snowflake Connector无法连接

AWS MSK自定义端口后Snowflake连接器连接失败问题

问题描述

我拥有一个AWS MSK集群,采用Goldman Sachs单NLB方案部署,集群启用IAM认证。当Advertised端口使用默认的9098时,Snowflake连接器运行正常;但将两个Broker的端口分别修改为9001/9002后,连接器无法连接到Broker并崩溃。同一VPC内的EC2实例连接集群时仍能正常运行,无异常。

修改端口所用命令

./kafka-configs \
--bootstrap-server $B1:9098 \
--entity-type brokers \
--entity-name 1 \
--alter \
--command-config client_iam.properties \
--add-config advertised.listeners=[CLIENT_IAM://$B1:9001,REPLICATION://b-1-internal.$KF_DOMAIN:9093,REPLICATION_SECURE://b-1-internal.$KF_DOMAIN:9095]

连接器日志中的错误信息

[Worker-03e3eee9d7f02cfe6] [2022-09-28 16:52:43,013] WARN [AdminClient clientId=adminclient-8] Connection to node 2 (b-2.XXX.XXX.c16.kafka.us-east-1.amazonaws.com/INTERNAL_IP) could not be established. Broker may not be available. (org.apache.kafka.clients.NetworkClient:782)
[Worker-03e3eee9d7f02cfe6] [2022-09-28 16:52:44,620] INFO [Worker clientId=connect-1, groupId=__amazon_msk_connect_cluster_XXX-json_631e64b8-7451-4c20-aabf-6a5fb80081a6-4] Reading to end of config log; current config state offset: 12 (org.apache.kafka.connect.runtime.distributed.DistributedHerder:1198)
[Worker-03e3eee9d7f02cfe6] [2022-09-28 16:52:47,620] WARN [Worker clientId=connect-1, groupId=__amazon_msk_connect_cluster_XXX-json_631e64b8-7451-4c20-aabf-6a5fb80081a6-4] Didn't reach end of config log quickly enough (org.apache.kafka.connect.runtime.distributed.DistributedHerder:1209)
[Worker-03e3eee9d7f02cfe6] java.util.concurrent.TimeoutException: Timed out waiting for future


[Worker-0c13b739d8fb9e210]  rest.advertised.host.name = INTERNAL_IP
[Worker-0c13b739d8fb9e210]  rest.advertised.listener = HTTP
[Worker-0c13b739d8fb9e210]  rest.advertised.port = null

Snowflake连接器配置

connector.class=com.snowflake.kafka.connector.SnowflakeSinkConnector
snowflake.topic2table.map=boosterflow:boosterflow
tasks.max=8
topics=boosterflow
snowflake.topic.sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler
buffer.flush.time=60
snowflake.url.name=XXXX.snowflakecomputing.com:443
snowflake.database.name=XXXXX
snowflake.schema.name=XXXXXX
buffer.count.records=10000
snowflake.user.name=XXXXX
snowflake.private.key="XXXXXXX"
snowflake.topic.sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required;
snowflake.topic.security.protocol=SASL_SSL
value.converter=com.snowflake.kafka.connector.records.SnowflakeJsonConverter
snowflake.topic.sasl.mechanism=AWS_MSK_IAM
buffer.size.bytes=5000000

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 08:15:44