如何仅记录AWS Kinesis最后重试失败的Put请求及保障数据投递
我有一个AWS Kinesis数据流分片,每秒最多支持1000条记录。多个应用的平均发送量远低于该阈值,但偶尔会出现持续1-2秒、每秒4000条记录的突发流量。
我采用AWS SDK v1配置客户端自动重试错误:
ClientConfiguration clientConfiguration = PredefinedClientConfigurations.defaultConfig(); clientConfiguration.setRetryMode(RetryMode.STANDARD);
该配置最多重试3次,且每次重试间隔递增,但我不清楚具体间隔时长。此外,我为putRecordAsync方法传入AsyncHandler,通过其onError方法记录错误。
目前该方法会在首次请求因ProvisionedThroughputExceededException失败时就记录错误,我希望仅在最后一次重试失败时才记录,以此确认数据确实未发送丢失。
请问如何实现这一需求?另外,在平均发送量低于1000条/秒的前提下,还有哪些方法可确保所有请求都能送达数据流?
一、仅在最后一次重试失败时记录错误
首先明确AWS SDK v1的STANDARD重试模式参数:默认最多重试3次(即总共发起4次请求:1次初始请求+3次重试),采用指数退避策略,初始间隔100ms,每次重试间隔翻倍(100ms → 200ms → 400ms)。
要实现仅最后一次重试失败才记录错误,可通过以下两种方式:
方式1:利用AmazonClientException的重试信息判断
在自定义AsyncHandler的onError方法中,通过异常的RetryInfo获取已重试次数,当次数达到最大重试上限(3次)时再记录错误:
AsyncHandler<PutRecordRequest, PutRecordResult> customHandler = new AsyncHandler<>() { @Override public void onSuccess(PutRecordRequest request, PutRecordResult result) { // 处理成功逻辑 } @Override public void onError(Exception exception) { if (exception instanceof AmazonClientException ace) { RetryInfo retryInfo = ace.getRetryInfo(); // 确认已耗尽所有重试次数 if (retryInfo != null && retryInfo.getRetryCount() == 3) { // 仅此时记录最终失败错误 log.error("所有重试已耗尽,数据发送失败", exception); } } else { // 非AWS客户端异常直接记录 log.error("发送请求发生未知错误", exception); } } };
方式2:自定义重试策略+监听重试事件
先明确配置自定义重试策略,确保重试次数可控,再通过RetryListener跟踪重试次数,在最后一次失败时记录:
// 1. 自定义重试策略,明确最大重试次数 RetryPolicy retryPolicy = new RetryPolicy( PredefinedRetryPolicies.DEFAULT_RETRY_CONDITION, PredefinedRetryPolicies.DEFAULT_BACKOFF_STRATEGY, 3, // 最大重试次数 true // 允许重试超时请求 ); // 2. 自定义重试监听器跟踪重试次数 class RetryCountListener extends RetryListenerAdapter { private int currentRetryCount = 0; @Override public void onRetry(RetryPolicy.RetryContext context) { currentRetryCount++; } public int getCurrentRetryCount() { return currentRetryCount; } } // 3. 配置客户端 RetryCountListener retryListener = new RetryCountListener(); ClientConfiguration clientConfig = PredefinedClientConfigurations.defaultConfig(); clientConfig.setRetryPolicy(retryPolicy); clientConfig.addRetryListener(retryListener); AmazonKinesisAsync kinesisClient = AmazonKinesisAsyncClientBuilder.standard() .withClientConfiguration(clientConfig) .build(); // 4. 在AsyncHandler中判断重试次数 AsyncHandler<PutRecordRequest, PutRecordResult> handler = new AsyncHandler<>() { @Override public void onSuccess(PutRecordRequest request, PutRecordResult result) { retryListener.currentRetryCount = 0; // 重置计数 } @Override public void onError(Exception exception) { if (retryListener.getCurrentRetryCount() == 3) { log.error("最终重试失败,数据未送达", exception); } retryListener.currentRetryCount = 0; // 重置计数 } };
二、确保所有请求送达Kinesis的其他方法
在平均流量低于阈值的前提下,可通过以下手段应对突发流量,保障数据送达:
使用批量写入接口:替换
putRecordAsync为putRecordsAsync,单次请求最多可发送500条记录(或5MB数据),减少请求频次,降低触发ProvisionedThroughputExceededException的概率,同时提升吞吐量利用率。切换到Kinesis按需模式:按需模式会自动根据流量调整分片吞吐量,无需手动配置分片数量,按实际使用量付费,完美适配低平均流量+高突发流量的场景,彻底避免吞吐量不足的问题。
本地流量整形:在应用端实现本地缓存队列(如
LinkedBlockingQueue),当突发流量到来时,先将请求缓存到本地,再以不超过1000条/秒的速率异步发送到Kinesis,削峰填谷,避免直接打满分片吞吐量。引入中间缓冲层:将应用请求先发送到Amazon SQS队列,再用AWS Lambda触发消费SQS消息,以Kinesis的吞吐量限制为基准控制发送速率,SQS作为缓冲层承接突发流量,确保数据不会丢失。
动态调整分片数量:如果突发流量频繁且可预测,可通过Kinesis API或控制台临时增加分片数量,待流量回落再减少分片,不过该方式需要手动或自动化调整,适合有规律的突发场景。
内容的提问来源于stack exchange,提问作者T3rm1

