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

Spring Kafka通过ConcurrentKafkaListenerContainerFactory实现自定义偏移管理与容器启停

Spring Kafka 动态容器自定义偏移量+自动关停实现方案

功能实现逻辑

1 消费完所有消息后自动关停容器

你已配置IdleEventInterval空闲事件触发间隔,只需补充两个逻辑:

  • 用本地缓存存储动态创建的容器实例,方便空闲事件触发时定位到对应容器
  • 空闲事件触发时,调用AdminClient查询目标分区的最大偏移量,和当前已消费的最大偏移量比对,一致则调用container.stop()关停容器

2 偏移量持久化到存储

  • 先创建偏移量存储表(以MySQL为例):
CREATE TABLE kafka_offset (
  id BIGINT AUTO_INCREMENT PRIMARY KEY,
  group_id VARCHAR(255) NOT NULL,
  topic VARCHAR(255) NOT NULL,
  partition INT NOT NULL,
  offset BIGINT NOT NULL,
  update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  UNIQUE KEY uk_g_t_p (group_id, topic, partition)
);
  • 修改消费监听逻辑,批量消费完成后,将当前批次的最大偏移量+1(下一次要消费的位置)写入上述表,全程关闭Kafka自动偏移量提交,完全自行管理偏移量。

3 启动时加载自定义偏移量

动态创建容器时,先从偏移量表查询对应groupId+topic+partition的偏移量记录,存在则直接指定起始偏移量创建TopicPartitionOffset实例,不存在则按照默认策略(比如从最早偏移量开始)消费。


完整修改后代码

// 补全对应类的import即可运行
@SpringBootApplication
public class KafkaApp{
    public static void main(String[] args) {
        SpringApplication.run(KafkaApp.class, args);
    }

    @Bean
    public NewTopic topic() {
        return TopicBuilder.name("testTopic").partitions(1).replicas(1).build();
    }
    
    // 配置自定义容器工厂,关闭自动偏移提交,设置手动ack
    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory(ConsumerFactory<String, String> consumerFactory) {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory);
        factory.setBatchListener(true);
        factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);
        // 关闭Kafka自动偏移提交
        Map<String, Object> props = consumerFactory.getConfigurationProperties();
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
        return factory;
    }

}

@Component
class Listener {
    private static final Logger log = LoggerFactory.getLogger(Listener.class);
    private static final Method otherListen;
    // 缓存动态创建的容器:key为groupId:topic:partition,value为容器实例
    private final Map<String, ConcurrentMessageListenerContainer<String, String>> containerCache = new ConcurrentHashMap<>();
    // 注入JDBC操作工具,可替换为你自己的DAO/Repository实现
    private final JdbcTemplate jdbcTemplate;

    static {
        try {
            // 方法参数新增ack对象和原始消费记录列表,用于拿偏移量和手动确认
            otherListen = Listener.class.getDeclaredMethod("otherListen", List.class, Acknowledgment.class, List.class);
        }
        catch (NoSuchMethodException | SecurityException ex) {
            throw new IllegalStateException(ex);
        }
    }

    private final ConcurrentKafkaListenerContainerFactory<String, String> factory;
    private final MessageHandlerMethodFactory methodFactory;
    private final KafkaAdmin admin;
    private final KafkaTemplate<String, String> template;

    public Listener(ConcurrentKafkaListenerContainerFactory<String, String> factory, KafkaAdmin admin,
                    KafkaTemplate<String, String> template, KafkaListenerAnnotationBeanPostProcessor<?, ?> bpp,
                    JdbcTemplate jdbcTemplate) {
        this.factory = factory;
        this.admin = admin;
        this.template = template;
        this.methodFactory = bpp.getMessageHandlerMethodFactory();
        this.jdbcTemplate = jdbcTemplate;
    }

