【问题标题】:Camel route from Kafka to JMS not working (JMS->Kafka->JMS)从 Kafka 到 JMS 的骆驼路线不起作用(JMS->Kafka->JMS)
【发布时间】:2021-01-05 12:48:37
【问题描述】:

在我的 Camel Router.java 中,我有下一条路线

from("jms:topic:test.source.topic?asyncConsumer=true")
            .log("Message: ${body}")
            .to("kafka:testing?brokers=192.168.0.100:9092");

from("kafka:testing?brokers=192.168.0.100:9092")
            .log("Message received from Kafka : ${body}")
            .log("    on the topic ${headers[kafka.TOPIC]}")
            .log("    on the partition ${headers[kafka.PARTITION]}")
            .log("    with the offset ${headers[kafka.OFFSET]}")
            .log("    with the key ${headers[kafka.KEY]}")
            // manually set JMSDeliveryMode (1 - NON_PERSISTENT, 2 - PERSISTENT)
            .process(new Processor() {
                @Override
                public void process(Exchange exchange) throws Exception {
                    exchange.getIn().setHeader("JMSDeliveryMode", "1");
               }
            })
            .to("jms:topic:test.sink.topic");

上述骆驼路线出现问题。如果我通过使用bin/kafka-console-producer.sh --topic testing --bootstrap-server localhost:9092 运行的Kafka 生产者向主题testing 发送一些消息,则从Kafka 到JMS 的路由工作正常。所以这些链接的骆驼路线有些问题。

在 Camel pom.xml 中有 Spring Boot camel-kafka-starter 和 camel-jms-starter 依赖项。

当我使用 Maven 启动 Spring Boot Camel 并将一些消息从 Kafka 生产者发送到 Kafka 代理 testing 主题时,我可以看到 Kafka 代理收到了该消息,并且上面的日志打印正常。

.to("jms:topic:test.sink.topic"); 行出现错误,我不知道是什么意思。

ERROR 28642 --- [aConsumer[testing]] o.a.c.p.e.DefaultErrorHandler: Failed delivery for (MessageId: ID-PCID on ExchangeId: ID-PCID). Exhausted after delivery attempt: 1 caught: org.springframework.jms.UncategorizedJmsException: Uncategorized exception occurred during JMS processing; nested exception is javax.jms.JMSException: AMQ139015: Illegal deliveryMode value: 0

org.springframework.jms.UncategorizedJmsException: Uncategorized exception occurred during JMS processing; nested exception is javax.jms.JMSException: AMQ139015: Illegal deliveryMode value: 0
    at org.springframework.jms.support.JmsUtils.convertJmsAccessException(JmsUtils.java:311) ~[spring-jms-5.2.6.RELEASE.jar:5.2.6.RELEASE]
    at org.springframework.jms.support.JmsAccessor.convertJmsAccessException(JmsAccessor.java:185) ~[spring-jms-5.2.6.RELEASE.jar:5.2.6.RELEASE]
    at org.springframework.jms.core.JmsTemplate.execute(JmsTemplate.java:507) ~[spring-jms-5.2.6.RELEASE.jar:5.2.6.RELEASE]
    at org.apache.camel.component.jms.JmsConfiguration$CamelJmsTemplate.send(JmsConfiguration.java:525) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsProducer.doSend(JmsProducer.java:438) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsProducer.processInOnly(JmsProducer.java:392) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsProducer.process(JmsProducer.java:155) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.processor.SendProcessor.process(SendProcessor.java:168) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.processor.errorhandler.RedeliveryErrorHandler$SimpleTask.run(RedeliveryErrorHandler.java:395) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.impl.engine.DefaultReactiveExecutor$Worker.schedule(DefaultReactiveExecutor.java:148) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.impl.engine.DefaultReactiveExecutor.scheduleMain(DefaultReactiveExecutor.java:60) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.processor.Pipeline.process(Pipeline.java:147) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.processor.CamelInternalProcessor.process(CamelInternalProcessor.java:286) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.impl.engine.DefaultAsyncProcessorAwaitManager.process(DefaultAsyncProcessorAwaitManager.java:83) ~[camel-base-3.4.0.jar:3.4.0]
    at org.apache.camel.support.AsyncProcessorSupport.process(AsyncProcessorSupport.java:40) ~[camel-support-3.4.0.jar:3.4.0]
    at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.doRun(KafkaConsumer.java:346) ~[camel-kafka-3.4.0.jar:3.4.0]
    at org.apache.camel.component.kafka.KafkaConsumer$KafkaFetchRecords.run(KafkaConsumer.java:222) ~[camel-kafka-3.4.0.jar:3.4.0]
    at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) ~[na:na]
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) ~[na:na]
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) ~[na:na]
    at java.base/java.lang.Thread.run(Thread.java:834) ~[na:na]
