【发布时间】: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 文档中的这句话让我有点怀疑,因为我们确实有相同的字段名称,我们必须以某种方式处理它们,而不需要大量的变通方法。
【问题讨论】: