【问题标题】:Flink and Kafka configuration for retries重试的 Flink 和 Kafka 配置
【发布时间】: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


【解决方案1】:

@sergiopf 重启策略如何工作?当执行抛出任何异常时,它是否会使用相同的数据重试相同的操作?

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2016-06-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-25
    • 1970-01-01
    • 2017-01-20
    • 1970-01-01
    • 2016-11-29
    相关资源
    最近更新 更多