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
相关产品推荐
相关产品推荐

