圆咕噜咕噜 发表于 4 天前

Flink转换算子

Apache Flink 是一个用于处理无界和有界数据的开源流处理框架。在 Flink 中,转换(Transformation)是数据流处理的焦点组件之一,它们界说了如何从输入数据集生成输出数据集。以下是 Flink 中一些常见的转换算子:

[*] Map: 将每个元素转换为另一个元素。通常用于一对一的数据转换。
DataStream<Integer> input = ...;
DataStream<Integer> result = input.map(new MapFunction<Integer, Integer>() {
    @Override
    public Integer map(Integer value) throws Exception {
      return value * 2;
    }
});
[*] FlatMap: 类似于 Map,但是可以产生多个输出元素。
DataStream<String> input = ...;
DataStream<String> result = input.flatMap(new FlatMapFunction<String, String>() {
    @Override
    public void flatMap(String value, Collector<String> out) throws Exception {
      for (String word : value.split(" ")) {
            out.collect(word);
      }
    }
});
[*] Filter: 过滤掉不符合条件的元素。
DataStream<Integer> input = ...;
DataStream<Integer> result = input.filter(new FilterFunction<Integer>() {
    @Override
    public boolean filter(Integer value) throws Exception {
      return value > 5;
    }
});
[*] KeyBy: 对数据流进行分组,以便后续操纵(如窗口操纵)能够按特定键处理数据。
DataStream<Tuple2<String, Integer>> input = ...;
KeyedStream<Tuple2<String, Integer>, String> keyed = input.keyBy(0);
[*] Reduce: 在每个分组内聚合元素,可以是滚动聚合或窗口内的聚合。
DataStream<Integer> input = ...;
DataStream<Integer> result = input.reduce(new ReduceFunction<Integer>() {
    @Override
    public Integer reduce(Integer value1, Integer value2) throws Exception {
      return value1 + value2;
    }
});
[*] Window: 界说窗口以对无界流中的数据进行分批处理。
DataStream<Integer> input = ...;
DataStream<Integer> result = input
    .keyBy((key) -> key)
    .timeWindow(Time.seconds(5))
    .reduce(new ReduceFunction<Integer>() {
      @Override
      public Integer reduce(Integer value1, Integer value2) throws Exception {
            return value1 + value2;
      }
    });
[*] Union: 合并两个或更多数据流。
DataStream<Integer> stream1 = ...;
DataStream<Integer> stream2 = ...;
DataStream<Integer> result = stream1.union(stream2);
[*] Connect: 连接两个不同范例的流,但不会将它们合并成一个流。
DataStream<Integer> stream1 = ...;
DataStream<String> stream2 = ...;
ConnectedStreams<Integer, String> connectedStreams = stream1.connect(stream2);
[*] Broadcast: 广播一个数据流到全部任务实例中。
DataStream<Integer> input = ...;
BroadcastStream<Integer> broadcastStream = input.broadcast();
这些只是 Flink 提供的一些根本转换算子。Flink 还支持更复杂的转换,例如通过 ProcessFunction 实现自界说逻辑,以及与其他外部体系的集成等高级特性。


[*] Map

[*]作用:将每个元素转换为另一个元素。
[*]示例:将每个整数乘以2。
[*]输入:1, 2, 3, 4, 5
[*]输出:2, 4, 6, 8, 10

[*] FlatMap

[*]作用:将每个元素转换为零个、一个或多个元素。
[*]示例:将每个字符串拆分成单词。
[*]输入:"hello world", "flink is great"
[*]输出:"hello", "world", "flink", "is", "great"

[*] Filter

[*]作用:过滤掉不符合条件的元素。
[*]示例:过滤出大于5的整数。
[*]输入:1, 2, 3, 4, 5, 6, 7, 8, 9, 10
[*]输出:6, 7, 8, 9, 10

[*] KeyBy

[*]作用:根据指定的键对数据流进行分组,以便后续操纵(如窗口操纵)能够按特定键处理数据。
[*]示例:按元组的第一个元素分组。
[*]输入:(a, 1), (b, 2), (a, 3), (b, 4)
[*]输出:分组后的数据流,按第一个元素分组。

[*] Reduce

[*]作用:在每个分组内聚合元素,可以是滚动聚合或窗口内的聚合。
[*]示例:在每个分组内将整数相加。
[*]输入:(a, 1), (b, 2), (a, 3), (b, 4)
[*]输出:(a, 4), (b, 6)

[*] Window

[*]作用:界说窗口以对无界流中的数据进行分批处理。
[*]示例:界说一个5秒的时间窗口,在每个窗口内将整数相加。
[*]输入:1, 2, 3, 4, 5, 6, 7, 8, 9, 10
[*]输出:窗口内的聚合效果,例如 15, 25, 35, ...

[*] Union

[*]作用:合并两个或更多数据流。
[*]示例:合并两个包罗整数的数据流。
[*]输入:stream1: 1, 2, 3, stream2: 4, 5, 6
[*]输出:1, 2, 3, 4, 5, 6

[*] Connect

[*]作用:连接两个不同范例的流,但不会将它们合并成一个流。
[*]示例:连接一个整数流和一个字符串流。
[*]输入:stream1: 1, 2, 3, stream2: "a", "b", "c"
[*]输出:ConnectedStreams<Integer, String>

[*] Broadcast

[*]作用:广播一个数据流到全部任务实例中。
[*]示例:广播一个包罗配置参数的数据流。
[*]输入:1, 2, 3
[*]输出:每个任务实例都收到 1, 2, 3


免责声明:如果侵犯了您的权益,请联系站长,我们会及时删除侵权内容,谢谢合作!更多信息从访问主页:qidao123.com:ToB企服之家,中国第一个企服评测及商务社交产业平台。
页: [1]
查看完整版本: Flink转换算子