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

Amazon Kinesis生产者单条记录发送失败,调试/多条发送正常

单条发送Amazon Kinesis Data Streams记录失败,但调试/多条发送正常的问题排查与解决

问题描述

仅发送单条测试记录时无法成功写入Kinesis Data Streams,但处于调试状态或发送多条记录时可正常工作。

实现代码

Kinesis Producer配置类

@Bean
public KinesisProducer kinesisProducer() {
    return new KinesisProducer(kinesisProducerConfiguration());
}

@Bean
public KinesisProducerConfiguration kinesisProducerConfiguration() {
    String accessKey = System.getenv("accessKey");
    String secretKey = System.getenv("secretKey");

    BasicAWSCredentials awsCredentials = new BasicAWSCredentials(accessKey, secretKey);

    return new KinesisProducerConfiguration()
            .setCredentialsProvider(new AWSStaticCredentialsProvider(awsCredentials))
            .setVerifyCertificate(false)
            .setRecordMaxBufferedTime(3000)
            .setMaxConnections(1)
            .setRequestTimeout(60000)
            .setRegion(region)
            .setRecordTtl(60000);
}

@Bean
public FutureCallback<UserRecordResult> futureCallback() {
  return new FutureCallback<>() {
        @Override
        public void onFailure(Throwable t) {
            log.info(Constants.KINESIS_FAILED + t.getMessage());
        }

        @Override
        public void onSuccess(UserRecordResult result) {
            System.out.println("here");
            log.info(Constants.KINESIS_SUCCESS + "shard id" + result.getShardId() + "sequence" + result.getSequenceNumber());
        }
    };
}

业务发送代码

ByteBuffer data = ByteBuffer.wrap(new JSONObject(logModel).toString().getBytes("UTF-8"));
var resultListenableFuture = kinesisProducer.addUserRecord(streamName, UUID.randomUUID().toString(), data);
Futures.addCallback(resultListenableFuture, futureCallback, MoreExecutors.directExecutor());

问题原因

核心在于Kinesis Producer Library (KPL)的异步缓冲机制:

  • 配置中的setRecordMaxBufferedTime(3000)指定了记录在缓冲区中最多等待3秒才会被批量发送
  • 单条发送场景下,如果程序在3秒内就终止(比如测试用例执行完毕、主线程退出),KPL后台发送线程还没来得及处理缓冲区中的这条记录,就被强制停止
  • 调试状态下程序不会立刻退出,或者多条发送时缓冲区快速达到批量阈值(默认500条/5MB),因此记录能被正常发送

解决方案

1. 手动触发缓冲区刷新

发送单条记录后调用flushSync(),强制立刻发送缓冲区中的所有记录:

var resultListenableFuture = kinesisProducer.addUserRecord(streamName, UUID.randomUUID().toString(), data);
Futures.addCallback(resultListenableFuture, futureCallback, MoreExecutors.directExecutor());
// 强制刷新缓冲区,同步等待发送完成
kinesisProducer.flushSync();

2. 调整缓冲时间(测试场景专用)

临时缩短缓冲时间,让单条记录更快被触发发送:

return new KinesisProducerConfiguration()
        // ...其他配置
        .setRecordMaxBufferedTime(100) // 缩短至100毫秒
        // ...其他配置

注意:生产环境不建议设置过小,会增加API调用频次,提升成本。

3. 延长程序运行时间

在测试代码中添加等待逻辑,给KPL足够的发送时间:

// 发送记录代码...
Thread.sleep(3500); // 等待时间大于配置的3秒缓冲时间

注意事项

  • KPL是异步批量发送组件,后台线程负责处理发送逻辑,单条发送时必须确保程序不会在发送完成前终止
  • 生产环境持续发送数据时,该问题通常不会出现,后续记录会触发批量发送机制

内容的提问来源于stack exchange,提问作者Erick Jhorman Romero

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 20:58:14