Spring Kafka中Idle Container的具体含义及触发条件咨询
Spring Kafka ListenerContainerIdleEvent 触发逻辑说明
首先明确:ListenerContainerIdleEvent 触发的「Idle」状态,指的是当前容器对应的消费者在设定的 idleEventInterval 时长内,没有从它分配的分区中拉取到任何可消费的消息,并非整个主题没有消息产生。
具体细节:
- 你的代码中,每个
ConcurrentMessageListenerContainer是通过createContainer(new TopicPartitionOffset(topic, partition))创建的,这意味着该容器仅负责消费指定的单个分区。当这个容器的消费者在设定的空闲周期内,对该分区轮询后没有获取到新消息,就会触发空闲事件。 - 即使主题的其他分区还有消息在产生或消费,只要当前容器负责的分区没有可拉取的新消息,就会触发
@EventListener标注的idle方法,和主题整体状态无关。 - 空闲事件的判断依据是消费者的轮询结果:每次轮询如果没有拿到记录,且连续的轮询空结果累计时长达到
idleEventInterval,就会触发事件。
你的代码示例:
public ConcurrentMessageListenerContainer<String, String> newContainer(String topic, int partition, int idlePeriod) { this.factory.setConsumerFactory(consumerFactory); this.factory.getContainerProperties().setIdleEventInterval(idlePeriod * 1000L); this.factory.getContainerProperties().setAckMode(AckMode.MANUAL_IMMEDIATE); ConcurrentMessageListenerContainer<String, String> container = this.factory .createContainer(new TopicPartitionOffset(topic, partition)); container.setupMessageListener((AcknowledgingMessageListener<String, String>) (record, acknowledgment) -> { kafkaService.proccessorConsumer(record); acknowledgment.acknowledge(); }); this.containers.put(service_name+"test", container); container.start(); return container; } @EventListener public void idle(ListenerContainerIdleEvent event) { log.warn("Idle period has been captured", event); kafkaService.processIdelConsumer(); }
内容的提问来源于stack exchange,提问作者kikicoder
相关产品推荐
相关产品推荐

