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
可正常执行的操作
- 执行以下命令可正常消费指定Topic数据:
kafka/bin/kafka-console-consumer.sh --consumer.config kafka/config/connect-custom.properties --topic <my_topic> --bootstrap-server <server>:<port>
- 使用本地标准配置运行以下命令,可正常创建自定义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
相关产品推荐
相关产品推荐

