【问题标题】:Apache Flink : Add side inputs for DataStream APIApache Flink:为 DataStream API 添加侧输入
【发布时间】: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

jira/browse/FLINK-6131

我的用例是:

使用缓慢演变的数据加入流:我们用于丰富的辅助输入随着时间的推移而演变(从数据库中读取数据)。这可以通过在处理主输入之前等待一些初始数据可用并在新数据到达时不断将新数据引入内部侧输入结构来完成。

【问题讨论】:

  • 您能否更新您的问题,说明为什么加入三个来源是不够的?还可以看看 [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


【解决方案1】:

根据最新的回复,@Arvid 的推荐实际上正是这里所需要的。

答案的核心:

即使stream1和stream2不同,您也可以轻松加入 类型。然后你可以将广播添加到结果中

链接到doc 和example,以及文档中的相关sn-p(示例太长,无法包含在此处):

import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
 
...

DataStream<Integer> orangeStream = ...
DataStream<Integer> greenStream = ...

orangeStream.join(greenStream)
    .where(<KeySelector>)
    .equalTo(<KeySelector>)
    .window(TumblingEventTimeWindows.of(Time.milliseconds(2)))
    .apply (new JoinFunction<Integer, Integer, String> (){
        @Override
        public String join(Integer first, Integer second) {
            return first + "," + second;
        }
    });

【讨论】:

  • 感谢丹尼斯提供的信息。一疑。在您的示例代码块中,orangestream 和 greenStream 都是相同的数据类型(都是整数)。如果一个是Integer,另一个是String类型,我们可以申请join()函数吗?
  • @Azhagesan 我现在无法检查,但从概念上讲,键需要“适合”,其余部分可以根据需要组合在一起。例如,示例已经展示了如何获取两个整数并生成一个字符串,我希望它是微不足道的,例如:获取一个整数和一个字符串并生成一个字符串。 -- 在最坏的情况下,您可以将字段显式转换为不同的数据类型。
  • 谢谢,丹尼斯。让我检查一下
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-07-14
  • 2021-04-16
  • 1970-01-01
  • 1970-01-01
  • 2020-02-11
  • 2018-09-13
相关资源
最近更新 更多