davsclaus commented on code in PR #27543:
URL: https://github.com/apache/camel/pull/27543#discussion_r4216654009
##########
components/camel-kafka/src/main/java/org/apache/camel/processor/keyvalue/kafka/KafkaKeyValueRepository.java:
##########
@@ -461,102 +406,27 @@ protected void doStart() throws Exception {
// if you need a larger working set to be visible locally.
this.cache = LRUCacheFactory.newLRUCache(maxCacheSize);
- if (consumerConfig == null) {
- consumerConfig = new Properties();
- StringHelper.notEmpty(bootstrapServers, "bootstrapServers");
- consumerConfig.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,
bootstrapServers);
- if (groupId != null) {
- consumerConfig.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
- }
- }
-
- if (producerConfig == null) {
- producerConfig = new Properties();
- StringHelper.notEmpty(bootstrapServers, "bootstrapServers");
- producerConfig.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
bootstrapServers);
- }
-
- ObjectHelper.notNull(consumerConfig, "consumerConfig");
- ObjectHelper.notNull(producerConfig, "producerConfig");
-
- consumerConfig.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,
Boolean.FALSE.toString());
- consumerConfig.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
StringDeserializer.class.getName());
- consumerConfig.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
ByteArrayDeserializer.class.getName());
-
- consumer = new KafkaConsumer<>(consumerConfig);
-
- producerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName());
- producerConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
ByteArraySerializer.class.getName());
- producerConfig.putIfAbsent(ProducerConfig.ACKS_CONFIG, "1");
- producerConfig.putIfAbsent(ProducerConfig.BATCH_SIZE_CONFIG, "0");
- producer = new KafkaProducer<>(producerConfig);
-
- poller = new TopicPoller();
- ServiceHelper.startService(poller);
- // populate cache on startup
- StopWatch watch = new StopWatch();
- LOG.info("Syncing KafkaKeyValueRepository from topic: {} starting",
topic);
- poller.run();
- LOG.info("Syncing KafkaKeyValueRepository from topic: {} complete:
{}", topic,
- TimeUtils.printDuration(watch.taken(), true));
-
- if (!startupOnly) {
- executorService
- =
camelContext.getExecutorServiceManager().newSingleThreadExecutor(this,
"KafkaKeyValueRepositorySync");
- LOG.info("Syncing KafkaKeyValueRepository from topic: {}
continuously using background thread", topic);
- executorService.submit(poller);
- }
+ consumerConfig = KafkaChangelog.consumerConfig(consumerConfig,
bootstrapServers, groupId);
+ producerConfig = KafkaChangelog.producerConfig(producerConfig,
bootstrapServers);
+ changelog = new KafkaChangelog<>(
+ "KafkaKeyValueRepository", topic, consumerConfig,
producerConfig, ByteArrayDeserializer.class,
+ ByteArraySerializer.class, pollDurationMs,
+ startupOnly, this::addToCache);
+ changelog.start(camelContext, this);
}
@Override
protected void doStop() throws Exception {
- ServiceHelper.stopService(poller);
- if (consumer != null) {
- consumer.wakeup();
- }
- if (executorService != null && camelContext != null) {
-
camelContext.getExecutorServiceManager().shutdownNow(executorService);
- executorService = null;
+ if (changelog != null) {
+ changelog.stop();
}
- IOHelper.close(consumer, "consumer", LOG);
- IOHelper.close(producer, "producer", LOG);
LOG.debug("Stopped KafkaKeyValueRepository. Cache counter: {}",
cacheCounter.get());
}
//
-------------------------------------------------------------------------
// TopicPoller inner class
Review Comment:
Nit: this section header is left without anything under it now that
`TopicPoller` moved to `KafkaChangelog`, so the three comment lines (and the
blank line after them) can go.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]