【问题标题】:Apache Flink - how to stop and resume stream processing on downstream failureApache Flink - 如何在下游故障时停止和恢复流处理
【发布时间】:2021-11-21 06:05:55
【问题描述】:

我有一个 Flink 应用程序,它使用具有多个分区的 Kafka 主题上的传入消息,进行一些处理,然后将它们发送到接收器,该接收器通过 HTTP 将它们发送到外部服务。有时下游服务处于下游,流处理需要停止,直到它恢复运行。

我正在考虑两种方法。

  1. Http 接收器发送输出消息失败时抛出异常。这将导致任务和作业根据配置的重启策略重启。最终,下游服务将返回,系统将从中断处继续。
  2. 让 Sink 休眠并重试失败;它可以持续执行此操作,直到下游服务返回为止。

根据我的理解和我的 PoC,有 1。由于接收器本身是外部状态,因此我将失去准确最少一次的保证。据我所知,您不能将简单的 HTTP 端点设置为事务性的,因为它需要实现 TwoPhaseCommitSinkFunction。

使用 2. 这不是一个问题,因为在接收器成功写入之前管道不会继续,我可以依靠整个系统的背压来暂停从 Kafka 源检索消息。

我的主要问题是:

  1. 您不能为简单的 HTTP 端点创建 TwoPhaseCommitSinkFunction 是否正确?
  2. 这两种策略中的哪一种最有意义?
  3. 我是否缺少更简单的明显解决方案?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    我认为你可以在 Flink 中尝试 AsyncIO - https://nightlies.apache.org/flink/flink-docs-master/docs/dev/datastream/operators/asyncio/

    在请求的所有操作完成后,尝试让 HTTP 端点发送响应,例如在 http 服务器中,请求的过程已经完成,结果已经提交给 DB。然后在 AsyncIO 运算符中使用 http 异步客户端。 AsyncIO 操作员将等待,直到操作员收到响应。如果发生任何错误,Flink 流式管道将失败并根据恢复策略重新启动管道。

    所有未收到响应的 HTTP 端点请求都将在 AsyncIO 操作符的内部缓冲区中,一旦流式传输管道失败,缓冲区中未决的请求将保存在检查点状态。内部缓冲区满时也会触发背压。

    【讨论】:

    • 我很想用水槽来做这件事,但这似乎是个好方法,谢谢
    • 已接受此答案。我发现的唯一可接受的替代方法是使用接收器的内部状态管理 HTTP 请求队列,随着时间的推移使用线程池发送它们,但异步操作符似乎等效但已经准备好。
    • 这是异步操作符:)
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-08-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-10-26
    相关资源
    最近更新 更多