欧美一级特黄大片做受成人-亚洲成人一区二区电影-激情熟女一区二区三区-日韩专区欧美专区国产专区

Flink1.10中SQL、HiveCatalog與事件時間整合的示例分析

這篇文章將為大家詳細講解有關Flink 1.10中SQL、HiveCatalog與事件時間整合的示例分析,小編覺得挺實用的,因此分享給大家做個參考,希望大家閱讀完這篇文章后可以有所收獲。

渦陽網站建設公司創(chuàng)新互聯(lián),渦陽網站設計制作,有大型網站制作公司豐富經驗。已為渦陽超過千家提供企業(yè)網站建設服務。企業(yè)網站搭建\外貿網站建設要多少錢,請找那個售后服務好的渦陽做網站的公司定做!

Flink 1.10 與 1.9 相比又是個創(chuàng)新版本,在我們感興趣的很多方面都有改進,特別是 Flink SQL。本文用根據埋點日志計算 PV、UV 的簡單示例來體驗 Flink 1.10 的兩個重要新特性:

  • 一是 SQL DDL 對事件時間的支持;
  • 二是 Hive Metastore 作為 Flink 的元數(shù)據存儲(即 HiveCatalog)。

這兩點將會為我們構建實時數(shù)倉提供很大的便利。  

添加依賴項

示例采用 Hive 版本為 1.1.0,Kafka 版本為 0.11.0.2。  
 

 
要使 Flink 與 Hive 集成以使用 HiveCatalog,需要先將以下 JAR 包放在 ${FLINK_HOME}/lib 目錄下。  

 
  • flink-connector-hive_2.11-1.10.0.jar
  • flink-shaded-hadoop-2-uber-2.6.5-8.0.jar
  • hive-metastore-1.1.0.jar
  • hive-exec-1.1.0.jar
  • libfb303-0.9.2.jar

 
后三個 JAR 包都是 Hive 自帶的,可以在 ${HIVE_HOME}/lib 目錄下找到。前兩個可以通過   阿里云 Maven    搜索 GAV 找到并手動下載(groupId 都是org.apache.flink)。  

 
再在 pom.xml 內添加相關的 Maven 依賴。  

 
 
Maven 下載:
    https://maven.aliyun.com/mvn/search
 
   
     
   
   
   
<properties>    <scala.bin.version>2.11</scala.bin.version>    <flink.version>1.10.0</flink.version>    <hive.version>1.1.0</hive.version>  </properties>
 <dependencies>    <dependency>      <groupId>org.apache.flink</groupId>      <artifactId>flink-table-api-scala_${scala.bin.version}</artifactId>      <version>${flink.version}</version>    </dependency>    <dependency>      <groupId>org.apache.flink</groupId>      <artifactId>flink-table-api-scala-bridge_${scala.bin.version}</artifactId>      <version>${flink.version}</version>    </dependency>    <dependency>      <groupId>org.apache.flink</groupId>      <artifactId>flink-table-planner-blink_${scala.bin.version}</artifactId>      <version>${flink.version}</version>    </dependency>    <dependency>      <groupId>org.apache.flink</groupId>      <artifactId>flink-sql-connector-kafka-0.11_${scala.bin.version}</artifactId>      <version>${flink.version}</version>    </dependency>    <dependency>      <groupId>org.apache.flink</groupId>      <artifactId>flink-connector-hive_${scala.bin.version}</artifactId>      <version>${flink.version}</version>    </dependency>    <dependency>      <groupId>org.apache.flink</groupId>      <artifactId>flink-json</artifactId>      <version>${flink.version}</version>    </dependency>    <dependency>      <groupId>org.apache.hive</groupId>      <artifactId>hive-exec</artifactId>      <version>${hive.version}</version>    </dependency>  </dependencies>
           
             
           
最后,找到 Hive 的配置文件 hive-site.xml,準備工作就完成了。              
           

           

注冊 HiveCatalog、創(chuàng)建數(shù)據庫

不多廢話了,直接上代碼,簡潔易懂。  

 
   
     
   
   
   
