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

如何在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();
    }
}

关键逻辑说明

  1. 线程安全保障:使用ConcurrentLinkedQueue和AtomicInteger处理多线程环境下的缓存和计数(SQS Listener默认是多线程调用)
  2. 双触发机制:
    • 数量触发:每添加一条记录就计数,达到BATCH_SIZE时立即触发批量发送
    • 超时触发:通过ScheduledExecutorService每隔固定时间检查缓存队列,若有未发送的记录则立即发送
  3. 优雅关闭:在应用关闭时(即使是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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 22:42:29