【问题标题】:@KafkaListener get all messages from particular kafka topic@KafkaListener 获取来自特定 kafka 主题的所有消息
【发布时间】:2021-10-26 07:39:03
【问题描述】:

我有一个 @KafkaListener 方法来获取主题中的所有消息,但是对于 @Scheduled 方法起作用的每个间隔时间,我只得到一条消息。如何一次获取主题的所有消息?

这是我的课;

@Slf4j
@Service
public class KafkaConsumerServiceImpl implements KafkaConsumerService {

    @Autowired
    private SimpMessagingTemplate webSocket;

    @Autowired
    private KafkaListenerEndpointRegistry registry;

    @Autowired
    private BrokerProducerService brokerProducerService;

    @Autowired
    private GlobalConfig globalConfig;

    @Override
    @KafkaListener(id = "snapshotOfOutagesId", topics = Constants.KAFKA_TOPIC, groupId = "snapshotOfOutages", autoStartup = "false")
    public void consumeToSnapshot(ConsumerRecord<String, OutageDTO> cr, @Payload String content) {
        log.info("Received content from Kafka notification to notification-snapshot topic: {}", content);
        MessageListenerContainer listenerContainer = registry.getListenerContainer("snapshotOfOutagesId");
        JSONObject jsonObject= new JSONObject(content);
        Map<String, Object> outageMap = jsonToMap(jsonObject);
        brokerProducerService.sendMessage(globalConfig.getTopicProperties().getSnapshotTopicName(),
                outageMap.get("outageId").toString(), toJson(outageMap));
        listenerContainer.stop();
    }

    @Scheduled(initialDelayString = "${scheduler.kafka.snapshot.monitoring}",fixedRateString = "${scheduler.kafka.snapshot.monitoring}")
    private void consumeWithScheduler() {
        MessageListenerContainer listenerContainer = registry.getListenerContainer("snapshotOfOutagesId");
        if (listenerContainer != null){
            listenerContainer.start();
        }
    }

这是我在 application.yml 中的 kafka 属性;

kafka:
  streams:
    common:
      configs:
        "[bootstrap.servers]": 192.168.99.100:9092
        "[client.id]": event
        "[producer.id]": event-producer
        "[max.poll.interval.ms]": 300000
        "[group.max.session.timeout.ms]": 300000
        "[session.timeout.ms]": 200000
        "[auto.commit.interval.ms]": 1000
        "[auto.offset.reset]": latest
        "[group.id]": event-consumer-group
        "[max.poll.records]": 1

还有我的 KafkaConfiguration 类;

    @Bean
    public Map<String, Object> consumerConfigs() {
        Map<String, Object> props = new HashMap<>(globalConfig.getBrokerProperties().getConfigs());
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest");
        return props;
    }

    @Bean
    public ConsumerFactory<String, String> consumerFactory() {
        return new DefaultKafkaConsumerFactory<>(consumerConfigs(), new StringDeserializer(), new StringDeserializer());
    }

    @Bean
    public KafkaListenerContainerFactory<ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

【问题讨论】:

  • 我是否理解正确,您想阅读完整主题一次然后停止消费?

标签: spring-boot apache-kafka kafka-consumer-api apache-kafka-streams spring-kafka


【解决方案1】:

你目前正在做的是:

  1. 创建一个侦听器,但尚未启动它 (autoStartup = false)
  2. 当计划的作业启动时,启动容器(将开始使用来自主题的第一条消息)
  3. 当第一条消息被消费时,您停止容器(导致不再消费任何消息)

因此,您所描述的行为确实不足为奇。

@KafkaListener 不需要计划任务即可开始使用消息。我想你可以删除autoStartup = false并删除预定的作业,之后监听器将一一消费该主题的所有消息,并等待新的消息出现在该主题上。

另外,我注意到的其他一些事情:

这些属性适用于 Kafka Streams,对于常规 Spring Kafka,您需要如下属性:

spring:
  kafka:
    bootstrap-servers: localhost:9092
    consumer:
      auto-offset-reset: earliest
      ...etc

另外:为什么要使用@Payload String content 而不是已经序列化的cr.getVaue()

【讨论】:

  • 感谢您的详细回答@moffeltje。如果我删除 Scheduled annotation,如何让 KafkaListener 在 15 分钟内工作一次?我应该怎么做才能一次收到所有消息而不是一一收到?顺便说一句,我已经在您的消息之前删除了 String 内容,因为我可以获得您使用 cr.getValue() 提到的值。
  • 为什么要每 15 分钟阅读一次所有邮件?对我来说,这听起来不像是一个合适的事件流架构。您可以查看 Spring Kafka 文档中的批处理消费者,但我不确定它是否可以完全满足您的需求。
猜你喜欢
  • 1970-01-01
  • 2018-10-27
  • 1970-01-01
  • 2017-05-29
  • 1970-01-01
  • 1970-01-01
  • 2021-12-10
  • 2021-07-09
  • 2019-11-15
相关资源
最近更新 更多