【发布时间】: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 生产者无法发送消息
可能还有其他情况。
现在我的问题是,如何在这些情况下按照上述方式确保一致性,而实际上我必须执行以下操作:
- Stream-Operator 6) 检测到问题(数据库不可用)
-
DBFunction对象的 DB 连接必须恢复,这可能需要几分钟后才能成功 - 这意味着必须暂停整体处理,最好是暂停整个管道,以便将传入消息大量加载到内存中
- 恢复数据库后继续处理。处理必须与在 1) 处遇到问题的消息完全一致
现在我知道至少有 2 个关于故障处理的工具:
- kafka 消费者偏移量
- apache flink 检查点
但是,在搜索文档时,我看不到如何在单个运算符中的流处理中间使用其中任何一个。
那么,对于流式应用程序中的细粒度错误处理和恢复,推荐的策略是什么?
【问题讨论】:
标签: apache-spark error-handling apache-kafka stream apache-flink