从ActiveMQ转Kafka:如何实现先生产消息再延迟消费?
Kafka 消息生产后延迟消费的实现方案
背景与需求
作为ActiveMQ资深用户,转用Kafka时希望实现以下逻辑:
- 向主题提交100条消息
- 等待任意时长
- 从该主题消费这100条消息,且保证单消费者独占消费
原测试代码存在问题:必须先启动消费者并等待10秒注册完成,才能发送消息,否则消费者无法接收消息。需要调整实现,支持先生产消息,再启动消费者也能完整消费所有消息。
修改思路
- 调整执行顺序:先完成消息生产,再启动消费者
- 配置消费者偏移重置策略:添加
auto.offset.reset=earliest,让新消费者从主题的最早偏移量开始消费,覆盖启动前生产的消息 - 控制消费终止条件:统计消费到的消息数量,达到100条后主动停止消费,避免无限循环
- 移除不必要的等待:删除原代码中
Thread.sleep(10000)的强制等待逻辑
修改后的完整代码
import java.time.Duration; import java.util.Collections; import java.util.Properties; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; 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.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringDeserializer; import org.apache.kafka.common.serialization.StringSerializer; import org.junit.After; import org.junit.Before; import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testcontainers.containers.KafkaContainer; import org.testcontainers.utility.DockerImageName; public class KafkaTest { private static final Logger LOG = LoggerFactory.getLogger(KafkaTest.class); public static final String MY_GROUP_ID = "my-group-id"; public static final String TOPIC = "topic"; private static final int MESSAGE_COUNT = 100; KafkaContainer kafka = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:6.2.1")); @Before public void before() { kafka.start(); } @After public void after() { kafka.close(); } @Test public void testPipes() throws ExecutionException, InterruptedException { ExecutorService es = Executors.newCachedThreadPool(); // 1. 先生产100条消息 Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka.getBootstrapServers()); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); Future<Void> producerFuture = es.submit(() -> { try (KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps)) { for (int counter = 0; counter <= MESSAGE_COUNT; counter++) { String msg = "Message " + counter; producer.send(new ProducerRecord<>(TOPIC, msg)).get(); // 同步发送确保消息写入Kafka LOG.info("Sent message: {}", msg); } } catch (Exception e) { LOG.error("Failed to send message by the producer", e); } return null; }); // 等待生产完成 producerFuture.get(); LOG.info("All {} messages have been produced", MESSAGE_COUNT); // 2. 启动消费者消费消息 Properties consumerProps = new Properties(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka.getBootstrapServers()); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, MY_GROUP_ID); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 关键配置:消费者启动时从最早的偏移量开始消费 consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); Future<Void> consumerFuture = es.submit(() -> { int consumedCount = 0; try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps)) { consumer.subscribe(Collections.singletonList(TOPIC)); while (consumedCount <= MESSAGE_COUNT) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { LOG.info("Thread: {}, Topic: {}, Partition: {}, Offset: {}, key: {}, value: {}", Thread.currentThread().getName(), record.topic(), record.partition(), record.offset(), record.key(), record.value().toUpperCase()); consumedCount++; } } LOG.info("Successfully consumed {} messages", consumedCount); } catch (Exception e) { LOG.error("Consumer error", e); } return null; }); // 等待消费完成 consumerFuture.get(); es.shutdown(); } }
关键修改说明
- 偏移重置配置:
auto.offset.reset=earliest是核心,当消费者组第一次消费主题时,会从主题的起始位置拉取所有历史消息,解决了"先生产后消费"时消息丢失的问题 - 同步发送消息:生产者发送时调用
.get()确保消息成功写入Kafka,避免因异步发送未完成就启动消费者导致的漏消费 - 消费终止逻辑:通过
consumedCount统计消费数量,达到目标后退出循环,避免无限阻塞 - 执行顺序调整:先完成消息生产,再启动消费者,完全匹配需求中的"生产→等待→消费"流程
内容的提问来源于stack exchange,提问作者Nicholas DiPiazza
相关产品推荐
相关产品推荐

