You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

从ActiveMQ转Kafka:如何实现先生产消息再延迟消费?

Kafka 消息生产后延迟消费的实现方案

背景与需求

作为ActiveMQ资深用户,转用Kafka时希望实现以下逻辑:

  • 向主题提交100条消息
  • 等待任意时长
  • 从该主题消费这100条消息,且保证单消费者独占消费

原测试代码存在问题:必须先启动消费者并等待10秒注册完成,才能发送消息,否则消费者无法接收消息。需要调整实现,支持先生产消息,再启动消费者也能完整消费所有消息。

修改思路

  1. 调整执行顺序:先完成消息生产,再启动消费者
  2. 配置消费者偏移重置策略:添加auto.offset.reset=earliest,让新消费者从主题的最早偏移量开始消费,覆盖启动前生产的消息
  3. 控制消费终止条件:统计消费到的消息数量,达到100条后主动停止消费,避免无限循环
  4. 移除不必要的等待:删除原代码中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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 02:36:20