Flink Transformation算子是DataStream API中用于数据转换的核心操作,主要分为以下几类:

基本转换算子

‌Map‌:一对一转换,对每个元素应用函数并输出新元素。例如将字符串转为大写或提取日志字段。
‌FlatMap‌:一对多转换,将单个元素拆分为零个或多个元素,常用于文本分词或嵌套结构展开。
‌Filter‌:根据条件过滤元素,仅保留满足布尔表达式的数据。

分组与聚合算子

KeyBy‌:按指定Key哈希分区,为聚合操作提供数据局部性支持。
Reduce‌:滚动聚合,基于前一次结果和当前元素计算新值,如累加或极值统计。
Window‌:基于时间或数量的窗口操作,支持滚动、滑动等窗口类型,需配合聚合函数使用。

多流操作算子

Union‌:合并多个同类型DataStream,不消除重复数据。
Connect‌:连接不同类型DataStream,生成ConnectedStreams,后续可通过CoMap/CoFlatMap处理。
‌Split/Select‌:将流按条件拆分为多个子流,再通过Select选择特定子流。

物理分区算子

Shuffle‌:随机均匀重分布数据,避免倾斜。
‌Rebalance‌:轮询分配数据到下游任务,实现负载均衡。
Broadcast‌:将数据广播到所有并行任务。

特殊转换算子

Project‌(仅限Tuple类型):选择字段子集,类似SQL的SELECT操作。
‌Iterate‌:迭代反馈流,用于实现循环逻辑。

例子

数据对象定义

package com.atguigu.wc.pojo;

import java.util.Objects;

public class WaterSensor {
    // 水位传感器类型
    public String id;
    // 传感器记录时间戳
    public Long ts;
    // 水位记录
    public Integer vc;

    public WaterSensor(){

    }

    public WaterSensor(String id, Long ts, Integer vc) {
        this.id = id;
        this.ts = ts;
        this.vc = vc;
    }


    public String getId() {
        return id;
    }

    public void setId(String id) {
        this.id = id;
    }

    public Long getTs() {
        return ts;
    }

    public void setTs(Long ts) {
        this.ts = ts;
    }

    public Integer getVc() {
        return vc;
    }

    public void setVc(Integer vc) {
        this.vc = vc;
    }

    @Override
    public String toString() {
        return "WaterSensor{" +
                "id='" + id + '\'' +
                ", ts=" + ts +
                ", vc=" + vc +
                "}";
    }

    @Override
    public boolean equals(Object o) {
        if (this == o) return true;
        if (o == null || getClass() != o.getClass()) return false;
        WaterSensor that = (WaterSensor) o;
        return  Objects.equals(id, that.id) && Objects.equals(ts, that.ts) && Objects.equals(vc, that.vc);
    }

    @Override
    public int hashCode() {
        return Objects.hash(id, ts, vc);
    }

}

基本转换算子

‌Map‌

// 方式一:传入匿名类,实现MapFunction
waterSensorDataStreamSource.map(new MapFunction<WaterSensor, String>() {
        @Override
        public String map(WaterSensor waterSensor) throws Exception {
              return waterSensor.id;
        }
}).print();

// 方式二:传入MapFunction实现类
waterSensorDataStreamSource.map(new UserMap()).print();

‌FlatMap‌

public class TransFlatMap {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<WaterSensor> waterSensorDataStreamSource = env.fromElements(
                new WaterSensor("sensor_1", 1L, 1),
                new WaterSensor("sensor_1", 2L, 2),
                new WaterSensor("sensor_2", 2L, 2),
                new WaterSensor("sensor_3", 3L, 3)
        );

        waterSensorDataStreamSource.flatMap(new MyFlatMap()).print();

        env.execute();


    }

    private static class MyFlatMap implements FlatMapFunction<WaterSensor,String> {
        @Override
        public void flatMap(WaterSensor waterSensor, Collector<String> collector) throws Exception {
            if(waterSensor.id.equals("sensor_1")){
                collector.collect("sensor_1 " + String.valueOf(waterSensor.vc));
            }else if(waterSensor.id.equals("sensor_2")){
                collector.collect("sensor_2_one "+String.valueOf(waterSensor.ts));
                collector.collect("sensor_2_two "+String.valueOf(waterSensor.vc));
            }
        }
    }
}

‌Filter‌

public class TransFilter {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<WaterSensor> waterSensorDataStreamSource = env.fromElements(
                new WaterSensor("sensor_1", 1L, 1),
                new WaterSensor("sensor_1", 2L, 2),
                new WaterSensor("sensor_2", 2L, 2),
                new WaterSensor("sensor_3", 3L, 3)
        );


        // 方式一:传入匿名类实现FilterFunction
        waterSensorDataStreamSource.filter(new FilterFunction<WaterSensor>() {
            @Override
            public boolean filter(WaterSensor waterSensor) throws Exception {
                return waterSensor.id.equals("sensor_1");
            }
        }).print();

