【问题标题】:What is the best way to do metadata join in Flink datastream?在 Flink 数据流中加入元数据的最佳方法是什么?
【发布时间】:2021-04-05 21:25:13
【问题描述】:

我们有一个 kafka 事件流,我们希望使用 MySQL DB 中的一些元数据来丰富它。

元数据每隔几个小时就会更改一次。本质上,我们希望定期读取数据库并使用这些新元数据不断丰富事件。

一种方法可能是将广播状态与周期性源一起使用,该源每隔几分钟/小时读取一次数据库。广播此流并使用它来加入。但问题可能是广播流的第一次读取可能晚于从 Kafka Stream 读取的某些消息。

有没有更好的办法?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    您可以为此使用 Flink SQL。根据具体要求,您可以针对来自 MySQL DB 的 CDC 流执行 time-versioned joins,或者针对 MySQL 执行 lookup joins(可能启用 optional cache)。

    另见https://github.com/ververica/flink-cdc-connectors

    更新:

    如果您想使用 DataStream API,但担心某些 kafka 消息可能会在广播流中的相应数据可用之前被处理,您可以:

    • 在扩充函数的 open() 方法中,对 MySQL 执行初始查询以预加载元数据
    • 如果在进行连接时广播数据仍然不可用,请使用在open() 期间获取的数据,或使用硬连线到代码中的一些默认值

    或者,您可以使用状态处理器 API 来引导广播状态的值。

    【讨论】:

    • 嘿@David 再次感谢,鉴于这个问题stackoverflow.com/questions/54748158/… 像这样的能力会很棒。我试图通过代码,似乎 StateDescriptor 接受一个默认值,只是 MapDescriptor 没有公开它。如果我们现在重写该类并公开默认值,可以吗?
    • 对不起,我不明白你的问题。
    • David ci.apache.org/projects/flink/flink-docs-master/api/java/org/… 可以看出StateDiscriptor 可以接受可以设置状态的默认值。所以在我的问题中,如果我能够为广播状态设置默认值,我应该能够使广播状态解决方案也正确工作吗?
    • 为状态描述符提供一个默认值并没有做任何你不能做的事情。当您需要广播状态的值时,您可以检查它是否为空并使用其他值。请参阅上面的扩展答案。
    猜你喜欢
    • 2020-04-16
    • 2017-03-26
    • 1970-01-01
    • 2016-03-05
    • 2011-02-19
    • 1970-01-01
    • 2021-05-12
    • 1970-01-01
    • 2013-06-07
    相关资源
    最近更新 更多