【问题标题】:Join a static and a dynamic Kafka source in Flink在 Flink 中加入静态和动态 Kafka 源
【发布时间】:2020-03-20 02:27:06
【问题描述】:

今天,我想谈一个关于 Flink 的概念性话题,而不是技术性话题。

在我们的例子中,我们确实有两个需要连接的 Kafka 主题 A 和 B。连接应该总是包括来自主题 A 的 所有 元素,以及来自主题 B 的所有新元素。实现这一点有两种可能性:总是创建一个新的消费者并开始消费主题 A 的所有元素,或者将主题 A 中的所有元素保持在一个状态内,一旦被消费。 目前,技术方法是通过加入两个 DataStreams,这很快向我们展示了它在这个用例中的局限性,因为没有窗口就不可能加入流(很公平)。来自主题 A 的元素最终会丢失,如果窗口继续移动并且我感觉定期重置消费者会绕过 Flink 引入的复杂逻辑。

我现在正在寻找的另一种方法是使用 Table API,听起来它最适合这项工作,实际上可以无限期地保持所有元素的状态。

但是我的问题:在深入研究 Table API 之前,我只想知道有一种更优雅的方式,我想确定这是否是解决此问题的最佳解决方案,或者是否有更合适的 Flink我不知道的概念?

编辑:我忘了说:我们不使用 POJO,而是保持它的通用性,这意味着传入的数据被标识为Tuple2<K,V>,其中K,V 是每个GenericRecord 的实例。序列化/反序列化的相应模式是在运行时从模式注册表中获得的。我不知道,在这种情况下,SQL 构造在多大程度上会成为瓶颈。 此外,Both tables must have distinct field names 文档中的这句话让我有点怀疑,因为我们确实有相同的字段名称,我们必须以某种方式处理它们,而不需要大量的变通方法。

【问题讨论】:

    标签: apache-kafka apache-flink


    【解决方案1】:

    如果 A 真的是静态的,那么如果你能以某种方式完全摄取 A,无论是进入 Flink 状态还是进入内存,然后将 B 流过 A,从而产生连接结果而不必存储 B,那么成本会更低。

    至少有两种方法可以使用 Flink 完成此任务。一种在this answer 中描述,另一种涉及使用State Processor API

    使用第二种方法,您可以将 A 保持在 key-partitioned Flink 状态。通过使用状态处理器 API,您可以引导一个包含所需状态的保存点,这样通过从该保存点开始您的作业,A 已经完全加载并立即可用。

    this gist 中有一个引导键控状态的简单示例。创建保存点后,您需要实现一个使用它来计算连接的流式作业——这可以通过 RichFlatMapFunction 完成。

    在不使用 Table API 的情况下实现连接的另一种替代方法是简单地使用 RichCoFlatMapFunction 或 KeyedCoProcessFunction 自行滚动。你会在 Flink 培训中找到 examples 。这些示例都没有真正符合您的要求,但它们给出了一般的味道。但是,我认为这没有任何优势——如果您要进行完全动态/动态连接,不妨使用 Table API。

    【讨论】:

    • 感谢您的回答!学习新技术在某种程度上是一个迭代过程,你经常让事情发挥作用,但随着时间的推移学习新的、更有效的方法,这导致废弃旧方法并从头开始重新设计使用(例如 Flink)。为了避免进度“重置”过大,我想从概念上解决如何以最佳方式实现这一目标的问题,而不是避免使用 Table API。
    • 顺便说一句,我说 A 是静态的时犯了一个错误,它不是真正静态的,而是偶尔发生变化,与快速和实际的流式相比B 的变化。从您的回答看来,Table API 是最合适的,我可能会看一下自定义 RichCoFlatMapFunction,但我认为实现自己的连接函数可能会带来一个全新的复杂层,这已经被 pre - Flink 的现有概念。
    • 在 A 不是真正静态的情况下,您可以考虑使用自定义分区器的策略,然后偶尔为 A 重新加载数据。如果 A 很大,这可能是值得的,因为这将减少 flink 必须管理的状态量。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-06-30
    • 1970-01-01
    • 2019-01-24
    • 1970-01-01
    • 2023-04-08
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多