        // 方式二:传入FilterFunction实现类
        waterSensorDataStreamSource.filter(new UserFilter()).print();
        env.execute();

    }

    private static class UserFilter implements FilterFunction<WaterSensor> {
        @Override
        public boolean filter(WaterSensor waterSensor) throws Exception {
            return waterSensor.id.equals("sensor_1");
        }
    }
}

分组与聚合算子

KeyBy‌

public class TransKeyBy {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<WaterSensor> waterSensorDataStreamSource = env.fromElements(
                new WaterSensor("sensor_1", 1L, 1),
                new WaterSensor("sensor_1", 2L, 2),
                new WaterSensor("sensor_2", 2L, 2),
                new WaterSensor("sensor_3", 3L, 3)
        );


        // 方式一:使用Lambda表达式
        KeyedStream<WaterSensor, String> waterSensorStringKeyedStream = waterSensorDataStreamSource.keyBy(e -> e.id);

        // 方式二:使用匿名类实现KeySelector
        waterSensorStringKeyedStream.keyBy(new KeySelector<WaterSensor, String>() {
            @Override
            public String getKey(WaterSensor waterSensor) throws Exception {
                return waterSensor.id;
            }
        });

        env.execute();
    }
}

‌Reduce‌

public class TransReduce {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env.socketTextStream("hadoop101", 9999)
                .map(new WaterSensorMapFunction())
                .keyBy(WaterSensor::getId)
                .reduce(new ReduceFunction<WaterSensor>() {
                    @Override
                    public WaterSensor reduce(WaterSensor waterSensor, WaterSensor t1) throws Exception {
                        int maxVc = Math.max(waterSensor.getVc(), t1.getVc());

                        if(waterSensor.getVc() > t1.getVc()) {
                            waterSensor.setVc(maxVc);
                            return waterSensor;
                        }else{
                            t1.setVc(maxVc);
                            return t1;
                        }


                    }
                }).print();

        env.execute();
    }
}

Window‌:TODO基于时间或数量的窗口操作,支持滚动、滑动等窗口类型,需配合聚合函数使用。

多流操作算子

Union‌

public class UnionExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        DataStreamSource<Integer> integerDataStreamSource = env.fromElements(1, 2, 3);
        DataStreamSource<Integer> integerDataStreamSource1 = env.fromElements(2, 2, 3);
        DataStreamSource<String> stringDataStreamSource = env.fromElements("2", "2", "3");

        integerDataStreamSource.union(integerDataStreamSource1,stringDataStreamSource.map(Integer::valueOf)).print();
        env.execute();
    }
}

Connect‌

public class ConnectDemo {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        SingleOutputStreamOperator<Integer> source1 = env.socketTextStream("hadoop101", 7777).map(Integer::valueOf);

        DataStreamSource<String> source2 = env.socketTextStream("hadoop101", 9999);
		// connect 类型可以不同
        ConnectedStreams<Integer, String> connect = source1.connect(source2);

        SingleOutputStreamOperator<String> map = connect.map(new CoMapFunction<Integer, String, String>() {
            @Override
            public String map1(Integer integer) throws Exception {
                return "来源于数字流:" + integer;
            }

            @Override
            public String map2(String s) throws Exception {
                return "来源于字母流:" + s;
            }
        });
        map.print();
        env.execute();
    }
}

‌Split/Select‌

public class SplitStreamByOutputTag {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<String> ds = env.socketTextStream("hadoop101", 9999);
        SingleOutputStreamOperator<WaterSensor> map = ds.map(new WaterSensorMapFunction());
        
        OutputTag<WaterSensor> outputTag = new OutputTag<>("s1",Types.POJO(WaterSensor.class));
        OutputTag<WaterSensor>  outputTag1 = new OutputTag<>("s2", Types.POJO(WaterSensor.class));

        // 返回的都是主流数据
        SingleOutputStreamOperator<WaterSensor> ds1 = map.process(new ProcessFunction<WaterSensor, WaterSensor>() {
            @Override
            public void processElement(WaterSensor waterSensor, ProcessFunction<WaterSensor, WaterSensor>.Context context, Collector<WaterSensor> collector) throws Exception {
                if("s1".equals(waterSensor.getId())){
                    context.output(outputTag,waterSensor);
                }else if("s2".equals(waterSensor.getId())){
                    context.output(outputTag1,waterSensor);
                }else {
                    collector.collect(waterSensor);
                }
            }
        });
        // 分别输出
        ds1.print();
        ds1.getSideOutput(outputTag).print();
        ds1.getSideOutput(outputTag1).print();


        env.execute();
    }
}

物理分区算子

Shuffle‌

public class ShuffleExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);

        env.socketTextStream("hadoop101", 9999)
                .shuffle().print();
        env.execute();
    }
}

‌Rebalance‌

public class RebalanceExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);

        env.socketTextStream("hadoop101", 9999)
                .rebalance().print();
        env.execute();
    }
}

Broadcast‌

public class BroadCastExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);

        env.socketTextStream("hadoop101", 9999)
                .broadcast().print();
        env.execute();
    }
}
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