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

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地址,而是直接打在本地回环地址。

解决方案

二选一即可:

  1. 禁用Quarkus默认的Elasticsearch自动配置,让Camel组件自行根据配置创建客户端。在集成的application.properties中添加配置:
    quarkus.elasticsearch.enabled=false
    
    该方案下可以保留代码中对组件、端点的hostAddresses配置,注意地址前建议显式添加http://前缀(即http://elasticsearch-rlam-service:9200),避免组件默认使用HTTPS协议发起连接导致异常。
  2. 不修改组件配置,直接通过Quarkus配置项指定正确的Elasticsearch地址,复用自动创建的客户端。在application.properties中添加配置:
    quarkus.elasticsearch.hosts=elasticsearch-rlam-service:9200
    quarkus.elasticsearch.protocol=http
    
    该方案下可以删除代码中手动实例化注册ElasticsearchComponent的冗余逻辑。

排查时可以在路由启动阶段添加日志,打印Camel上下文内Elasticsearch组件实际使用的连接地址,确认配置生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 10:42:20