Spring Boot下如何验证Kafka Topic是否非空(含至少1条消息)
判断Kafka Topic是否存在消息的实现方案(Spring Boot + Kafka)
方法一:通过AdminClient查询分区偏移量(批量检测)
这种方法适合批量检测已有Topic是否包含消息,核心逻辑是对比每个分区的起始偏移量(beginningOffset)和当前末尾偏移量(endOffset):如果两者差值大于0,说明该分区存在消息;只要Topic下有一个分区有消息,即可判定该Topic收到过消息。
修改你现有代码的实现思路:
- 遍历每个Topic的所有分区
- 调用AdminClient的
listOffsets方法获取分区的起始与末尾偏移量 - 计算偏移量差值,判断Topic是否存在消息
修改后的示例代码:
@ReadOperation public List<TopicManifest> kafkaTopic() throws ExecutionException, InterruptedException { ListTopicsOptions listTopicsOptions = new ListTopicsOptions(); listTopicsOptions.listInternal(true); ListTopicsResult listTopicsResult = adminClient.listTopics(listTopicsOptions); Set<String> topics = listTopicsResult.names().get().stream() .filter(topic -> !topic.startsWith("_")) .collect(Collectors.toSet()); DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(topics); Map<String, KafkaFuture<TopicDescription>> topicNameValues = describeTopicsResult.topicNameValues(); return topicNameValues.entrySet().stream().map(entry -> { try { TopicDescription topicDescription = entry.getValue().get(); String topicName = entry.getKey(); int partitionCount = topicDescription.partitions().size(); boolean hasMessages = false; // 遍历每个分区,检查偏移量 for (TopicPartitionInfo partitionInfo : topicDescription.partitions()) { TopicPartition tp = new TopicPartition(topicName, partitionInfo.partition()); // 获取起始偏移量 long earliestOffset = adminClient.listOffsets(Collections.singletonMap(tp, OffsetSpec.earliest())) .all().get().get(tp); // 获取末尾偏移量 long latestOffset = adminClient.listOffsets(Collections.singletonMap(tp, OffsetSpec.latest())) .all().get().get(tp); if (latestOffset > earliestOffset) { hasMessages = true; break; // 只要有一个分区有消息,直接判定Topic有消息 } } return TopicManifest.builder() .name(topicName) .noOfPartitions(partitionCount) .hasMessages(hasMessages) // 新增字段标识是否有消息 .build(); } catch (InterruptedException | ExecutionException e) { e.printStackTrace(); } return null; }).filter(Objects::nonNull) // 过滤异常导致的null结果 .collect(Collectors.toList()); }
注意事项:
- 需要给
TopicManifest新增hasMessages布尔类型字段 - 该方法属于非实时检测,适合定时巡检场景
- 频繁调用会增加Kafka集群负载,建议合理控制检测频率
方法二:实时监听消息(首次消息触发)
如果需要实时感知第一条消息到达的事件,最直接的方式是为目标Topic创建消费者监听,当首次收到消息时标记该Topic为已收到消息。
方式1:使用Spring Kafka的@KafkaListener
@Component public class TopicMessageListener { // 用ConcurrentHashMap维护Topic的首次消息状态,保证线程安全 private final Map<String, Boolean> topicHasMessageStatus = new ConcurrentHashMap<>(); // 通过正则匹配监听所有非内部Topic @KafkaListener(topicPattern = "^(?!_).*$", groupId = "topic-monitor-group") public void listen(ConsumerRecord<String, Object> record) { String topic = record.topic(); // 仅在首次收到消息时更新状态 topicHasMessageStatus.putIfAbsent(topic, true); } // 供指标导出器调用的状态查询方法 public boolean hasTopicReceivedMessages(String topic) { return Boolean.TRUE.equals(topicHasMessageStatus.get(topic)); } }
方式2:手动创建消费者(更高灵活性)
如果需要更精细的控制逻辑,可以手动创建消费者,订阅目标Topic并轮询消息:
@Component public class TopicMessageMonitor { private final Map<String, Boolean> topicStatus = new ConcurrentHashMap<>(); private final Consumer<String, Object> consumer; public TopicMessageMonitor(ConsumerFactory<String, Object> consumerFactory) { this.consumer = consumerFactory.createConsumer("topic-monitor-group"); // 订阅所有非内部Topic consumer.subscribe(Pattern.compile("^(?!_).*$")); // 启动后台线程轮询消息 new Thread(this::pollMessages).start(); } private void pollMessages() { while (!Thread.currentThread().isInterrupted()) { ConsumerRecords<String, Object> records = consumer.poll(Duration.ofSeconds(1)); if (!records.isEmpty()) { records.forEach(record -> { topicStatus.putIfAbsent(record.topic(), true); }); } } consumer.close(); } public boolean isTopicHasMessages(String topic) { return Boolean.TRUE.equals(topicStatus.get(topic)); } }
注意事项:
- 消费者组ID需唯一,避免影响其他业务消费流程
- 使用
topicPattern可自动订阅新增的非内部Topic - 可将
topicStatus的状态直接暴露为自定义指标,供导出器使用
内容的提问来源于stack exchange,提问作者Hari Rao
相关产品推荐
相关产品推荐

