【发布时间】:2020-09-20 11:54:07
【问题描述】:
在我的 Java 应用程序中,我有 三个 DataStreams。例如,一个流数据是从 Kafka 消费的,另一个流数据是从 Apache Nifi 消费的。对于这两个流的对象类型是不同的。例如Stream-1对象类型为Person,Stream-2对象类型为Address。
第三个是广播流(因为这个数据是从Kafka消费的)。
现在我想在 Job 类中组合 Stream-1 和 Stream-2,并希望在任务流程元素中进行拆分。如何实现?
注意: Stream-1 是主流,Stream-2 是侧输入。 MainStream 不断从 Kafka 获取数据。对于 Side Input,最初在应用程序启动时,所有表数据都从 DB 加载,然后在表数据更新时(不频繁)读取新数据。
示例结构:
DataStream<Person> stream-1 = env.addSource(read data from kafka)....
DataStream<Address> stream-2 = env.addSource(read data from nifi)....
BroadcastStream<String> BroadCastStream = stream-3.broadcast(read data from kafka);
我被称为以下链接。
FLIP-17 Side Inputs for DataStream API
我的用例是:
使用缓慢演变的数据加入流:我们用于丰富的辅助输入随着时间的推移而演变(从数据库中读取数据)。这可以通过在处理主输入之前等待一些初始数据可用并在新数据到达时不断将新数据引入内部侧输入结构来完成。
【问题讨论】:
-
您能否更新您的问题,说明为什么加入三个来源是不够的?还可以看看 [temporal joins|ci.apache.org/projects/flink/flink-docs-stable/dev/table/….
-
加入不够。为什么,因为在我的情况下,每种流类型都是不同的。 Join 仅适用于相同类型的流。
-
你到底是什么意思?即使它们具有不同的类型,您也可以轻松地加入 stream1 和 stream2。然后您可以将广播添加到结果中。
-
您好,Arvid,感谢您提供详细信息。让我试试看。
-
嗨 Arvid Heise,我是 Flink 的新手。你有加入不同流类型和广播的示例代码部分吗?
标签: java apache-flink apache-nifi flink-streaming flink-batch