如何在Spring应用无SQS请求时批量发送记录到Kinesis Stream?
解决方案:实现基于数量和超时的Kinesis批量发送
针对你的需求,我们可以通过本地缓存队列+原子计数器+定时任务的组合,同时实现「累积指定数量触发批量发送」和「超时自动触发发送」的逻辑,具体实现如下:
修改后的KPLService代码
@Service public class KPLService { private Logger log = LoggerFactory.getLogger(this.getClass()); @Autowired public KinesisProducer producer; // 线程安全的缓存队列,存储待发送的记录 private final ConcurrentLinkedQueue<MappedSqsRecord> recordQueue = new ConcurrentLinkedQueue<>(); // 原子计数器,记录当前缓存的记录数 private final AtomicInteger recordCount = new AtomicInteger(0); // 批量触发阈值(可通过配置动态调整) private final int BATCH_SIZE = 500; // 超时触发间隔(单位:秒,可通过配置动态调整) private final int FLUSH_INTERVAL_SECONDS = 10; private ScheduledExecutorService scheduler; @PostConstruct public void initBatchScheduler() { // 初始化单线程定时任务,定期检查并发送缓存的记录 scheduler = Executors.newSingleThreadScheduledExecutor(); scheduler.scheduleAtFixedRate( this::flushPendingRecords, FLUSH_INTERVAL_SECONDS, FLUSH_INTERVAL_SECONDS, TimeUnit.SECONDS ); } public void addRecord(MappedSqsRecord record) { recordQueue.add(record); // 每添加一条记录就计数,达到阈值则触发批量发送 if (recordCount.incrementAndGet() >= BATCH_SIZE) { flushPendingRecords(); } } /** * 批量发送缓存中的所有记录 */ private void flushPendingRecords() { List<MappedSqsRecord> batch = new ArrayList<>(); MappedSqsRecord record; // 一次性取出队列中所有记录 while ((record = recordQueue.poll()) != null) { batch.add(record); } if (!batch.isEmpty()) { try { // 批量添加到KPL for (MappedSqsRecord r : batch) { producer.addUserRecord( "kinesis-stream-name", r.uuid, ByteBuffer.wrap(new ObjectMapper().writeValueAsBytes(r)) ); } // 强制刷新,确保记录发送到Kinesis producer.flushSync(); log.info("成功批量发送{}条记录到Kinesis Stream", batch.size()); } catch (JsonProcessingException e) { log.error("序列化记录失败", e); // 可选:将失败的记录存入死信队列或重试队列,避免数据丢失 } finally { // 重置计数器 recordCount.set(0); } } } @PreDestroy public void onShutdown() { // 关闭定时任务 scheduler.shutdown(); try { if (!scheduler.awaitTermination(5, TimeUnit.SECONDS)) { scheduler.shutdownNow(); } } catch (InterruptedException e) { scheduler.shutdownNow(); } // 应用关闭前最后一次发送缓存的记录 flushPendingRecords(); } }
关键逻辑说明
- 线程安全保障:使用
ConcurrentLinkedQueue和AtomicInteger处理多线程环境下的缓存和计数(SQS Listener默认是多线程调用) - 双触发机制:
- 数量触发:每添加一条记录就计数,达到
BATCH_SIZE时立即触发批量发送 - 超时触发:通过
ScheduledExecutorService每隔固定时间检查缓存队列,若有未发送的记录则立即发送
- 数量触发:每添加一条记录就计数,达到
- 优雅关闭:在应用关闭时(即使是Fargate重启/销毁),先关闭定时任务,再发送剩余缓存的记录,避免数据丢失
优化建议
- 可将
BATCH_SIZE和FLUSH_INTERVAL_SECONDS通过@Value注解从配置文件读取,方便动态调整:@Value("${kinesis.batch.size:500}") private int BATCH_SIZE; @Value("${kinesis.flush.interval:10}") private int FLUSH_INTERVAL_SECONDS; - 增加失败重试机制:如果Kinesis发送失败,可将记录存入SQS死信队列或本地重试队列,避免数据丢失
- 结合KPL原生配置:KPL本身支持通过
KinesisProducerConfiguration设置batchSize和lingerMs,但上述方案更灵活可控,可根据业务需求选择组合使用
内容的提问来源于stack exchange,提问作者Rue Honda
相关产品推荐
相关产品推荐