val streamEnv = StreamExecutionEnvironment.getExecutionEnvironment    streamEnv.setParallelism(5)    streamEnv.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
   val tableEnvSettings = EnvironmentSettings.newInstance()        .useBlinkPlanner()        .inStreamingMode()        .build()    val tableEnv = StreamTableEnvironment.create(streamEnv, tableEnvSettings)
   val catalog = new HiveCatalog(      "rtdw",                   // catalog name      "default",                // default database      "/Users/lmagic/develop",  // Hive config (hive-site.xml) directory      "1.1.0"                   // Hive version    )    tableEnv.registerCatalog("rtdw", catalog)    tableEnv.useCatalog("rtdw")
   val createDbSql = "CREATE DATABASE IF NOT EXISTS rtdw.ods"    tableEnv.sqlUpdate(createDbSql)
           

           

創(chuàng)建 Kafka 流表并指定事件時間

我們的埋點日志存儲在指定的 Kafka topic 里,為 JSON 格式,簡化版 schema 大致如下。  
   
     
   
   
   

           
"eventType": "clickBuyNow",    "userId": "97470180",    "shareUserId": "",    "platform": "xyz",    "columnType": "merchDetail",    "merchandiseId": "12727495",    "fromType": "wxapp",    "siteId": "20392",    "categoryId": "",    "ts": 1585136092541
           

           
其中 ts 字段就是埋點事件的時間戳(毫秒)。在 Flink 1.9 時代,用 CREATE TABLE 語句創(chuàng)建流表時是無法指定事件時間的,只能默認用處理時間。而在 Flink 1.10 下,可以這樣寫。  

 
   
     
   
   
   
CREATE TABLE rtdw.ods.streaming_user_active_log (  eventType STRING COMMENT '...',  userId STRING,  shareUserId STRING,  platform STRING,  columnType STRING,  merchandiseId STRING,  fromType STRING,  siteId STRING,  categoryId STRING,  ts BIGINT,  procTime AS PROCTIME(), -- 處理時間  eventTime AS TO_TIMESTAMP(FROM_UNIXTIME(ts / 1000, 'yyyy-MM-dd HH:mm:ss')), -- 事件時間  WATERMARK FOR eventTime AS eventTime - INTERVAL '10' SECOND -- 水印) WITH (  'connector.type' = 'kafka',  'connector.version' = '0.11',  'connector.topic' = 'ng_log_par_extracted',  'connector.startup-mode' = 'latest-offset', -- 指定起始offset位置  'connector.properties.zookeeper.connect' = 'zk109:2181,zk110:2181,zk111:2181',  'connector.properties.bootstrap.servers' = 'kafka112:9092,kafka113:9092,kafka114:9092',  'connector.properties.group.id' = 'rtdw_group_test_1',  'format.type' = 'json',  'format.derive-schema' = 'true', -- 由表schema自動推導解析JSON  'update-mode' = 'append')
           

                         
Flink SQL 引入了計算列(computed column)的概念,其語法為 column_name AS computed_column_expression,它的作用是在表中產生數(shù)據源 schema 不存在的列,并且可以利用原有的列、各種運算符及內置函數(shù)。比如在以上 SQL 語句中,就利用內置的 PROCTIME() 函數(shù)生成了處理時間列,并利用原有的 ts 字段與 FROM_UNIXTIME()、TO_TIMESTAMP() 兩個時間轉換函數(shù)生成了事件時間列。  

 
為什么 ts 字段不能直接用作事件時間呢?因為 Flink SQL 規(guī)定時間特征必須是 TIMESTAMP(3) 類型,即形如"yyyy-MM-ddTHH:mm:ssZ"格式的字符串,Unix 時間戳自然是不行的,所以要先轉換一波。  

 
既然有了事件時間,那么自然要有水印。Flink SQL 引入了 WATERMARK FOR rowtime_column_name AS watermark_strategy_expression 的語法來產生水印,有以下兩種通用的做法:  

 
  • 單調不減水?。▽?DataStream API 的 AscendingTimestampExtractor)

 
   
     
   
   
   
WATERMARK FOR rowtime_column AS rowtime_column - INTERVAL '0.001' SECOND
           

           
  • 有界亂序水印(對應 DataStream API 的 BoundedOutOfOrdernessTimestampExtractor)
WATERMARK FOR rowtime_column AS rowtime_column - INTERVAL 'n' TIME_UNIT
 

 
上文的 SQL 語句中就是設定了 10 秒的亂序區(qū)間。如果看官對水印、AscendingTimestampExtractor 和 BoundedOutOfOrdernessTimestampExtractor 不熟的話,可以參見之前的   這篇   ,就能理解為什么會是這樣的語法了。  

 
https://www.jianshu.com/p/c612e95a5028
下面來正式建表。  

 
    val createTableSql =      """        |上文的SQL語句        |......      """.stripMargin    tableEnv.sqlUpdate(createTableSql)
 

     
