【发布时间】:2022-11-11 05:06:18
【问题描述】:
我已经为重试完成了 Flink 配置,它正在工作
env.setRestartStrategy(RestartStrategies.failureRateRestart(
3, // number of restart attempts
Time.of(30, TimeUnit.SECONDS),
Time.of(30, TimeUnit.SECONDS) // delay
));
但是我正在使用基于 FlinkKafkaConsumer 的另一种配置来接收消息,我不知道要配置重试。
例如 Spring 有自己的 ErrorHandler,我希望 FlinkKafkaConsumer 和 FlinkKafkaProducer 有类似的东西。
factory.setErrorHandler(new SeekToCurrentErrorHandler(
new DeadLetterPublishingRecoverer(template), 3));
两者都适合,重启策略和FlinkKafka消费者?如果 FlinkKafkaConsumer 可以配置为重试,我可以只使用一个还是应该配置 RestartStrategy?
【问题讨论】:
-
你说的另一个基于 FlinkKafkaConsumer 的配置是什么意思,你能提供一个例子吗?
-
例如 Spring 有自己的 ErrorHandler(我添加到帖子中),我希望 FlinkKafkaConsumer 和 FlinkKafkaProducer 有类似的东西。
标签: java apache-kafka apache-flink