Apache Flink:asyncInvoke()中如何获取重试次数及处理重试超限?
Flink Async I/O 重试相关问题解答
一、在asyncInvoke()中获取当前重试次数
Flink原生AsyncFunction的asyncInvoke()API没有直接提供重试次数的获取方式,但可以通过自定义AsyncRetryStrategy结合请求上下文传递的方式实现:
- 定义带重试次数的请求包装类
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; } }
- 自定义
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; } }
- 在
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
相关产品推荐
相关产品推荐

