process算子:处理每个keyBy(分区)输入到窗口的批量数据流(为KeyedStream类型数据流)
示例环境
1java.version: 1.8.x 2flink.version: 1.11.1
示例数据源 (项目码云下载)
Process.java
1import com.flink.examples.DataSource; 2import org.apache.flink.api.java.functions.KeySelector; 3import org.apache.flink.api.java.tuple.Tuple3; 4import org.apache.flink.streaming.api.datastream.DataStream; 5import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; 6import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction; 7import org.apache.flink.streaming.api.windowing.windows.GlobalWindow; 8import org.apache.flink.util.Collector; 9import java.util.Iterator; 10import java.util.List; 11/** 12 * @Description process算子:处理每个keyBy(分区)输入到窗口的批量数据流(为KeyedStream类型数据流) 13 */ 14public class Process { 15 16 /** 17 * 遍历集合,分别打印不同性别的总人数与年龄之和 18 * @param args 19 * @throws Exception 20 */ 21 public static void main(String[] args) throws Exception { 22 final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); 23 List<Tuple3<String, String, Integer>> tuple3List = DataSource.getTuple3ToList(); 24 DataStream<String> dataStream = env.fromCollection(tuple3List) 25 .keyBy((KeySelector<Tuple3<String, String, Integer>, String>) k -> k.f1) 26 //按数量窗口滚动,每3个输入数据流,计算一次 27 .countWindow(3) 28 //处理每keyBy后的窗口数据流,process方法通常应用于KeyedStream类型的数据流处理 29 .process(new ProcessWindowFunction<Tuple3<String, String, Integer>, String, String, GlobalWindow>() { 30 /** 31 * 处理窗口数据集合 32 * @param s 从keyBy里返回的key值 33 * @param context 窗口的上下文 34 * @param input 从窗口获取的所有分区数据流 35 * @param out 输出数据流对象 36 * @throws Exception 37 */ 38 @Override 39 public void process(String s, Context context, Iterable<Tuple3<String, String, Integer>> input, Collector<String> out) throws Exception { 40 Iterator<Tuple3<String, String, Integer>> iterator = input.iterator(); 41 int total = 0; 42 int i = 0; 43 while (iterator.hasNext()){ 44 Tuple3<String, String, Integer> tuple3 = iterator.next(); 45 total += tuple3.f2; 46 i ++ ; 47 } 48 out.collect(s + "共:"+i+"人,平均年龄:" + total/i); 49 } 50 }); 51 dataStream.print(); 52 env.execute("flink Process job"); 53 } 54}
打印结果
14> girl共:3人,平均年龄:24 22> man共:3人,平均年龄:26