Flink实战教程

Stella981
• 阅读 863

目录:

  • 自定义函数

  • 单个eval方法

  • 多个eval方法

  • 不固定参数

  • 通过注解指定返回类型

  • 注册函数

  • 构造数据源

  • 查询

  • left join

  • join

  • 多种类型参数

  • 不固定参数类型

今天我们来聊聊flink sql中另外一种自定义函数-TableFuntion.
TableFuntion 可以有0个、一个、多个输入参数,他的返回值可以是任意行,每行可以有多列数据.

实现自定义TableFunction需要继承TableFunction类,然后定义一个public类型的eval方法。结合官网的例子具体来讲解一下。

自定义函数

单个eval方法

   public static class Split extends TableFunction<Tuple2<String,Integer>> {        private String separator = ",";        public Split(String separator) {            this.separator = separator;        }        public void eval(String str) {            for (String s : str.split(separator)) {                collect(new Tuple2<String,Integer>(s, s.length()));            }        }    }

来解释一下:

  1. 这个函数接收一个字符串类型的入参,将传进来的字符串用指定分隔符拆分,然后返回值是一组Tuple2,每个Tuple2包含单词以及其长度.

  2. TableFunction是一个泛型类,需要指定返回值类型

  3. 不同于标量函数,eval方法没有返回值,使用collect方法来收集对象。

多个eval方法

 /**  * 注册多个eval方法,接收long或者string类型的参数,然后将他们转成string类型  */ public static class DuplicatorFunction extends TableFunction<String>{  public void eval(Long i){   eval(String.valueOf(i));  }  public void eval(String s){   collect(s);  } }

不固定参数

 /**  * 接收不固定个数的int型参数,然后将所有数据依次返回  */ public static class FlattenFunction extends TableFunction<Integer>{  public void eval(Integer... args){   for (Integer i: args){    collect(i);   }  } }

通过注解指定返回类型

 /**  * 通过注册指定返回值类型,flink 1.11 版本开始支持  */    @FunctionHint(output = @DataTypeHint("ROW< i INT, s STRING >"))    class DuplicatorFunction extends TableFunction<Row> {      public void eval(Integer i, String s) {        collect(Row.of(i, s));        collect(Row.of(i, s));      }    }

注册函数

这里使用blink的planner,然后把上述三个函数都注册了

  StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();  EnvironmentSettings bsSettings = EnvironmentSettings.newInstance()                                                      .useBlinkPlanner()                                                      .inStreamingMode()                                                      .build();  StreamTableEnvironment tEnv = StreamTableEnvironment.create(env, bsSettings);  tEnv.registerFunction("split", new Split(" "));  tEnv.registerFunction("duplicator", new DuplicatorFunction());  tEnv.registerFunction("flatten", new FlattenFunction());

构造数据源

        List<Tuple2<Long,String>> ordersData = new ArrayList<>();        ordersData.add(Tuple2.of(2L, "Euro"));        ordersData.add(Tuple2.of(1L, "US Dollar"));        ordersData.add(Tuple2.of(50L, "Yen"));        ordersData.add(Tuple2.of(3L, "Euro"));        DataStream<Tuple2<Long,String>> ordersDataStream = env.fromCollection(ordersData);        Table orders = tEnv.fromDataStream(ordersDataStream, "amount, currency, proctime.proctime");        tEnv.registerTable("Orders", orders);

查询

left join

 Table result = tEnv.sqlQuery(    "SELECT o.currency, T.word, T.length FROM Orders as o LEFT JOIN LATERAL TABLE(split(currency)) as T(word, length) ON TRUE");  tEnv.toAppendStream(result, Row.class).print();

解释一下:

  • 有两种使用方式, 使用 join的时候用LATERAL TABLE ,使用left join的时候用LATERAL TABLE .... ON TRUE.

  • 给TableFuntion返回的数据起一个别名:T(word, length),其中T是表的别名,word和length是字段别名,所以我们可以在sql中通过o.currency, T.word, T.length来查询字段。

join

 String sql = "SELECT o.currency, T.word, T.length FROM Orders as o ," +               " LATERAL TABLE(split(currency)) as T(word, length)";

多种类型参数

  String sql2 = "SELECT * FROM Orders as o , " +                "LATERAL TABLE(duplicator(amount))," +                "LATERAL TABLE(duplicator(currency))";

不固定参数类型

 String sql3 = "SELECT * FROM Orders as o , " +                "LATERAL TABLE(flatten(100,200,300))";

今天这个TableFuntion我们就先讲到这里,后续我们通过自定义的TableFuntion来实现一个mysql维表和hbase维表功能,用来在流式数据中补全字段信息. 完整代码请参考:

https://github.com/zhangjun0x01/bigdata-examples/blob/master/flink/src/main/java/sql/function/CustomTableFunction.java

更多精彩内容,欢迎关注我的公众号【大数据技术与应用实战】!

Flink实战教程

本文分享自微信公众号 - 大数据技术与应用实战(bigdata_bigdata)。
如有侵权,请联系 support@oschina.cn 删除。
本文参与“OSC源创计划”,欢迎正在阅读的你也加入,一起分享。

点赞
收藏
评论区
推荐文章
blmius blmius
3年前
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
皕杰报表之UUID
​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为
待兔 待兔
4个月前
手写Java HashMap源码
HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22
Jacquelyn38 Jacquelyn38
3年前
2020年前端实用代码段,为你的工作保驾护航
有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )
Wesley13 Wesley13
3年前
mysql设置时区
mysql设置时区mysql\_query("SETtime\_zone'8:00'")ordie('时区设置失败,请联系管理员!');中国在东8区所以加8方法二:selectcount(user\_id)asdevice,CONVERT\_TZ(FROM\_UNIXTIME(reg\_time),'08:00','0
Stella981 Stella981
3年前
HIVE 时间操作函数
日期函数UNIX时间戳转日期函数: from\_unixtime语法:   from\_unixtime(bigint unixtime\, string format\)返回值: string说明: 转化UNIX时间戳(从19700101 00:00:00 UTC到指定时间的秒数)到当前时区的时间格式举例:hive   selec
Wesley13 Wesley13
3年前
00:Java简单了解
浅谈Java之概述Java是SUN(StanfordUniversityNetwork),斯坦福大学网络公司)1995年推出的一门高级编程语言。Java是一种面向Internet的编程语言。随着Java技术在web方面的不断成熟,已经成为Web应用程序的首选开发语言。Java是简单易学,完全面向对象,安全可靠,与平台无关的编程语言。
Stella981 Stella981
3年前
Django中Admin中的一些参数配置
设置在列表中显示的字段,id为django模型默认的主键list_display('id','name','sex','profession','email','qq','phone','status','create_time')设置在列表可编辑字段list_editable
Wesley13 Wesley13
3年前
MySQL部分从库上面因为大量的临时表tmp_table造成慢查询
背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_
Python进阶者 Python进阶者
10个月前
Excel中这日期老是出来00:00:00,怎么用Pandas把这个去除
大家好,我是皮皮。一、前言前几天在Python白银交流群【上海新年人】问了一个Pandas数据筛选的问题。问题如下:这日期老是出来00:00:00,怎么把这个去除。二、实现过程后来【论草莓如何成为冻干莓】给了一个思路和代码如下:pd.toexcel之前把这