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
相关产品推荐
相关产品推荐

