flink api-datastream api-transformation算子
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();
}
}
更多推荐



所有评论(0)