Apache Flink AsyncRetryStrategy重试范围问询:全方法还是单API?
Flink AsyncRetryStrategy 重试粒度问题
示例代码
AsyncRetryStrategy asyncRetryStrategy = new AsyncRetryStrategies.FixedDelayRetryStrategyBuilder(3, 100L) // maxAttempts=3, fixedDelay=100ms .ifResult(RetryPredicates.EMPTY_RESULT_PREDICATE) .ifException(RetryPredicates.HAS_EXCEPTION_PREDICATE) .build(); // apply the async I/O transformation with retry DataStream<Tuple2<String, String>> resultStream = AsyncDataStream.unorderedWaitWithRetry(stream, new AsyncDatabaseRequest(), 1000, TimeUnit.MILLISECONDS, 100, asyncRetryStrategy);
此示例来自Flink官方文档。
问题
若在AsyncDatabaseRequest的asyncInvoke方法内发起5次异步API调用,请问Flink的AsyncRetryStrategy中的断言是针对输出流的Tuple2<String, String>类型元素,重试整个asyncInvoke方法,还是针对这5次异步API调用的各自输出,仅重试其中失败的调用?
解答
Flink的AsyncRetryStrategy的重试粒度是整个asyncInvoke方法:
- 当
asyncInvoke方法执行过程中抛出异常,或者其最终输出的Tuple2<String, String>元素触发了重试断言(比如空结果匹配EMPTY_RESULT_PREDICATE),Flink会重新执行整个asyncInvoke方法,也就是重新发起全部5次异步API调用。 - 它不会单独识别并重试5次异步调用中的某一次失败请求,重试判断的对象是
asyncInvoke方法的整体执行结果与状态。
内容的提问来源于stack exchange,提问作者Vijayant
相关产品推荐
相关产品推荐

