Flink 系例 之 Process

process算子:处理每个keyBy(分区)输入到窗口的批量数据流(为KeyedStream类型数据流)

示例环境

1java.version: 1.8.x 2flink.version: 1.11.1

示例数据源 (项目码云下载)

Flink 系例 之 搭建开发环境与数据

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
点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

Flink 系例 之 Process - HelloWorld