You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 22:53:09