Я не могу использовать сообщения Kafka, работающие в локальном режиме с использованием верблюда Apache.Я не просматриваю какие-либо потребительские свойства Kafka при запуске приложения.Ниже приведен мой потребительский код Kafka
import org.apache.camel.CamelContext;
import org.apache.camel.builder.RouteBuilder;
import org.apache.camel.component.properties.PropertiesComponent;
import org.apache.camel.impl.DefaultCamelContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public final class MessageConsumerClient {
private static final Logger LOG = LoggerFactory.getLogger(MessageConsumerClient.class);
private MessageConsumerClient() {
}
public static void main(String[] args) throws Exception {
LOG.info("About to run Kafka-camel integration...");
CamelContext camelContext = new DefaultCamelContext();
// Add route to send messages to Kafka
camelContext.addRoutes(new RouteBuilder() {
public void configure() {
PropertiesComponent pc = getContext().getComponent("properties", PropertiesComponent.class);
pc.setLocation("classpath:application.properties");
log.info("About to start route: Kafka Server -> Log ");
from("kafka:{{consumer.topic}}?brokers={{kafka.host}}:{{kafka.port}}"
+ "&maxPollRecords={{consumer.maxPollRecords}}" + "&consumersCount={{consumer.consumersCount}}"
+ "&seekTo={{consumer.seekTo}}" + "&groupId={{consumer.group}}").routeId("FromKafka")
.log("${body}");
}
});
camelContext.start();
Thread.sleep(5 * 60 * 1000);
camelContext.stop();
}
}
Я использую следующие зависимости в pom.xml:
1.Camel core - 2.21.1 2.Camel Kafka - 2.21.1
Pom.xml
<dependencies>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-core</artifactId>
<version>2.21.1</version>
</dependency>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-kafka</artifactId>
<version>2.21.1</version>
</dependency>
</dependencies>
Любая помощь будет принята с благодарностью, так как я не могу просматривать журналы при выполнении.