    @KafkaListener(id = "myId", topics = "testTopic")
    public void listen(String topicName) {
        int partition = 0;
        String groupId = "group.for." + topicName;
        String cacheKey = String.join(":", groupId, topicName, String.valueOf(partition));
        
        try (AdminClient client = AdminClient.create(this.admin.getConfigurationProperties())) {
            NewTopic topic = TopicBuilder.name(topicName).partitions(1).replicas(1).build();
            client.createTopics(List.of(topic)).all().get(10, TimeUnit.SECONDS);
        }
        catch (Exception e) {
            log.error("创建topic失败", e);
            return;
        }
        
        // 启动时查询自定义偏移量
        Long startOffset = getOffsetFromDB(groupId, topicName, partition);
        TopicPartitionOffset tpo;
        if (startOffset != null) {
            // 存在自定义偏移量,从指定偏移量开始消费
            tpo = new TopicPartitionOffset(topicName, partition, startOffset);
        } else {
            // 不存在则从最早偏移量开始
            tpo = new TopicPartitionOffset(topicName, partition, OffsetResetStrategy.EARLIEST);
        }
        
        ConcurrentMessageListenerContainer<String, String> container =
                this.factory.createContainer(tpo);
        BatchMessagingMessageListenerAdapter<String, String> adapter =
                new BatchMessagingMessageListenerAdapter<>(this, otherListen);
        adapter.setHandlerMethod(new HandlerAdapter(
                this.methodFactory.createInvocableHandlerMethod(this, otherListen)));
        FilteringBatchMessageListenerAdapter<String, String> filtered =
                new FilteringBatchMessageListenerAdapter<>(adapter, record -> !record.key().equals("foo"));
        container.getContainerProperties().setMessageListener(filtered);
        container.getContainerProperties().setGroupId(groupId);
        container.setBeanName(topicName + ".container");
        container.getContainerProperties().setIdleEventInterval(3000L);
        // 容器存入缓存
        containerCache.put(cacheKey, container);
        // 容器关停后自动移除缓存
        container.stopAbortable(() -> containerCache.remove(cacheKey));
        container.start();
        
        IntStream.range(0, 10).forEach(i -> this.template.send(topicName, 0, i % 2 == 0 ? "foo" : "bar", "test" + i));
    }

    void otherListen(List<String> others, Acknowledgment ack, List<ConsumerRecord<String, String>> records) {
        log.info("接收到消息: {}", others);
        if (records.isEmpty()) {
            ack.acknowledge();
            return;
        }
        // 消费完成后更新偏移量到数据库
        String groupId = "group.for." + records.get(0).topic();
        int partition = records.get(0).partition();
        String topic = records.get(0).topic();
        // 取当前批次最大偏移量,+1为下一次要消费的位置
        long maxOffset = records.stream().mapToLong(ConsumerRecord::offset).max().getAsLong() + 1;
        saveOffsetToDB(groupId, topic, partition, maxOffset);
        ack.acknowledge();
    }
    
    @EventListener
    public void eventHandler(ListenerContainerIdleEvent event) {
        log.info("{} 毫秒未收到消息", event.getIdleTime());
        // 获取当前事件对应的容器消费的分区信息
        Collection<TopicPartition> assigned = event.getConsumer().assignment();
        if (assigned.isEmpty()) {
            return;
        }
        TopicPartition tp = assigned.iterator().next();
        String groupId = event.getContainerProperties().getGroupId();
        String cacheKey = String.join(":", groupId, tp.topic(), String.valueOf(tp.partition()));
        ConcurrentMessageListenerContainer<String, String> container = containerCache.get(cacheKey);
        if (container == null) {
            return;
        }
        // 比对当前偏移量和分区最大偏移量,一致则关停容器
        try (AdminClient client = AdminClient.create(admin.getConfigurationProperties())) {
            Map<TopicPartition, Long> endOffsets = client.endOffsets(List.of(tp)).all().get(10, TimeUnit.SECONDS);
            Long endOffset = endOffsets.get(tp);
            Long currentOffset = getOffsetFromDB(groupId, tp.topic(), tp.partition());
            if (currentOffset != null && currentOffset.equals(endOffset)) {
                log.info("所有消息消费完成,关停容器:{}", container.getBeanName());
                container.stop();
            }
        } catch (Exception e) {
            log.error("检查偏移量失败", e);
        }
    }
    
    // 从数据库查询偏移量
    private Long getOffsetFromDB(String groupId, String topic, int partition) {
        try {
            return jdbcTemplate.queryForObject("SELECT offset FROM kafka_offset WHERE group_id = ? AND topic = ? AND partition = ?",
                    Long.class, groupId, topic, partition);
        } catch (EmptyResultDataAccessException e) {
            return null;
        }
    }
    
    // 保存偏移量到数据库
    private void saveOffsetToDB(String groupId, String topic, int partition, long offset) {
        jdbcTemplate.update("INSERT INTO kafka_offset (group_id, topic, partition, offset) VALUES (?, ?, ?, ?) " +
                        "ON DUPLICATE KEY UPDATE offset = ?",
                groupId, topic, partition, offset, offset);
    }
}

注意:如果你的topic是多分区,需要调整创建容器时的TopicPartitionOffset配置,以及偏移量比对、存储逻辑,遍历所有分配的分区处理即可。

内容的提问来源于stack exchange,提问作者Sumit Sood

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 12:09:03