【发布时间】:2019-08-24 05:56:59
【问题描述】:
我正在使用带有 spring-kafka 的 Spring Boot 2.x(不是 spring-integration-kafka)
我有多个用 @KafkaListener 注释的 bean ...每个都消耗一个主题...所以既然我有 12 个主题,那么我还需要有 12 个 KafkaConsumers bean ...我想知道我是否可以以编程方式/动态创建这些 bean ...也许使用 KafkaListenerEndpointRegistry 来动态创建消费者容器。
注意:我需要批量消费消息...所以也许我可以使用 BatchMessageListener?
当前代码:
@KafkaListener(
id = COUNTRY,
containerFactory = KAFKA_LISTENER_FACTORY_BEAN_NAME,
topics = {TOPIC},
groupId = GROUP_ID,
clientIdPrefix = CLIENT_ID,
errorHandler = VALIDATION_ERROR_HANDLER_BEAN_NAME
)
@Override
public void consume(final List<MessageDTO> messages,
@Header(KafkaHeaders.RECEIVED_TOPIC) final List<String> topics,
@Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) final List<String> messagesKey,
@Header(KafkaHeaders.RECEIVED_PARTITION_ID) final List<Integer> partitionIds,
@Header(KafkaHeaders.RECEIVED_TIMESTAMP) final List<Long> timestamps,
@Header(KafkaHeaders.OFFSET) final List<Long> offsets) {
(...)
}
每个主题消费者都有自己的实现,具体取决于主题。请你们指点我的博客/伪代码/git线程/答案吗?
【问题讨论】:
-
每个主题都在同一个集群中?那么有效载荷呢?它们有什么不同吗?
-
有效载荷结构是相同的......每个国家都有一个主题。然后每个主题至少有一个消费者,因为实施取决于国家/地区
标签: java spring spring-boot kafka-consumer-api spring-kafka