【问题标题】:Controlled/manual error/recovery handling in stream-based applications基于流的应用程序中的受控/手动错误/恢复处理
【发布时间】:2019-06-23 23:46:32
【问题描述】:

我正在开发一个基于Apache Flink 的应用程序,它使用Apache Kafka 进行输入和输出。可能这个应用程序会被移植到Apache Spark,所以我也把它添加为标签,问题还是一样。

我的要求是所有通过 kafka 接收的传入消息必须按顺序处理,并安全地存储在持久层(数据库)中,并且不得丢失任何消息。

此应用程序中的流式传输部分相当琐碎/小,因为主要逻辑将归结为:

environment.addSource(consumer)    // 1) DataStream[Option[Elem]]
  .filter(_.isDefined)             // 2) discard unparsable messages
  .map(_.get)                      // 3) unwrap Option
  .map(InputEvent.fromXml(_))      // 4) convert from XML to internal representation
  .keyBy(_.id)                     // 5) assure in-order processing on logical-key level
  .map(new DBFunction)             // 6) database lookup, store of update and additional enrichment
  .map(InputEvent.toXml(_))        // 7) convert back to XML
  .addSink(producer)               // 8) attach kafka producer sink

现在,在此管道期间,可能会发生几种错误情况

  • 数据库变得不可用(关闭、表空间已满,...)
  • 由于逻辑错误(来自列格式)而无法存储更改
  • 由于代理不可用,kafka 生产者无法发送消息

可能还有其他情况。

现在我的问题是,如何在这些情况下按照上述方式确保一致性,而实际上我必须执行以下操作:

  1. Stream-Operator 6) 检测到问题(数据库不可用)
  2. DBFunction 对象的 DB 连接必须恢复,这可能需要几分钟后才能成功
  3. 这意味着必须暂停整体处理,最好是暂停整个管道,以便将传入消息大量加载到内存中
  4. 恢复数据库后继续处理。处理必须与在 1) 处遇到问题的消息完全一致

现在我知道至少有 2 个关于故障处理的工具:

  1. kafka 消费者偏移量
  2. apache flink 检查点

但是,在搜索文档时,我看不到如何在单个运算符中的流处理中间使用其中任何一个。

那么,对于流式应用程序中的细粒度错误处理和恢复,推荐的策略是什么?

【问题讨论】:

    标签: apache-spark error-handling apache-kafka stream apache-flink


    【解决方案1】:

    几点:

    keyBy 不会帮助确保按顺序处理。如果有的话,它可能会交错来自不同 Kafka 分区的事件(可能在每个分区中是有序的),从而在以前不存在的地方产生无序。如果不了解您打算使用多少个 FlinkKafkaConsumer 实例、每个将使用多少个分区、密钥如何在 Kafka 分区中分布以及您为什么这么认为,很难更具体地评论您如何保证按顺序处理一个 keyBy 是必要的——但如果你设置正确,保持顺序可能是可以实现的。 reinterpretAsKeyedStream 在这里可能会有所帮助,但此功能难以理解,并且难以正确使用。

    您可以使用 Flink 的 AsyncFunction 以容错、精确一次的方式管理与外部数据库的连接。

    Flink 不支持以系统方式进行细粒度恢复——它的检查点是整个分布式集群状态的全局快照,旨在在恢复期间作为一个整体的、自洽的快照使用。如果您的工作失败,通常唯一的办法是从检查点重新启动,这将涉及倒回输入队列(到存储在检查点中的偏移量),重放自这些偏移量以来的事件,重新发出数据库查找(异步功能将自动执行),并使用 kafka 事务来实现端到端的恰好一次语义。但是,在令人尴尬的并行作业的情况下,有时可以利用fine-grained recovery

    【讨论】:

    • 据我了解,keyBy 导致始终由相同的运算符实例处理相同的密钥,无论并行度如何,因此确保了运算符中的有序处理-范围。当然,keyBy 将永远无法将通过 kafka out of order 提供的东西有序化。我这方面的要求很简单:必须按顺序处理具有相同逻辑键的消息,如何设置 kafka 主题、消费者和 flink-application 由我决定。目前,恰好有 1 个FlinkKafkaConsumer011 实例。
    • ... 和该主题的 1 个分区。关于恢复:当然,该应用程序将是一个 24/7 的在线应用程序,并且很可能在任何时候都没有人工支持。因此,仅依靠在失败的情况下从检查点重新启动(= 手动任务)并不是一种可行的方法。我宁愿尝试实现一个在所有情况下都稳定的实现,它不需要用户交互,或者在最坏的情况下,如果错误无法恢复,就会让一切停止。这就是我的问题所涉及的,如何自动停止一切。
    • 更准确地说,我对 flink 在以下情况下会发生什么缺乏了解:1)操作员线程因陷入循环而被 阻塞尝试恢复数据库连接 2) 操作线程被阻塞,因为检测到一些不可恢复的错误,必须停止处理。这个话题有点复杂,因为我所有的操作员都是无国籍的,但我很可能会遵循 Fabian Hueske 的建议,如下所示:stackoverflow.com/questions/54986886/… 并实施交易。
    • 从检查点重新启动是一个全自动过程,您可以对自动重新启动策略进行相当程度的控制。 ci.apache.org/projects/flink/flink-docs-release-1.7/dev/… 无需人工干预。
    • 如果你确保对于任何给定的 key,它的所有事件都在同一个 kafka 分区中,那么它们自然会保持有序,因为它们都会被同一个 FlinkKafkaConsumer011 实例消费。在这种情况下,keyBy 不会将它们与来自其他分区的相同键的事件一起洗牌(这会导致无序)。
    猜你喜欢
    • 2012-12-14
    • 2017-05-14
    • 1970-01-01
    • 2021-11-26
    • 2011-04-24
    • 1970-01-01
    • 1970-01-01
    • 2014-02-20
    • 1970-01-01
    相关资源
    最近更新 更多