【问题标题】:How to load a compacted topic in memory before starting the context如何在启动上下文之前在内存中加载压缩主题
【发布时间】:2019-09-24 19:53:51
【问题描述】:

我在 kafka 中使用压缩主题,我在应用程序启动时将其加载到 HashMap 中。 然后我正在听一个普通主题的消息,并使用从压缩主题构造的 HashMap 处理它们。

在开始收听其他主题之前,如何确保已压缩的主题已完全读取并且 HashMap 已完全初始化? (RestController 也一样)

【问题讨论】:

    标签: spring spring-boot spring-kafka


    【解决方案1】:

    实现SmartLifecycle 并在start() 中加载地图。确保phase 早于任何其他需要地图的对象。

    【讨论】:

    • 谢谢,但是在 KafkaListener 中,我可以访问分区偏移量吗?为了知道我是否通过读取所有分区直到最后一个偏移量来完成初始化地图?即使消息到达其他侦听器类,KafkaListener 也只会为该特定 bean 触发?
    • 您不应该为此使用KafkaListener;使用原始的Consumerpoll() 直到你没有更多的记录。但是,是的,您可以获得偏移量等;添加@Header参数。
    • KafkaListener 在这种情况下不起作用?还是使用 Consuler.poll() 更好?而且我仍然需要在更新发生时更新地图(新记录),以便消费者保持活力
    • 问题是所有@KafkaListeners都是在同一个阶段启动的; start() 不会阻塞,它只是启动消费者。您还必须侦听空闲事件才能确定何时完成。在poll() 不再返回记录之前轮询消费者更简单(当然,在您寻求开始之后)。您也可以随时使用@KafkaListener,以获取未来的更新,只需确保在关闭初始消费者之前提交偏移量并将侦听器放在同一个消费者组中。
    • 您可以将所有其他侦听器的 autoStartup 属性设置为 false;并在构建地图后通过注册表手动启动它们。
    【解决方案2】:

    这是一个老问题,我知道,但我想提供一个更完整的解决方案代码示例,当我自己遇到这个问题时,我最终得到了这个解决方案。

    这个想法是,就像 Gary 在他自己的答案的 cmets 中提到的那样,在初始化期间使用侦听器不是正确的事情 - 之后会出现。然而,Garry 的SmartLifecycle 想法的替代方案是InitializingBean,我发现实现起来不太复杂,因为它只是一种方法:afterPropertiesSet()

    @Slf4j
    @Configuration
    @RequiredArgsConstructor
    public class MyCacheInitializer implements InitializingBean {
    
        private final ApplicationProperties applicationProperties; // A custom ConfigurationProperties-class
        private final KafkaProperties kafkaProperties;
        private final ConsumerFactory<String, Bytes> consumerFactory;
        private final MyKafkaMessageProcessor messageProcessor;
    
        @Override
        public void afterPropertiesSet() {
            String topicName = applicationProperties.getKafka().getConsumer().get("my-consumer").getTopic();
            Duration pollTimeout = kafkaProperties.getListener().getPollTimeout();
    
            try (Consumer<String, Bytes> consumer = consumerFactory.createConsumer()) {
                consumer.subscribe(List.of(topicName));
    
                log.info("Starting to cache the contents of {}", topicName);
    
                ConsumerRecords<String, Bytes> records;
    
                do {
                    records = consumer.poll(pollTimeout);
                    records.forEach(messageProcessor::process);
                } while (!records.isEmpty());
            }
    
            log.info("Completed caching {}", topicName);
        }
    }
    

    为简洁起见,我使用了 Lombok 的 @Slf4j@RequiredArgsConstructor 注释,但它们可以很容易地替换。 ApplicationProperties 类只是我获取我感兴趣的主题名称的方式。它可以替换为其他内容,但我的实现使用 Lombok 的 @Data 注释,看起来像这样:

    @Data
    @Configuration
    @ConfigurationProperties(prefix = "app")
    public class ApplicationProperties {
    
        private Kafka kafka = new Kafka();
    
        @Data
        public static class Kafka {
            private Map<String, KafkaConsumer> consumer = new HashMap<>();
        }
    
        @Data
        public static class KafkaConsumer {
            private String topic;
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-02-23
      • 2023-03-18
      • 2021-10-29
      • 2011-11-15
      • 1970-01-01
      • 1970-01-01
      • 2012-06-02
      • 1970-01-01
      相关资源
      最近更新 更多