在 Flink 流中使用静态 DataSet 丰富 DataStream

Vij*_*sal 5 data-analysis bigdata apache-flink flink-streaming

我正在编写一个 Flink 流程序,其中我需要使用一些静态数据集(信息库,IB)来丰富用户事件的 DataStream。

例如,假设我们有一个买家的静态数据集,我们有一个传入的事件点击流,对于每个事件,我们想要添加一个布尔标志,指示事件的执行者是否是买家。

实现此目的的理想方法是按用户 id 对传入流进行分区,让数据集中的买家集再次按用户 id 进行分区,然后在此数据集中查找流中的每个事件。

由于 Flink 不允许在流程序中使用 DataSets,我该如何实现上述功能?

另一种选择可能是使用托管运营商状态来存储买家集,但我如何保持此状态按用户 ID 分发,以避免在单个事件查找中进行网络输入/输出?在内存状态后端的情况下,状态是否保持由某个键分布,还是在所有操作员子任务中复制?

在 Flink 流程序中实现上述丰富需求的正确设计模式是什么?

Dav*_*son 5

我会通过 user_id 键控流,并使用 RichFlatMap 进行丰富。在 RichFlatMap 的 open() 方法中,您可以为该用户加载静态买家标志并将其缓存在布尔字段中。