【发布时间】:2020-03-20 20:17:08
【问题描述】:
我是 Apache 骆驼的新手。我们正在做 POC 以使用 Camel 开发 kafka 消费者。下面是示例代码。
context.addRoutes(new RouteBuilder(){
@Override
public void configure() throws Exception {
// TODO Auto-generated method stub
from("kafka:{{consumer.topic}}?brokers={{kafka.host}}:{{kafka.port}}"
+ "&consumersCount={{consumer.consumersCount}}"
+ "&seekTo={{consumer.seekTo}}"
+ "&groupId={{consumer.group}}")
.process(new Processor() {
@Override
public void process(Exchange exchange) throws Exception {
Message message = exchange.getIn();
Object data = message.getBody();
System.out.println(data);
}
})
.to("seda:end");
});
context.start();
ConsumerTemplate template=context.createConsumerTemplate();
String info=template.receiveBody("seda:end",String.class);
System.out.println(info);
}
我遇到以下问题:
- 上下文在启动后立即停止。
- 如果我使用消费者模板轮询到端点,它不会打印任何内容,而在 .process() 内部,当我在无限循环中启动上下文时,我能够打印 kafka 消息。为什么消费者模板无法打印。
【问题讨论】:
-
github.com/apache/camel-examples/blob/master/examples/…有一个Kafka的例子,start方法没有阻塞,看这个博客tomd.xyz/camel-standalone-example如果你想继续运行Camel独立,有camel-main,看github.com/apache/camel-examples/tree/master/examples/…
-
@ClausIbsen - 感谢您的链接。根据链接 - github.com/apache/camel-examples/tree/master/examples/...,我已经添加了 Thread.sleep(),这将使骆驼在“Thread.sleep()”中指定的时间内保持运行。但我想让消费者永远运行,除非有一些维护问题。使用“Thread.sleep()”总是会限制消费者运行的时间。如果我想永远运行,是否需要将 context.start() 放入无限循环中?
标签: java apache-kafka apache-camel