maxBy聚合:获取一组数据流算子中最大的记录行(和max的区别,max是返回计算字段的最大值)
示例环境
1java.version: 1.8.x 2flink.version: 1.11.1
示例数据源 (项目码云下载)
MaxBy.java
1import com.flink.examples.DataSource; 2import org.apache.flink.api.common.typeinfo.Types; 3import org.apache.flink.api.java.functions.KeySelector; 4import org.apache.flink.api.java.tuple.Tuple3; 5import org.apache.flink.streaming.api.datastream.DataStream; 6import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; 7import java.util.List; 8/** 9 * @Description maxBy聚合:获取一组数据流算子中最大的记录行(和max的区别,max是返回计算字段的最大值) 10 */ 11public class MaxBy { 12 /** 13 * 遍历集合,返回每个性别分区下最大年龄数据记录 14 * @param args 15 * @throws Exception 16 */ 17 public static void main(String[] args) throws Exception { 18 final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); 19 List<Tuple3<String, String, Integer>> tuple3List = DataSource.getTuple3ToList(); 20 DataStream<Tuple3<String, String, Integer>> dataStream = env.fromCollection(tuple3List) 21 .returns(Types.TUPLE(Types.STRING, Types.STRING,Types.INT)) 22 .keyBy((KeySelector<Tuple3<String, String, Integer>, String>) k ->k.f1) 23 //按数量窗口滚动,每3个输入数据流,计算一次 24 .countWindow(3) 25 //注意:计算变量为f2 26 .maxBy(2); 27 dataStream.print(); 28 env.execute("flink MaxBy job"); 29 } 30}
打印结果
14> (刘六,girl,32) 22> (吴八,man,30)