【发布时间】:2019-09-24 19:53:51
【问题描述】:
我在 kafka 中使用压缩主题,我在应用程序启动时将其加载到 HashMap 中。 然后我正在听一个普通主题的消息,并使用从压缩主题构造的 HashMap 处理它们。
在开始收听其他主题之前,如何确保已压缩的主题已完全读取并且 HashMap 已完全初始化? (RestController 也一样)
【问题讨论】:
标签: spring spring-boot spring-kafka
我在 kafka 中使用压缩主题,我在应用程序启动时将其加载到 HashMap 中。 然后我正在听一个普通主题的消息,并使用从压缩主题构造的 HashMap 处理它们。
在开始收听其他主题之前,如何确保已压缩的主题已完全读取并且 HashMap 已完全初始化? (RestController 也一样)
【问题讨论】:
标签: spring spring-boot spring-kafka
实现SmartLifecycle 并在start() 中加载地图。确保phase 早于任何其他需要地图的对象。
【讨论】:
KafkaListener;使用原始的Consumer 和poll() 直到你没有更多的记录。但是,是的,您可以获得偏移量等;添加@Header参数。
KafkaListener 在这种情况下不起作用?还是使用 Consuler.poll() 更好?而且我仍然需要在更新发生时更新地图(新记录),以便消费者保持活力
@KafkaListeners都是在同一个阶段启动的; start() 不会阻塞,它只是启动消费者。您还必须侦听空闲事件才能确定何时完成。在poll() 不再返回记录之前轮询消费者更简单(当然,在您寻求开始之后)。您也可以随时使用@KafkaListener,以获取未来的更新,只需确保在关闭初始消费者之前提交偏移量并将侦听器放在同一个消费者组中。
这是一个老问题,我知道,但我想提供一个更完整的解决方案代码示例,当我自己遇到这个问题时,我最终得到了这个解决方案。
这个想法是,就像 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;
}
}
【讨论】: