Spring Batch中限制Kafka消息消费数量或速率的方法
解决方案
方案一:限制队列容量为100万,满量暂停消费、空闲恢复
ConcurrentLinkedQueue是无界队列,无法直接限制容量,建议替换为有界阻塞队列,配合Spring Kafka的容器暂停/恢复机制实现流量控制:
代码实现
import org.springframework.kafka.config.KafkaListenerEndpointRegistry; import org.springframework.kafka.listener.MessageListenerContainer; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; import java.util.List; import java.util.concurrent.LinkedBlockingQueue; @Component public class KafkaConsumer { // 有界阻塞队列,容量设为100万 private final LinkedBlockingQueue<Message> queue = new LinkedBlockingQueue<>(1000000); @Autowired private KafkaListenerEndpointRegistry registry; // 启动队列监控线程,空闲时恢复消费 @PostConstruct public void startQueueMonitor() { new Thread(() -> { while (!Thread.currentThread().isInterrupted()) { try { // 当队列剩余容量大于单批次最大消息数(假设为1000)时恢复消费 if (queue.remainingCapacity() > 1000) { MessageListenerContainer container = registry.getListenerContainer("topic-kafka-listener"); if (container != null && !container.isRunning()) { container.start(); } } Thread.sleep(100); // 每隔100ms检查一次队列状态 } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }).start(); } @KafkaListener( topics = "topic", id = "topic-kafka-listener", groupId = "batch-processor", containerFactory = "kafkaListenerContainerFactory" ) public void receive(List<Message> messages) { try { // 批量写入队列,put方法会阻塞直到有空闲空间 for (Message msg : messages) { queue.put(msg); } // 若队列剩余容量不足以容纳下一批次,暂停消费容器 if (queue.remainingCapacity() < messages.size()) { MessageListenerContainer container = registry.getListenerContainer("topic-kafka-listener"); if (container != null && container.isRunning()) { container.pause(); } } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }
关键说明
- 用
LinkedBlockingQueue替代无界队列,直接限制最大容量 - 通过
KafkaListenerEndpointRegistry获取消费容器,实现暂停/恢复 - 监控线程定期检查队列空闲状态,自动恢复消费
方案二:限制消费速率为每秒10万条
使用令牌桶算法精准控制消费速度,推荐用Guava的RateLimiter实现:
代码实现
import com.google.common.util.concurrent.RateLimiter; import org.springframework.stereotype.Component; import java.util.List; import java.util.concurrent.ConcurrentLinkedQueue; @Component public class KafkaConsumer { private final ConcurrentLinkedQueue<Message> queue = new ConcurrentLinkedQueue<>(); // 初始化令牌桶,每秒生成10万令牌(对应每秒消费10万条消息) private final RateLimiter rateLimiter = RateLimiter.create(100000.0); @KafkaListener( topics = "topic", id = "topic-kafka-listener", groupId = "batch-processor", containerFactory = "kafkaListenerContainerFactory" ) public void receive(List<Message> messages) { // 预占对应数量的令牌,若令牌不足则阻塞等待 rateLimiter.acquire(messages.size()); queue.addAll(messages); } }
关键说明
RateLimiter.create(100000.0)表示每秒生成10万令牌,每个令牌对应一条消息acquire(messages.size())会根据当前批次消息数获取令牌,自动控制消费速率- 若需分布式环境下的全局速率控制,需改用Redis等分布式限流方案
内容的提问来源于stack exchange,提问作者PrathamBhatTech
相关产品推荐
相关产品推荐

