【发布时间】:2021-11-21 06:05:55
【问题描述】:
我有一个 Flink 应用程序,它使用具有多个分区的 Kafka 主题上的传入消息,进行一些处理,然后将它们发送到接收器,该接收器通过 HTTP 将它们发送到外部服务。有时下游服务处于下游,流处理需要停止,直到它恢复运行。
我正在考虑两种方法。
- Http 接收器发送输出消息失败时抛出异常。这将导致任务和作业根据配置的重启策略重启。最终,下游服务将返回,系统将从中断处继续。
- 让 Sink 休眠并重试失败;它可以持续执行此操作,直到下游服务返回为止。
根据我的理解和我的 PoC,有 1。由于接收器本身是外部状态,因此我将失去准确最少一次的保证。据我所知,您不能将简单的 HTTP 端点设置为事务性的,因为它需要实现 TwoPhaseCommitSinkFunction。
使用 2. 这不是一个问题,因为在接收器成功写入之前管道不会继续,我可以依靠整个系统的背压来暂停从 Kafka 源检索消息。
我的主要问题是:
- 您不能为简单的 HTTP 端点创建 TwoPhaseCommitSinkFunction 是否正确?
- 这两种策略中的哪一种最有意义?
- 我是否缺少更简单的明显解决方案?
【问题讨论】:
标签: apache-flink flink-streaming