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

Kafka Connect Standalone连接外部Broker断开问题排查与解决

Kafka Connect对接Snowflake时Broker断开问题的解决

环境配置

  • Ubuntu 20.04
  • Java OpenJDK 1.8.0_362
  • Kafka 3.2.1
  • snowflake-kafka-connector-1.9.1

可正常执行的操作

  1. 执行以下命令可正常消费指定Topic数据:
kafka/bin/kafka-console-consumer.sh --consumer.config kafka/config/connect-custom.properties --topic <my_topic> --bootstrap-server <server>:<port>
  1. 使用本地标准配置运行以下命令,可正常创建自定义Producer:
kafka/bin/connect-standalone.sh kafka/config/connect-standalone.properties kafka/config/connect-snowflake-kafka-connector.properties 

异常情况

使用connect-custom.properties配置运行以下命令对接生产环境Broker、导入数据到Snowflake时,出现Broker断开日志,无后续操作:

kafka/bin/connect-standalone.sh kafka/config/connect-custom.properties kafka/config/connect-snowflake-kafka-connector.properties 

日志信息

[2023-04-11 14:54:06,975] INFO [snowflakesink|task-0] [Consumer clientId=connector-consumer-snowflakesink-0, groupId=connect-snowflakesink] Node -1 disconnected. (org.apache.kafka.clients.NetworkClient:935)
[2023-04-11 14:54:06,975] INFO [snowflakesink|task-0] [Consumer clientId=connector-consumer-snowflakesink-0, groupId=connect-snowflakesink] Cancelled in-flight API_VERSIONS request with correlation id 32 due to node -1 being disconnected (elapsed time since creation: 79ms, elapsed time since send: 79ms, request timeout: 30000ms) (org.apache.kafka.clients.NetworkClient:341)
[2023-04-11 14:54:06,975] WARN [snowflakesink|task-0] [Consumer clientId=connector-consumer-snowflakesink-0, groupId=connect-snowflakesink] Bootstrap broker <server>:<port> (id: -1 rack: null) disconnected (org.apache.kafka.clients.NetworkClient:1063)

原配置文件(connect-custom.properties)

bootstrap.servers=<server>:<port>

security.protocol=SSL
ssl.truststore.type=JKS
ssl.keystore.type=JKS
ssl.truststore.location=/path/to/truststore
ssl.truststore.password=***
ssl.keystore.location=/path/to/keystore
ssl.keystore.password=***
ssl.key.password=***
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true
offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000

排查与解决

已尝试更换Java、Kafka版本,添加插件路径,问题依旧,推测为配置文件问题。最终通过在connect-custom.properties中添加Consumer级别的SSL配置解决问题,修改后的配置如下:

bootstrap.servers=<server>:<port>

security.protocol=SSL
ssl.truststore.type=JKS
ssl.keystore.type=JKS
ssl.truststore.location=/path/to/truststore
ssl.truststore.password=***
ssl.keystore.location=/path/to/keystore
ssl.keystore.password=***
ssl.key.password=***

consumer.security.protocol=SSL
consumer.ssl.truststore.type=JKS
consumer.ssl.keystore.type=JKS
consumer.ssl.truststore.location=/path/to/truststore
consumer.ssl.truststore.password=***
consumer.ssl.keystore.location=/path/to/keystore
consumer.ssl.keystore.password=***
consumer.ssl.key.password=***

key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
key.converter.schemas.enable=true
value.converter.schemas.enable=true
offset.storage.file.filename=/tmp/connect.offsets
offset.flush.interval.ms=10000

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 04:53:18