執(zhí)行完畢后,我們還可以去到 Hive 執(zhí)行 DESCRIBE FORMATTED ods.streaming_user_active_log 語句,能夠發(fā)現(xiàn)該表并沒有事實上的列,而所有屬性(包括 schema、connector、format 等等)都作為元數(shù)據記錄在了 Hive Metastore 中。  

 
Flink 1.10中SQL、HiveCatalog與事件時間整合的示例分析  
Flink 1.10中SQL、HiveCatalog與事件時間整合的示例分析  

 
Flink SQL 創(chuàng)建的表都會帶有一個標記屬性 is_generic=true,圖中未示出。  

 

開窗計算 PV、UV

用30秒的滾動窗口,按事件類型來分組,查詢語句如下。  

 
   
     
   
   
   
SELECT eventType,TUMBLE_START(eventTime, INTERVAL '30' SECOND) AS windowStart,TUMBLE_END(eventTime, INTERVAL '30' SECOND) AS windowEnd,COUNT(userId) AS pv,COUNT(DISTINCT userId) AS uvFROM rtdw.ods.streaming_user_active_logWHERE platform = 'xyz'GROUP BY eventType, TUMBLE(eventTime, INTERVAL '30' SECOND)
           

                         
關于窗口在 SQL 里的表達方式請參見   官方文檔   。1.10 版本 SQL 的官方文檔寫的還是比較可以的。    
SQL 文檔:
    https://ci.apache.org/projects/flink/flink-docs-release-1.10/dev/table/sql/queries.html#group-windows
 
懶得再輸出到一個結果表了,直接轉換成流打到屏幕上。    
    val queryActiveSql =      """        |......        |......      """.stripMargin    val result = tableEnv.sqlQuery(queryActiveSql)
   result        .toAppendStream[Row]        .print()        .setParallelism(1)
 

關于“Flink 1.10中SQL、HiveCatalog與事件時間整合的示例分析”這篇文章就分享到這里了,希望以上內容可以對大家有一定的幫助,使各位可以學到更多知識,如果覺得文章不錯,請把它分享出去讓更多的人看到。

網頁題目:Flink1.10中SQL、HiveCatalog與事件時間整合的示例分析
分享URL:http://aaarwkj.com/article6/igosog.html

成都網站建設公司_創(chuàng)新互聯(lián),為您提供品牌網站建設、建站公司、云服務器、手機網站建設營銷型網站建設、品牌網站設計

廣告

聲明:本網站發(fā)布的內容(圖片、視頻和文字)以用戶投稿、用戶轉載內容為主,如果涉及侵權請盡快告知,我們將會在第一時間刪除。文章觀點不代表本網站立場,如需處理請聯(lián)系客服。電話:028-86922220;郵箱:631063699@qq.com。內容未經允許不得轉載,或轉載時需注明來源: 創(chuàng)新互聯(lián)

成都seo排名網站優(yōu)化
说中文字幕的黄色大网站| 日本中文有码在线观看| 国产精品久久电影观看| 久久最新最热视频精品| 亚洲av永久国产剧情| 国产精品一区欧美精品| 亚洲日本欧美一区二区| 国产精品午夜福利天堂| 国产三级在线视频不卡| 中文字幕在线一区国产精品| 精品国产一区=区三区乱码| 国产精品色网在线播放| 日本中文字幕一二三四区| 日本一区二区三区精彩视频| 未满十八禁止下载软件| av在线免费观看大全| 中文字幕人妻丝袜一区一三区| 国产日韩精品专区一区| 加勒比av免费在线播放| 亚洲欧美日韩制服另类| 人妻系列日本在线播放| 日韩一区二区三精品| 国模在线视频一区二区| 国内一级黄色片免费观看| 国产精品亚洲二区三区| 国产精品一区二区三区久久| 天天干天天干夜夜操| 青青草针对华人在线视频| 欧美国产成人精品一区| 亚洲精品色婷婷一区二区| 国产亚洲精品国产福利久久| 亚洲综合中文字幕经典av在线| 蜜臀久久精品国产综合| 午夜福利视频在线一区| 欧洲亚洲国产一区二区| 丁香六月综合激情啪啪啪| 97在线观看视频在线观看| 中文字幕人妻熟女人妻| 亚洲黄色成人免费观看| av真人青青小草一区二区欧美| 十八禁真人无摭挡观看|