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

Apache Flink:asyncInvoke()中如何获取重试次数及处理重试超限?

一、在asyncInvoke()中获取当前重试次数

Flink原生AsyncFunction的asyncInvoke()API没有直接提供重试次数的获取方式,但可以通过自定义AsyncRetryStrategy结合请求上下文传递的方式实现:

  1. 定义带重试次数的请求包装类
public class RetryAwareRequest<T> {
    private final T originalRequest;
    private final int retryCount;

    public RetryAwareRequest(T originalRequest, int retryCount) {
        this.originalRequest = originalRequest;
        this.retryCount = retryCount;
    }

    // Getters
    public T getOriginalRequest() { return originalRequest; }
    public int getRetryCount() { return retryCount; }
}
  1. 自定义AsyncRetryStrategy传递重试次数
    在重试逻辑中,将当前重试次数(attemptNumber参数,从1开始计数)绑定到请求上下文:
public class CustomAsyncRetryStrategy<T> implements AsyncRetryStrategy<T> {
    private final int maxRetries;

    public CustomAsyncRetryStrategy(int maxRetries) {
        this.maxRetries = maxRetries;
    }

    @Override
    public boolean canRetry(T input, Throwable throwable) {
        // 根据业务规则判断是否需要重试(比如异常类型)
        return throwable instanceof IOException;
    }

    @Override
    public CompletableFuture<T> retry(T input, Throwable throwable, int attemptNumber) {
        // 包装原始请求,携带当前重试次数
        RetryAwareRequest<T> wrappedRequest = new RetryAwareRequest<>((T) input, attemptNumber);
        return CompletableFuture.completedFuture((T) wrappedRequest);
    }

    @Override
    public int getMaxRetryAttempts() {
        return maxRetries;
    }
}
  1. 在asyncInvoke()中解析重试次数
    在自定义AsyncFunction中,判断请求是否为包装类,取出重试次数:
public class MyAsyncProcessor extends AsyncFunction<OriginalRequest, Result> {
    @Override
    public void asyncInvoke(OriginalRequest input, ResultFuture<Result> resultFuture) throws Exception {
        int currentRetryCount = 0;
        OriginalRequest actualRequest = input;

        // 解析重试次数
        if (input instanceof RetryAwareRequest) {
            RetryAwareRequest<OriginalRequest> wrappedReq = (RetryAwareRequest<OriginalRequest>) input;
            actualRequest = wrappedReq.getOriginalRequest();
            currentRetryCount = wrappedReq.getRetryCount();
        }

        // 使用重试次数做自定义逻辑(比如日志、降级处理)
        System.out.printf("Processing request, retry count: %d%n", currentRetryCount);

        // 执行异步业务调用
        asyncService.call(actualRequest)
                .thenAccept(result -> resultFuture.complete(Collections.singleton(result)))
                .exceptionally(e -> {
                    resultFuture.completeExceptionally(e);
                    return null;
                });
    }
}

二、重试超限后的处理API及自定义逻辑

当重试次数超过设定上限时,Flink会触发AsyncRetryStrategy中的onRetryFailure()方法(Flink 1.13+版本支持),你可以重写该方法接入自定义逻辑:

重写onRetryFailure()实现自定义失败处理

@Override
public void onRetryFailure(T input, Throwable throwable, int attemptNumber) {
    // attemptNumber为最后一次重试的次数(等于maxRetries)
    // 自定义逻辑示例:记录错误日志、上报监控、保存失败请求到外部存储
    System.err.printf("Retry exhausted after %d attempts. Request: %s, Error: %s%n", 
                      attemptNumber, input.toString(), throwable.getMessage());
    // alertService.sendAlert("Async I/O Retry Failed", input, throwable);
    // failureStorage.saveFailedRequest(input, throwable);
}

低版本Flink兼容方案(无onRetryFailure())

如果使用的Flink版本低于1.13,可以在canRetry()方法中判断重试次数上限,执行自定义逻辑后返回false终止重试:

@Override
public boolean canRetry(T input, Throwable throwable) {
    int currentRetryCount = getCurrentRetryCount(input); // 从请求上下文解析当前重试次数
    if (currentRetryCount >= maxRetries) {
        // 执行自定义失败逻辑
        handleRetryExhausted(input, throwable, currentRetryCount);
        return false;
    }
    return throwable instanceof IOException;
}

private void handleRetryExhausted(T input, Throwable throwable, int retryCount) {
    // 自定义失败处理逻辑
}

内容的提问来源于stack exchange,提问作者Chuni Lal Kukreja

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 19:10:29