Caused by: javax.jms.JMSException: AMQ139015: Illegal deliveryMode value: 0
    at org.apache.activemq.artemis.jms.client.ActiveMQMessage.setJMSDeliveryMode(ActiveMQMessage.java:450) ~[artemis-jms-client-2.12.0.jar:2.12.0]
    at org.apache.camel.component.jms.JmsMessageHelper.setJMSDeliveryMode(JmsMessageHelper.java:435) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsBinding.appendJmsProperty(JmsBinding.java:393) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsBinding.appendJmsProperties(JmsBinding.java:371) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsBinding.makeJmsMessage(JmsBinding.java:346) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsProducer$2.createMessage(JmsProducer.java:325) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsConfiguration$CamelJmsTemplate.doSendToDestination(JmsConfiguration.java:561) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.apache.camel.component.jms.JmsConfiguration$CamelJmsTemplate.lambda$send$0(JmsConfiguration.java:527) ~[camel-jms-3.4.0.jar:3.4.0]
    at org.springframework.jms.core.JmsTemplate.execute(JmsTemplate.java:504) ~[spring-jms-5.2.6.RELEASE.jar:5.2.6.RELEASE]
    ... 19 common frames omitted

当从 JMS 向 Kafka 发送一些消息时,路由工作正常。

【问题讨论】:

  • 我已经用完整的错误记录、骆驼路线和有关问题的更多信息更新了问题。使用 -e 和 -X 运行 maven 并没有提供有关错误的更多详细信息。
  • 我根据堆栈跟踪更新了我的答案。

标签: spring-boot apache-kafka apache-camel jms activemq-artemis


【解决方案1】:

以下是相关的错误信息:

AMQ139015: Illegal deliveryMode value: 0

此错误消息意味着有东西调用了 JMS 方法 javax.jms.Message#setJMSDeliveryMode,其值为 00 的值无效。有效的交付模式值由javax.jms.DeliveryMode 定义:

堆栈跟踪表明这个无效值是由 Camel 根据输入消息设置的。见this code。为了解决这个问题,您需要确定在 org.apache.camel.Message 上设置标头 JMSDeliveryMode 的位置。

【讨论】:

  • 我没有确定标题 JMSDeliveryMode 的确切更改位置。在从 JMS 到 Kafka 的骆驼路线中,我添加了 getHeader() 的日志并看到了JMSDeliveryMode=1。在从 Kafka 到 JMS 的骆驼路线中也做了同样的事情,看到了 {JMSDeliveryMode=[B@6069793c ...... kafka.HEADERS=RecordHeaders(headers = [RecordHeader(key = JMSDeliveryMode, value = [0, 0, 0, 1]) ...
  • 作为一种解决方案(我认为不是一个好的解决方案),在从 Kafka 到 JMS 的 Camel 路线中,我添加了一个进程并手动设置 JMSDeliveryMode,就像上面有问题的更新代码一样。
猜你喜欢
  • 2012-11-10
  • 1970-01-01
  • 2012-12-01
  • 1970-01-01
  • 2018-06-09
  • 2013-05-23
  • 2018-03-07
  • 2013-04-19
  • 1970-01-01
相关资源
最近更新 更多