ICode9

精准搜索请尝试: 精确搜索
首页 > 数据库> 文章详细

Flink-Sql自定义UDF

2019-07-22 15:42:32  阅读:1588  来源: 互联网

标签:itemId 自定义 unionId Flink rankIndex 1563641998 UDF time action


最近尝试使用flink的table-sql,发现没有from_unixtime函数,只能自定义该udf。
原始kafka消息日志

{"action":"exposure","itemId":"16c65063e51d4d834722bf1a4b1d6378@TT@1576","rankIndex":14,"time":"1563641998","unionId":"ohmdTtymqiQw5aSxIt3ejxeAqpgs"}
{"action":"exposure","itemId":"15af1bc74e1ce2d7d0c12a0968618f1c@TT@16","rankIndex":10,"time":"1563641998","unionId":"ohmdTt_gZk2UkbbWsXBARMsTl1mI"}
{"action":"exposure","itemId":"15af1bc74e1ce2d7d0c12a0968618f1c@FT@287","rankIndex":21,"time":"1563641998","unionId":"ohmdTt_gZk2UkbbWsXBARMsTl1mI"}
{"action":"exposure","itemId":"15af1bc74e1ce2d7d0c12a0968618f1c@TT@12","rankIndex":22,"time":"1563641998","unionId":"ohmdTt_gZk2UkbbWsXBARMsTl1mI"}
{"action":"exposure","itemId":"b6f42135e217f70e97e214faf818ff07@TT@1523","rankIndex":10,"time":"1563641998","unionId":"ohmdTtzivXuT9u3oWFO5daAxziI0"}
{"action":"exposure","itemId":"b6f42135e217f70e97e214faf818ff07@TT@2759","rankIndex":25,"time":"1563641998","unionId":"ohmdTtzivXuT9u3oWFO5daAxziI0"}
{"action":"exposure","itemId":"15af1bc74e1ce2d7d0c12a0968618f1c@TT@68","rankIndex":2,"time":"1563641998","unionId":"ohmdTt_gZk2UkbbWsXBARMsTl1mI"}
{"action":"exposure","itemId":"16c65063e51d4d834722bf1a4b1d6378@TT@2045","rankIndex":13,"time":"1563641998","unionId":"ohmdTtymqiQw5aSxIt3ejxeAqpgs"}
{"action":"exposure","itemId":"16c65063e51d4d834722bf1a4b1d6378@TT@982","rankIndex":17,"time":"1563641998","unionId":"ohmdTtymqiQw5aSxIt3ejxeAqpgs"}
{"action":"exposure","itemId":"b6f42135e217f70e97e214faf818ff07@TT@1498","rankIndex":28,"time":"1563641998","unionId":"ohmdTtzivXuT9u3oWFO5daAxziI0"}

我们要格式化就是time字段。
自定义udf函数

import org.apache.flink.table.functions.ScalarFunction;

import java.text.SimpleDateFormat;
import java.util.Date;

public class FromUnixTimeUDF extends ScalarFunction {
    public String DATE_FORMAT;

    public FromUnixTimeUDF() {
        this.DATE_FORMAT = "yyyy-MM-dd HH:mm:ss";
    }

    public FromUnixTimeUDF(String dateFormat) {
        this.DATE_FORMAT = dateFormat;
    }

    public String eval(String longTime) {
        try {
            SimpleDateFormat sdf = new SimpleDateFormat(DATE_FORMAT);
            Date date = new Date(Long.parseLong(longTime) * 1000);
            return sdf.format(date);
        } catch (Exception e) {
            return null;
        }
    }
}

主程序main函数

		final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment();
    	StreamTableEnvironment tEnv = TableEnvironment.getTableEnvironment(env);
        tEnv.registerFunction("from_unixtime", new FromUnixTimeUDF());
        tEnv.connect(initKafkaDescriptor()).withFormat(new Json().failOnMissingField(true).deriveSchema())
        .withSchema(initSchema()).inAppendMode().registerTableSource("transfer_plan_show");
        Table result = tEnv.sqlQuery("select unionId,itemId,action,from_unixtime(`time`) as creat_time,rankIndex as rank_index from transfer_plan_show");
        result.printSchema();
        tEnv.toAppendStream(result, Row.class).print();
        env.execute();

相关函数

	//链接kafka配置
    private Kafka initKafkaDescriptor(){
      Kafka kafkaDescriptor=  new Kafka().version("0.11").topic("transfer_plan_show")
                .startFromLatest().property("bootstrap.servers", KafkaConfig.KAFKA_BROKER_LIST)
                .property("group.id", "trafficwisdom-streaming");
      return kafkaDescriptor;
    }
	//根据json自定义schema
    private Schema initSchema(){
        Schema schema=new Schema().field("action", Types.STRING())
                .field("itemId",Types.STRING())
                .field("time",Types.STRING())
                .field("unionId",Types.STRING())
                .field("rankIndex",Types.INT());
        return schema;
    }

说明因为sql中time是关键字,所以加上加上两个反斜杠 ``.

标签:itemId,自定义,unionId,Flink,rankIndex,1563641998,UDF,time,action
来源: https://blog.csdn.net/baifanwudi/article/details/96861504

本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享;
2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关;
3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关;
4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除;
5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。

专注分享技术,共同学习,共同进步。侵权联系[81616952@qq.com]

Copyright (C)ICode9.com, All Rights Reserved.

ICode9版权所有