Camel Elasticsearch REST组件K8s环境报Connection refused错误
问题环境
- Kubernetes集群,基于Camel-K运行Camel 3.17版本
- 对接Elastic 7.7实例,使用官方Elasticsearch REST Component
- 所有组件通过Kubernetes Service暴露服务,已验证集群内其他Pod访问
elasticsearch-rlam-service:9200可正常连通,无全局网络故障 - 运行集成路由时触发
Connection refused(连接被拒绝)错误,无论手动实例化注册ElasticsearchComponent组件,还是仅在端点URI配置hostAddresses参数,问题均复现
问题复现代码
import org.apache.camel.builder.RouteBuilder; import org.apache.camel.impl.DefaultCamelContext; import org.apache.camel.component.elasticsearch.ElasticsearchComponent; public class Routes extends RouteBuilder { @Override public void configure() throws Exception { // 手动注册组件的逻辑即使删除,问题仍存在 ElasticsearchComponent elasticsearchComponent = new ElasticsearchComponent(); elasticsearchComponent.setHostAddresses("elasticsearch-rlam-service:9200"); getContext().addComponent("elasticsearch-rest", elasticsearchComponent); from("kafka:dbz?brokers={{kafka.bootstrap.address}}&groupId=apps&autoOffsetReset=earliest") .choice() .when().simple("${body} == 'null'") .log("Null!") .otherwise() .log("Message: ${body}") .to("elasticsearch-rest://elasticsearch?hostAddresses=elasticsearch-rlam-service:9200&operation=INDEX&indexName=dbz") .endChoice(); } }
错误堆栈
2022-06-20 00:11:35,837 INFO [route1] (Camel (camel-1) thread #1 - KafkaConsumer[dbz]) Null! 2022-06-20 00:11:35,849 INFO [route1] (Camel (camel-1) thread #1 - KafkaConsumer[dbz]) Message: {"last_name":"Ketchmar","id":1004,"first_name":"Anne","email":"annek@noanswer.org"} 2022-06-20 00:11:36,116 ERROR [org.apa.cam.pro.err.DefaultErrorHandler] (Camel (camel-1) thread #1 - KafkaConsumer[dbz]) Failed delivery for (MessageId: 274D73D47B2829F-0000000000000000 on ExchangeId: 274D73D47B2829F-0000000000000000). Exhausted after delivery attempt: 1 caught: java.net.ConnectException: Connection refused Message History (source location and message history is disabled) --------------------------------------------------------------------------------------------------------------------------------------- Source ID Processor Elapsed (ms) route1/route1 from[kafka://dbz?autoOffsetReset=earliest&brokers= 272 ... route1/to1 elasticsearch-rest://elasticsearch?hostAddresses=e 0 Stacktrace ---------------------------------------------------------------------------------------------------------------------------------------: java.net.ConnectException: Connection refused at org.elasticsearch.client.RestClient.extractAndWrapCause(RestClient.java:918) at org.elasticsearch.client.RestClient.performRequest(RestClient.java:299) at org.elasticsearch.client.RestClient.performRequest(RestClient.java:287) at org.elasticsearch.client.RestHighLevelClient.internalPerformRequest(RestHighLevelClient.java:1632) at org.elasticsearch.client.RestHighLevelClient.performRequest(RestHighLevelClient.java:1602) at org.elasticsearch.client.RestHighLevelClient.performRequestAndParseEntity(RestHighLevelClient.java:1572) at org.elasticsearch.client.RestHighLevelClient.index(RestHighLevelClient.java:989) at org.apache.camel.component.elasticsearch.ElasticsearchProducer.process(ElasticsearchProducer.java:170) at org.apache.camel.support.AsyncProcessorConverterHelper$ProcessorToAsyncProcessorBridge.process(AsyncProcessorConverterHelper.java:66) at org.apache.camel.processor.SendProcessor.process(SendProcessor.java:172) at org.apache.camel.processor.errorhandler.RedeliveryErrorHandler$SimpleTask.run(RedeliveryErrorHandler.java:471) at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:193) at org.apache.camel.impl.engine.DefaultReactiveExecutor.scheduleMain(DefaultReactiveExecutor.java:64) at org.apache.camel.processor.Pipeline.process(Pipeline.java:184) at org.apache.camel.impl.engine.CamelInternalProcessor.process(CamelInternalProcessor.java:399) at org.apache.camel.impl.engine.DefaultAsyncProcessorAwaitManager.process(DefaultAsyncProcessorAwaitManager.java:83) at org.apache.camel.support.AsyncProcessorSupport.process(AsyncProcessorSupport.java:41) at org.apache.camel.component.kafka.consumer.support.KafkaRecordProcessor.processExchange(KafkaRecordProcessor.java:109) at org.apache.camel.component.kafka.consumer.support.KafkaRecordProcessorFacade.processRecord(KafkaRecordProcessorFacade.java:120) at org.apache.camel.component.kafka.consumer.support.KafkaRecordProcessorFacade.processPolledRecords(KafkaRecordProcessorFacade.java:80) at org.apache.camel.component.kafka.KafkaFetchRecords.startPolling(KafkaFetchRecords.java:280) at org.apache.camel.component.kafka.KafkaFetchRecords.run(KafkaFetchRecords.java:181) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.net.ConnectException: Connection refused at java.base/sun.nio.ch.SocketChannelImpl.checkConnect(Native Method) at java.base/sun.nio.ch.SocketChannelImpl.finishConnect(SocketChannelImpl.java:777) at org.apache.http.impl.nio.reactor.DefaultConnectingIOReactor.processEvent(DefaultConnectingIOReactor.java:174) at org.apache.http.impl.nio.reactor.DefaultConnectingIOReactor.processEvents(DefaultConnectingIOReactor.java:148) at org.apache.http.impl.nio.reactor.AbstractMultiworkerIOReactor.execute(AbstractMultiworkerIOReactor.java:351) at org.apache.http.impl.nio.conn.PoolingNHttpClientConnectionManager.execute(PoolingNHttpClientConnectionManager.java:221) at org.apache.http.impl.nio.client.CloseableHttpAsyncClientBase$1.run(CloseableHttpAsyncClientBase.java:64) ... 1 more
根因
Camel K默认使用Quarkus运行时,当classpath中存在Elasticsearch REST客户端依赖时,Quarkus会自动创建一个默认连接localhost:9200的RestHighLevelClient Bean。Camel Elastisearch REST组件检测到上下文存在该Bean时,会直接复用这个实例,完全忽略代码或端点中配置的hostAddresses参数。由于集成Pod本地没有运行Elasticsearch服务,直接触发连接拒绝错误,这也是错误日志中ES请求耗时为0ms的原因——请求根本没有发往Kubernetes Service地址,而是直接打在本地回环地址。
解决方案
二选一即可:
- 禁用Quarkus默认的Elasticsearch自动配置,让Camel组件自行根据配置创建客户端。在集成的
application.properties中添加配置:
该方案下可以保留代码中对组件、端点的hostAddresses配置,注意地址前建议显式添加quarkus.elasticsearch.enabled=falsehttp://前缀(即http://elasticsearch-rlam-service:9200),避免组件默认使用HTTPS协议发起连接导致异常。 - 不修改组件配置,直接通过Quarkus配置项指定正确的Elasticsearch地址,复用自动创建的客户端。在
application.properties中添加配置:
该方案下可以删除代码中手动实例化注册ElasticsearchComponent的冗余逻辑。quarkus.elasticsearch.hosts=elasticsearch-rlam-service:9200 quarkus.elasticsearch.protocol=http
排查时可以在路由启动阶段添加日志,打印Camel上下文内Elasticsearch组件实际使用的连接地址,确认配置生效。
内容的提问来源于stack exchange,提问作者RLam
相关产品推荐
相关产品推荐

