• FlinkSQL: Create function using jar-located in HDFS


    序号作者版本时间备注
    1HamaWhite1.0.02022-12-05新增Flink UDF
    2HamaWhite1.0.12022-12-06新增Hive UDF

    一、基础信息

    1.1 组件版本

    • JDK:  1.8
    • Flink:    1.16.0
    • Hive:     3.1.2
    • Hadoop:  3.2.2

    二、准备工作

    2.1 新建Scalar Function

    1. package com.hw.flink.udf;
    2. import org.apache.flink.table.functions.ScalarFunction;
    3. public class MySuffixFunction extends ScalarFunction {
    4. public String eval(String input) {
    5. return input.concat("-HamaWhite");
    6. }
    7. }

    Maven依赖如下: 

    1. <dependencies>
    2. <dependency>
    3. <groupId>org.apache.flinkgroupId>
    4. <artifactId>flink-table-commonartifactId>
    5. <version>1.16.0version>
    6. <scope>providedscope>
    7. dependency>
    8. dependencies>

    2.2 构建Jar包

    把上述代码打包为 flink1.16-udf-demo.jar。

    2.3 上传HDFS

    1. $ hadoop fs -put flink1.16-udf-demo.jar .
    2. $ hadoop fs -ls /user/deploy
    3. -rw-r--r-- 1 deploy supergroup 1228 2022-12-05 20:03 /user/deploy/flink1.16-udf-demo.jar

    三、验证函数

    参考《flink-docs-release-1.16#create-function

    3.1 新建及测试函数

    1. EnvironmentSettings settings = EnvironmentSettings.inStreamingMode();
    2. TableEnvironment tableEnv = TableEnvironment.create(settings);
    3. // use the hive catalog
    4. String catalogName = "hive";
    5. HiveCatalog hiveCatalog = new HiveCatalog(catalogName, "flink_demo", "conf/hive");
    6. tableEnv.registerCatalog(catalogName, hiveCatalog);
    7. tableEnv.useCatalog(catalogName);
    8. tableEnv.executeSql("DROP FUNCTION IF EXISTS my_suffix_udf_jar");
    9. // create function my_suffix_udf_jar
    10. tableEnv.executeSql("CREATE FUNCTION IF NOT EXISTS my_suffix_udf_jar " +
    11. "AS 'com.hw.flink.udf.MySuffixFunction' " +
    12. "USING JAR 'hdfs://xxx.xxx.xxx.xxx:9000/user/deploy/flink1.16-udf-demo.jar'"
    13. );
    14. // use function my_suffix_udf_jar
    15. tableEnv.executeSql("select my_suffix_udf_jar('hw')").print();

    最后一行的运行结果是:

    1. +----+--------------------------------+
    2. | op | EXPR$0 |
    3. +----+--------------------------------+
    4. | +I | hw-HamaWhite |
    5. +----+--------------------------------+
    6. 1 row in set

    3.2 查看Hive Meta数据库

    上文用的是Hive Catalog,Flink创建函数的时候会把函数信息注册到Hive Meta数据库中,存储在FUNCS和FUNC_RU表中。查看数据如下:

    1. mysql> select * from FUNCS where FUNC_NAME like '%jar%' \G;
    2. *************************** 1. row ***************************
    3. FUNC_ID: 764
    4. CLASS_NAME: com.hw.flink.udf.MySuffixFunction
    5. CREATE_TIME: 1670242907
    6. DB_ID: 11
    7. FUNC_NAME: my_suffix_udf_jar
    8. FUNC_TYPE: 1
    9. OWNER_NAME: NULL
    10. OWNER_TYPE: GROUP
    11. 1 row in set (0.00 sec)
    1. mysql> select * from FUNC_RU \G;
    2. *************************** 1. row ***************************
    3. FUNC_ID: 764
    4. RESOURCE_TYPE: 1
    5. RESOURCE_URI: hdfs://xxx.xxx.xxx.xxx:9000/user/deploy/flink1.16-udf-demo.jar
    6. INTEGER_IDX: 0

    四、Flink中使用Hive UDF

    参考《flink-docs-release-1.16#hive_functions

    4.1  自定义Hive UDF

    1. package com.hw.hive.udf;
    2. import org.apache.hadoop.hive.ql.exec.UDF;
    3. public class MyPrefixFunction extends UDF {
    4. public String evaluate(String input) {
    5. return "HamaWhite-".concat(input);
    6. }
    7. }

    Maven依赖如下: 

    1. <dependencies>
    2. <dependency>
    3. <groupId>org.apache.hivegroupId>
    4. <artifactId>hive-execartifactId>
    5. <version>3.1.2version>
    6. <scope>providedscope>
    7. dependency>
    8. dependencies>

    4.2 构建Jar包

    把上述代码打包为 hive3-udf-demo.jar。

    4.3 上传HDFS

    $ hadoop fs -put hive3-udf-demo.jar .
    

    4.4 在Hive SQL 中创建UDF

    1. use flink_demo;
    2. CREATE FUNCTION my_prefix_udf_jar AS 'com.hw.hive.udf.MyPrefixFunction'
    3. USING JAR 'hdfs://xxx.xxx.xxx.xxx:9000/user/deploy/hive3-udf-demo.jar';

    4.5 Flink中测试Hive UDF

    在3.1的测试代码中增加测试代码:

    注: 经验证,官网中提到的 tableEnv.loadModule("hive", new HiveModule("3.1.2")) 不用写,Flink也能正常访问到Hive UDF。

    1. # use hive udf
    2. tableEnv.executeSql("select my_prefix_udf_jar('hw')").print();

     最后一行的运行结果是:

    1. +----+--------------------------------+
    2. | op | EXPR$0 |
    3. +----+--------------------------------+
    4. | +I | HamaWhite-hw |
    5. +----+--------------------------------+
    6. 1 row in set

    经测试,在Flink中是可以同时用Hive和Flink的UDF的。例如:

    1. # use flink udf
    2. tableEnv.executeSql("select my_suffix_udf_jar('hw')").print();
    3. # use hive udf
    4. tableEnv.executeSql("select my_prefix_udf_jar('hw')").print();
    5. # mixed use flink and hive udf
    6. tableEnv.executeSql("select my_suffix_udf_jar(my_prefix_udf_jar('hw'))").print();

    五、参考文献

    1. flink-docs-release-1.16#create-function
    2. flink-docs-release-1.16#hive_functions

  • 相关阅读:
    SpringBoot携带Jre绿色部署项目_免安装Jdk[Linux服务器]
    Calico IP In IP模拟组网
    加速乐源码(golang版本)
    Java互联网+公立医院绩效考核源码
    从助力跨境互通到保障农民工,区块链在大湾区做了什么? | 研讨会
    学习 Rust 的第十二天:如何使用向量
    YOLOv5训练自己的voc数据集
    阿里云视频上传实战
    java面试官:程序员,请你告诉我是谁把公司面试题泄露给你的?
    (原创)【B4A】一步一步入门03:APP名称、图标等信息修改
  • 原文地址:https://blog.csdn.net/xin_jmail/article/details/128192136