配置Kafka Connect连接ElasticSearch遇LOG4J2及连接拒绝错误
问题排查与解决
一、先处理Log4j2警告
ERROR StatusLogger Log4j2 could not find a logging implementation.
这个警告是因为Kafka Connect运行时缺少Log4j2核心实现依赖,解决方式:
- 检查Kafka安装目录下的
lib文件夹,确认是否存在log4j-core-*.jar和log4j-api-*.jar,若缺失需补充对应版本的jar包(版本要与Kafka自带的Log4j依赖一致)。 - 若通过Confluent Hub安装ElasticSearch连接器,确保插件目录下的依赖包含Log4j相关包,或者将Log4j核心jar放到Kafka的
lib目录中。
二、核心问题:ElasticSearch连接拒绝
ERROR Failed to create client to verify connection. (io.confluent.connect.elasticsearch.Validator:120) ElasticsearchException[java.util.concurrent.ExecutionException: java.net.ConnectException: Connection refused]; nested: ExecutionException[java.net.ConnectException: Connection refused]; nested: ConnectException[Connection refused]
按以下步骤逐一排查:
1. 验证ElasticSearch服务状态
在终端执行命令测试本地ES连接:
curl http://localhost:9200
- 若返回包含集群信息的JSON响应,说明ES服务正常。
- 若报错
Connection refused,说明ES未启动或端口配置错误:- 启动ElasticSearch服务,检查启动日志无报错。
- 查看ES配置文件
config/elasticsearch.yml中的http.port是否为9200,若修改过端口,同步调整elasticsearch.properties中的connection.url。
2. 修正connection.url格式
你的elasticsearch.properties中配置的是:
connection.url=localhost:9200
部分版本的ElasticSearch连接器要求完整URL格式(包含协议),改为:
connection.url=http://localhost:9200
若ES启用了HTTPS,需改为https://localhost:9200并补充证书相关配置(如有需要)。
3. 检查ElasticSearch网络绑定配置
查看elasticsearch.yml中的network.host配置:
- 若设置为
0.0.0.0或非127.0.0.1的特定IP,确认Kafka Connect所在机器能访问该IP。 - 本地测试场景建议设置为
127.0.0.1或localhost,确保本地连接正常。
4. 确认版本兼容性
- 你的Kafka版本是3.2.3,需确保Confluent ElasticSearch连接器版本与Kafka版本兼容(例如Confluent Platform 7.2.x对应Kafka 3.2.x)。
- 同时确认ElasticSearch版本与连接器版本匹配:ES 7.x需连接器5.x+,ES 8.x需连接器11.x+。
5. 优化其他配置细节
type.name:ElasticSearch 7.x及以上版本已废弃文档类型,默认使用_doc,建议将type.name改为_doc,避免后续写入数据出错。- 若要使用注释掉的Bonsai远程ES服务,取消对应配置的注释并注释本地
connection.url,同时确保网络能访问该远程地址。
三、重启测试
修改配置后,重新启动Kafka Connect:
connect-standalone.sh config/connect-standalone.properties config/elasticsearch.properties
若仍有问题,查看Kafka Connect日志文件(默认在logs/connect.log)获取更详细的错误信息。
内容的提问来源于stack exchange,提问作者punky
相关产品推荐
相关产品推荐

