【问题标题】:Interruption exception with Apache Camel route from/to Apache Kafka来自/到 Apache Kafka 的 Apache Camel 路由中断异常
【发布时间】:2020-06-10 21:10:05
【问题描述】:

我有一个正在运行的 Apache Kafka 消息代理实例,我想将其用作骆驼路线的起点/终点。开始路线时 - 似乎工作正常 - 我得到一个 InterruptedException 我不知道如何解决这个问题:

09:06:33.840 [Camel (camel-1) thread #1 - KafkaConsumer[OLOG_INBOUND]] DEBUG org.apache.kafka.clients.NetworkClient - [Consumer clientId=KafkaTrafoDataRoute, groupId=8563046a-15fa-48fe-858f-cfd68c7b921c] Completed connection to node -1. Fetching API versions.
09:06:33.840 [Camel (camel-1) thread #1 - KafkaConsumer[OLOG_INBOUND]] DEBUG org.apache.kafka.clients.NetworkClient - [Consumer clientId=KafkaTrafoDataRoute, groupId=8563046a-15fa-48fe-858f-cfd68c7b921c] Initiating API versions fetch from node -1.
09:06:33.849 [Camel (camel-1) thread #1 - KafkaConsumer[OLOG_INBOUND]] WARN org.apache.camel.component.kafka.KafkaConsumer - Interrupted while consuming OLOG_INBOUND-Thread 0 from kafka topic. Caused by: [org.apache.kafka.common.errors.InterruptException - java.lang.InterruptedException]
org.apache.kafka.common.errors.InterruptException: java.lang.InterruptedException
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.maybeThrowInterruptException(ConsumerNetworkClient.java:504)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:287)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:242)
    at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:218)
    at org.apache.kafka.clients.consumer.internals.AbstractCoordinator.ensureCoordinatorReady(AbstractCoordinator.java:230)
    at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:314)
    at org.apache.kafka.clients.consumer.KafkaConsumer.updateAssignmentMetadataIfNeeded(KafkaConsumer.java:1218)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1181)
    at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1115)
    at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.doRun(KafkaConsumer.java:293)
    at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.run(KafkaConsumer.java:215)
    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:834)
Caused by: java.lang.InterruptedException: null
    ... 16 common frames omitted
09:06:33.850 [Camel (camel-1) thread #1 - KafkaConsumer[OLOG_INBOUND]] INFO org.apache.camel.component.kafka.KafkaConsumer - Unsubscribing OLOG_INBOUND-Thread 0 from topic OLOG_INBOUND
09:06:33.850 [Camel (camel-1) thread #1 - KafkaConsumer[OLOG_INBOUND]] DEBUG org.apache.kafka.clients.consumer.KafkaConsumer - [Consumer clientId=KafkaTrafoDataRoute, groupId=8563046a-15fa-48fe-858f-cfd68c7b921c] Unsubscribed all topics or patterns and assigned partitions
09:06:33.850 [Camel (camel-1) thread #1 - KafkaConsumer[OLOG_INBOUND]] DEBUG org.apache.camel.component.kafka.KafkaConsumer - Closing OLOG_INBOUND-Thread 0

路由是从这样的主程序调用的(注意:我尽量不使用 Spring,因为它在我的情况下会导致太多问题):

public static void main(String[] args) {
    CamelContext camelContext = new DefaultCamelContext();

    try {
        camelContext.addRoutes(new MyKafkaRoute());
    } catch (Exception e) {
        // TODO Auto-generated catch block
        e.printStackTrace();
    }

    try {
        camelContext.start();
    } catch (Exception e) {
        // TODO Auto-generated catch block
        e.printStackTrace();
    }
}

这是MyKafkaRoute的代码:

public class MyKafkaRoute extends RouteBuilder {

    private String consumerEndpoint = "kafka:TOPIC_NAME?"    //
            + "brokers=server:port"                   //
            + "&clientId=myKafkaRoute";

    private String emitterEndpoint  = "kafka:TOPIC_NAME?"   //
            + "brokers=server:port"                   //
            + "&clientId=myKafkaRoute";

    @Override
    public void configure() throws Exception {

        from(consumerEndpoint) //
                .process(... processing ...) //
                .to(emitterEndpoint) //
                .onException(Exception.class) //
                .useOriginalMessage() //
                .handled(true) //
                .to("stream:out");
    }
}

【问题讨论】:

    标签: java apache-kafka apache-camel


    【解决方案1】:

    看起来(至少)其中一个端点可能不可用。除此之外,将骆驼路线修剪到 from()。 ... 。至()。 ... .to() 以更清楚地了解发送给谁的内容。

    【讨论】:

      【解决方案2】:

      我的问题是我自己的错:

      在上面的主程序中,我启动消费者

      camelContext.start();
      

      但在此语句之后,我的测试程序立即结束导致InterruptedExecption。我增强了代码:

      try {
          camelContext.start();
      } catch (Exception e) {
          // TODO Auto-generated catch block
          e.printStackTrace();
      }
      
      try {
          Thread.sleep(5 * 60 * 1000);
      } catch (InterruptedException e) {
          // TODO Auto-generated catch block
          e.printStackTrace();
      }
      
      try {
          camelContext.stop();
      } catch (Exception e) {
          // TODO Auto-generated catch block
          e.printStackTrace();
      }
      

      现在我可以看到正在发送到代理的消息。

      【讨论】:

      • 嗨@WolfiG,Apache Camel 听起来很有趣!你介意分享你的用例吗?提前致谢!
      • 我的用例如下:我在一家大型加速器设施工作。在这样的加速器中,可以触发束诊断仪器来读取加速器中的离子束电流。我使用 Apache Kafka 作为代理从其他应用程序获取消息,请求读取光束诊断(设备名称、时间等)。 apache 路由在加速器控制系统的框架中运行,并调用代码进行实际读出。读出的值被传回 Kafka,其他应用程序可以从中获取读出的数据。
      猜你喜欢
      • 1970-01-01
      • 2016-09-29
      • 1970-01-01
      • 2021-03-10
      • 1970-01-01
      • 2018-07-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多