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

如何仅记录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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 09:43:19