gnodet-bot commented on code in PR #27543:
URL: https://github.com/apache/camel/pull/27543#discussion_r4216076544


##########
components/camel-kafka/src/main/java/org/apache/camel/processor/idempotent/kafka/KafkaChangelog.java:
##########
@@ -0,0 +1,271 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.processor.idempotent.kafka;
+
+import java.time.Duration;
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.Consumer;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.RuntimeCamelException;
+import org.apache.camel.support.service.ServiceHelper;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.util.IOHelper;
+import org.apache.camel.util.ObjectHelper;
+import org.apache.camel.util.StopWatch;
+import org.apache.camel.util.StringHelper;
+import org.apache.camel.util.TimeUtils;
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.consumer.ConsumerRecords;
+import org.apache.kafka.clients.consumer.KafkaConsumer;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.clients.producer.Producer;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.PartitionInfo;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.WakeupException;
+import org.apache.kafka.common.serialization.Deserializer;
+import org.apache.kafka.common.serialization.Serializer;
+import org.apache.kafka.common.serialization.StringDeserializer;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * A Kafka topic used as the changelog of the local cache of a repository, 
shared by the Kafka repositories
+ * ({@link KafkaIdempotentRepository} and the KafkaKeyValueRepository).
+ * <p/>
+ * On start, the whole topic is read from the beginning and every record is 
applied to the cache; then, unless the
+ * repository syncs on startup only, the records written by the other 
instances keep being applied by a background
+ * thread. The changes of this instance are written to the topic synchronously.
+ *
+ * @param <V> the type of the record values
+ */
+public final class KafkaChangelog<V> {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(KafkaChangelog.class);
+
+    private final String name;
+    private final String topic;
+    private final Properties consumerConfig;
+    private final Properties producerConfig;
+    private final Class<? extends Deserializer<V>> valueDeserializer;
+    private final Class<? extends Serializer<V>> valueSerializer;
+    private final int pollDurationMs;
+    private final boolean startupOnly;
+    private final Consumer<ConsumerRecord<String, V>> recordHandler;
+
+    private CamelContext camelContext;
+    private org.apache.kafka.clients.consumer.Consumer<String, V> consumer;
+    private Producer<String, V> producer;
+    private TopicPoller poller;
+    private ExecutorService executorService;
+
+    /**
+     * @param name              the name of the repository, used in the logs 
and in the name of the background thread
+     * @param topic             the changelog topic
+     * @param consumerConfig    the properties of the consumer, see {@link 
#consumerConfig(Properties, String, String)}
+     * @param producerConfig    the properties of the producer, see {@link 
#producerConfig(Properties, String)}
+     * @param valueDeserializer the deserializer of the record values (the 
keys are strings)
+     * @param valueSerializer   the serializer of the record values (the keys 
are strings)
+     * @param pollDurationMs    the poll duration of the consumer
+     * @param startupOnly       whether to read the topic on start only, or to 
keep reading it in the background
+     * @param recordHandler     applies a record of the topic to the local 
cache
+     */
+    public KafkaChangelog(String name, String topic, Properties 
consumerConfig, Properties producerConfig,
+                          Class<? extends Deserializer<V>> valueDeserializer, 
Class<? extends Serializer<V>> valueSerializer,
+                          int pollDurationMs, boolean startupOnly, 
Consumer<ConsumerRecord<String, V>> recordHandler) {
+        this.name = name;
+        this.topic = topic;
+        this.consumerConfig = consumerConfig;
+        this.producerConfig = producerConfig;
+        this.valueDeserializer = valueDeserializer;
+        this.valueSerializer = valueSerializer;
+        this.pollDurationMs = pollDurationMs;
+        this.startupOnly = startupOnly;
+        this.recordHandler = recordHandler;
+    }
+
+    /**
+     * The consumer properties of a repository: the configured properties, 
otherwise properties with the bootstrap
+     * servers and the group id.
+     */
+    public static Properties consumerConfig(Properties consumerConfig, String 
bootstrapServers, String groupId) {
+        if (consumerConfig != null) {
+            return consumerConfig;
+        }
+        Properties answer = new Properties();
+        StringHelper.notEmpty(bootstrapServers, "bootstrapServers");
+        answer.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
+        if (groupId != null) {
+            answer.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
+        }
+        return answer;
+    }
+
+    /**
+     * The producer properties of a repository: the configured properties, 
otherwise properties with the bootstrap
+     * servers.
+     */
+    public static Properties producerConfig(Properties producerConfig, String 
bootstrapServers) {
+        if (producerConfig != null) {
+            return producerConfig;
+        }
+        Properties answer = new Properties();
+        StringHelper.notEmpty(bootstrapServers, "bootstrapServers");
+        answer.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
+        return answer;
+    }
+
+    /**
+     * Creates the consumer and the producer, reads the topic from the 
beginning, and starts the background thread
+     * unless the repository syncs on startup only.
+     *
+     * @param camelContext the CamelContext
+     * @param source       the owner of the background thread
+     */
+    public void start(CamelContext camelContext, Object source) {
+        this.camelContext = camelContext;
+
+        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, 
valueDeserializer.getName());
+        consumer = new KafkaConsumer<>(consumerConfig);
+
+        producerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
StringSerializer.class.getName());
+        producerConfig.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
valueSerializer.getName());
+        // set up the producer to remove all batching on send, we want all 
sends to be fully synchronous
+        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 to be ready
+        StopWatch watch = new StopWatch();
+        LOG.info("Syncing {} from topic: {} starting", name, topic);
+        poller.run();
+        LOG.info("Syncing {} from topic: {} complete: {}", name, topic, 
TimeUtils.printDuration(watch.taken(), true));
+
+        if (!startupOnly) {
+            // continue sync job in background
+            executorService = 
camelContext.getExecutorServiceManager().newSingleThreadExecutor(source, name + 
"Sync");
+            LOG.info("Syncing {} from topic: {} continuously using background 
thread", name, topic);
+            executorService.submit(poller);
+        }
+    }
+
+    /**
+     * Stops the background thread, and closes the consumer and the producer.
+     */
+    public void stop() {
+        ServiceHelper.stopService(poller);
+        if (consumer != null) {
+            consumer.wakeup();
+        }
+        if (executorService != null && camelContext != null) {
+            
camelContext.getExecutorServiceManager().shutdownNow(executorService);
+            executorService = null;
+        }
+        IOHelper.close(consumer, "consumer", LOG);
+        IOHelper.close(producer, "producer", LOG);
+    }
+
+    /**
+     * Writes a record to the topic, and waits until it is written.
+     */
+    public void send(String key, V value) {
+        try {
+            ObjectHelper.notNull(producer, "producer");
+            producer.send(new ProducerRecord<>(topic, key, value)).get(); // 
sync send
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+            throw new RuntimeCamelException(e);
+        } catch (ExecutionException e) {
+            throw new RuntimeCamelException(e);
+        }
+    }
+
+    /**
+     * To use the given producer, for testing.
+     */
+    void setProducer(Producer<String, V> producer) {
+        this.producer = producer;
+    }
+
+    private void populateCache() {
+        LOG.debug("Getting partitions of topic {}", topic);
+        List<PartitionInfo> partitionInfos = consumer.partitionsFor(topic);
+        Collection<TopicPartition> partitions = partitionInfos.stream()
+                .map(pi -> new TopicPartition(pi.topic(), pi.partition()))
+                .toList();
+
+        LOG.debug("Assigning consumer to partitions {}", partitions);
+        consumer.assign(partitions);
+
+        LOG.debug("Seeking consumer to beginning of partitions {}", 
partitions);
+        consumer.seekToBeginning(partitions);
+
+        Map<TopicPartition, Long> endOffsets = consumer.endOffsets(partitions);
+        LOG.debug("Consuming records from partitions {} till end offsets {}", 
partitions, endOffsets);
+        while (!KafkaConsumerUtil.isReachedOffsets(consumer, endOffsets)) {
+            ConsumerRecords<String, V> consumerRecords = 
consumer.poll(Duration.ofMillis(pollDurationMs));
+            for (ConsumerRecord<String, V> consumerRecord : consumerRecords) {
+                recordHandler.accept(consumerRecord);
+            }
+        }
+    }
+
+    private class TopicPoller extends ServiceSupport implements Runnable {
+
+        private final AtomicBoolean init = new AtomicBoolean();
+
+        @Override
+        public void run() {
+            if (init.compareAndSet(false, true)) {
+                // sync cache on startup
+                LOG.debug("TopicPoller populating cache on startup");
+                populateCache();
+                LOG.debug("TopicPoller populated cache on startup complete");
+                return;
+            }
+
+            LOG.debug("TopicPoller running");
+            while (isRunAllowed()) {
+                try {
+                    ConsumerRecords<String, V> consumerRecords = 
consumer.poll(Duration.ofMillis(pollDurationMs));
+                    for (ConsumerRecord<String, V> consumerRecord : 
consumerRecords) {
+                        recordHandler.accept(consumerRecord);
+                    }
+                } catch (WakeupException e) {
+                    LOG.debug("TopicPoller woken up during shutdown");
+                } catch (Exception e) {
+                    LOG.warn("TopicPoller error syncing due to: " + 
e.getMessage() + ". This exception is ignored.", e);

Review Comment:
   ⚠️ **Eager string concatenation in `LOG.warn`:** The message is built via 
`+` even when `WARN` is disabled. Prefer the parameterized form so the 
concatenation is deferred:
   
   ```suggestion
                       LOG.warn("TopicPoller error syncing due to: {}. This 
exception is ignored.", e.getMessage(), e);
   ```



##########
components/camel-kafka/src/main/java/org/apache/camel/processor/idempotent/kafka/KafkaChangelog.java:
##########
@@ -0,0 +1,271 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.processor.idempotent.kafka;
+
+import java.time.Duration;
+import java.util.Collection;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.function.Consumer;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.RuntimeCamelException;
+import org.apache.camel.support.service.ServiceHelper;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.util.IOHelper;
+import org.apache.camel.util.ObjectHelper;
+import org.apache.camel.util.StopWatch;
+import org.apache.camel.util.StringHelper;
+import org.apache.camel.util.TimeUtils;
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.ConsumerRecord;
+import org.apache.kafka.clients.consumer.ConsumerRecords;
+import org.apache.kafka.clients.consumer.KafkaConsumer;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.clients.producer.Producer;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.PartitionInfo;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.errors.WakeupException;
+import org.apache.kafka.common.serialization.Deserializer;
+import org.apache.kafka.common.serialization.Serializer;
+import org.apache.kafka.common.serialization.StringDeserializer;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * A Kafka topic used as the changelog of the local cache of a repository, 
shared by the Kafka repositories
+ * ({@link KafkaIdempotentRepository} and the KafkaKeyValueRepository).
+ * <p/>
+ * On start, the whole topic is read from the beginning and every record is 
applied to the cache; then, unless the
+ * repository syncs on startup only, the records written by the other 
instances keep being applied by a background
+ * thread. The changes of this instance are written to the topic synchronously.
+ *
+ * @param <V> the type of the record values
+ */
+public final class KafkaChangelog<V> {
+

Review Comment:
   💡 **Package placement:** `KafkaChangelog` is shared infrastructure — it's 
imported by `KafkaKeyValueRepository` which lives in 
`processor.keyvalue.kafka`, yet the class itself lives under 
`processor.idempotent.kafka`. This creates an awkward cross-package dependency 
on the idempotent package for something that isn't idempotent-specific. 
Consider placing it in a neutral parent package (e.g. 
`processor.kafka.changelog` or `processor.kafka.common`) so neither consumer 
has to reach into the other's namespace.



-- 
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]

Reply via email to