在大資料處理中,實時資料分析是一個重要的需求,隨著資料量的不斷增長,對于實時分析的挑戰也在不斷加大,傳統的批處理方式已經不能滿足實時資料處理的需求,需要一種更加高效的技術來解決這個問題,Apache Hudi(Hadoop Upserts Deletes and Incremental Processing)就是這樣一種技術,提供了高效的實時資料倉庫管理功能,
本文將介紹袋鼠云基于 Hudi 構建資料湖的整體方案架構及其在實時資料倉庫處理方面的特點,并且為大家展示一個使用 Apache Hudi 的簡單示例,便于新手上路,
Apache Hudi 介紹
Apache Hudi 是一個開源的資料湖存盤系統,可以在 Hadoop 生態系統中提供實時資料倉庫處理功能,Hudi 最早由 Uber 開發,后來成為 Apache 頂級專案,
Hudi 主要特性
· 支持快速插入和更新操作,以便在資料倉庫中實時處理資料;
· 提供增量查詢功能,可有效提高資料分析效率;
· 支持時間點查詢,以便查看資料在某一時刻的狀態;
· 與 Apache Spark、Hive 等大資料分析工具兼容,
Hudi 架構
Apache Hudi 的架構包括以下幾個主要組件:
· Hudi 資料存盤:Hudi 資料存盤是 Hudi 的核心組件,負責存盤資料,資料存盤有兩種型別:Copy-On-Write(COW)和 Merge-On-Read(MOR);
· Copy-On-Write:COW 存盤型別會在對資料進行更新時,創建一個新的資料檔案副本,將更新的資料寫入副本中,之后,新的資料檔案副本會替換原始資料檔案;
· Merge-On-Read:MOR 存盤型別會在查詢時,將更新的資料與原始資料進行合并,這種方式可以減少資料存盤的寫入延遲,但會增加查詢的計算量;
· Hudi 索引:Hudi 索參考于維護資料記錄的位置資訊,索引有兩種型別:內置索引(如 Bloom 過濾器)和外部索引(如 HBase 索引);
· Hudi 查詢引擎:Hudi 查詢引擎負責處理查詢請求,Hudi 支持多種查詢引擎,如 Spark SQL、Hive、Presto 等,

Hudi 的使用場景
Apache Hudi 可以幫助企業和組織實作實時資料處理和分析,實時資料處理需要快速地處理和查詢資料,同時還需要保證資料的一致性和可靠性,
Apache Hudi 的增量資料處理、ACID 事務性保證、寫時合并等技術特性可以幫助企業更好地實作實時資料處理和分析,基于 Hudi 的特性可以在一定程度上在實時數倉的構建程序中承擔上下游資料鏈路的對接(類似 Kafka 的角色),既能實作增量的資料處理,也能為批流一體的處理提供存盤基礎,
Hudi 的優勢和劣勢
● 優勢
· 高效處理大規模資料集;
· 支持實時資料更新和查詢;
· 實作了增量寫入機制,提高了資料訪問效率;
· Hudi 可以與流處理管道集成;
· Hudi 提供了時間旅行功能,允許回溯資料的歷史版本,
● 劣勢
· 在讀寫資料時需要付出額外的代價;
· 操作比較復雜,需要使用專業的編程語言和工具,
Hudi 在袋鼠云資料湖平臺上的實踐
Hudi 在袋鼠云資料湖的技術架構
Hudi 在袋鼠云的資料湖平臺上主要對資料湖管理提供助力:
· 元資料的接入,讓用戶可以快速的對表進行管理;
· 資料快速接入,包括對符合條件的原有表資料進行轉換,快速搭建資料湖能力;
· 湖表的管理,監控小檔案定期進行合并,提升表的查詢性能,內在豐富的表操作功能,包括 time travel ,孤兒檔案清理,過期快照清理等;
· 索引構建,提供多種索引包括 bloom filter,zorder 等,提升計算引擎的查詢性能,

Hudi 使用示例
在介紹了 Hudi 的基本資訊和袋鼠云資料湖平臺的結構之后,我們來看一個使用示例,替換 Flink 在記憶體中的 join 程序,
在 Flink 中對多流 join 往往是比較頭疼的場景,需要考慮 state ttl 時間設定,設定太小資料經常關聯不上,設定太大記憶體又需要很高才能保留,我們通過 Hudi 的方式來換個思路實作,
● 構建 catalog
public String createCatalog(){
String createCatalog = "CREATE CATALOG hudi_catalog WITH (\n" +
" 'type' = 'hudi',\n" +
" 'mode' = 'hms',\n" +
" 'default-database' = 'default',\n" +
" 'hive.conf.dir' = '/hive_conf_dir',\n" +
" 'table.external' = 'true'\n" +
")";
return createCatalog;
}
● 創建 hudi 表
public String createHudiTable(){
String createTable = "CREATE TABLE if not exists hudi_catalog.flink_db.test_hudi_flink_join_2 (\n" +
" id int ,\n" +
" name VARCHAR(10),\n" +
" age int ,\n" +
" address VARCHAR(10),\n" +
" dt VARCHAR(10),\n" +
" primary key(id) not enforced\n" +
")\n" +
"PARTITIONED BY (dt)\n" +
"WITH (\n" +
" 'connector' = 'hudi',\n" +
" 'table.type' = 'MERGE_ON_READ',\n" +
" 'changelog.enabled' = 'true',\n" +
" 'index.type' = 'BUCKET',\n" +
" 'hoodie.bucket.index.num.buckets' = '2',\n" +
String.format(" '%s' = '%s',\n", FlinkOptions.PRECOMBINE_FIELD.key(), FlinkOptions.NO_PRE_COMBINE) +
" 'write.payload.class' = '" + PartialUpdateAvroPayload.class.getName() + "'\n" +
");";
return createTable;
}
● 更新 hudi 表的 flink_db.test_hudi_flink_join_2 的 id, name, age, dt 列
01 從 kafka 中讀取 topic1
public String createKafkaTable1(){
String kafkaSource1 = "CREATE TABLE source1\n" +
"(\n" +
" id INT,\n" +
" name STRING,\n" +
" age INT,\n" +
" dt String,\n" +
" PROCTIME AS PROCTIME()\n" +
") WITH (\n" +
" 'connector' = 'kafka'\n" +
" ,'topic' = 'join_topic1'\n" +
" ,'properties.bootstrap.servers' = 'localhost:9092'\n" +
" ,'scan.startup.mode' = 'earliest-offset'\n" +
" ,'format' = 'json'\n" +
" ,'json.timestamp-format.standard' = 'SQL'\n" +
" )";
return kafkaSource1;
}
02 從 kafka 中讀取 topic2
public String createKafkaTable2(){
String kafkaSource2 = "CREATE TABLE source2\n" +
"(\n" +
" id INT,\n" +
" name STRING,\n" +
" address string,\n" +
" dt String,\n" +
" PROCTIME AS PROCTIME()\n" +
") WITH (\n" +
" 'connector' = 'kafka'\n" +
" ,'topic' = 'join_topic2'\n" +
" ,'properties.bootstrap.servers' = 'localhost:9092'\n" +
" ,'scan.startup.mode' = 'earliest-offset'\n" +
" ,'format' = 'json'\n" +
" ,'json.timestamp-format.standard' = 'SQL'\n" +
" )";
return kafkaSource2;
}
● 執行插入邏輯1
String insertSQL = "insert into hudi_catalog.flink_db.test_hudi_flink_join_2(id,name,age,dt) " +
"select id, name,age,dt from source1";
● 通過 spark 查詢資料
20230323090605515 20230323090605515_1_186 45 1 c990a618-896c-4627-8243-baace65c7ad6-0_0-21-26_20230331101342388.parquet 45 xc 45 NULL 1
20230323090605515 20230323090605515_1_179 30 1 c990a618-896c-4627-8243-baace65c7ad6-0_0-21-26_20230331101342388.parquet 30 xc 30 NULL 1
● 執行插入邏輯2
String insertSQL = "insert into hudi_catalog.flink_db.test_hudi_flink_join_2(id,name,address,dt) " +
"select id, name, address,dt from source2";
● 運行成功
運行成功后在 spark 中查詢對應的表資料:
20230323090605515 20230323090605515_1_186 45 1 c990a618-896c-4627-8243-baace65c7ad6-0_0-21-26_20230331101342388.parquet 45 xc 45 xc:address45 1
20230323090605515 20230323090605515_1_179 30 1 c990a618-896c-4627-8243-baace65c7ad6-0_0-21-26_20230331101342388.parquet 30 xc 30 xc:address30 1
可以發現在第二次資料運行之后,表資料的對應欄位 address 已經更新,達到了類似在 Flink 中直接執行 join 的效果,
`insert into hudi_catalog.flink_db.test_hudi_flink_join_2
select a.id, a.name, a.age,b.address a.dt from source1 a left join source2 b on a.id = b.id `
《數堆疊產品白皮書》:https://www.dtstack.com/resources/1004?src=https://www.cnblogs.com/DTinsight/archive/2023/05/24/szsm
《資料治理行業實踐白皮書》下載地址:https://www.dtstack.com/resources/1001?src=https://www.cnblogs.com/DTinsight/archive/2023/05/24/szsm
想了解或咨詢更多有關袋鼠云大資料產品、行業解決方案、客戶案例的朋友,瀏覽袋鼠云官網:https://www.dtstack.com/?src=https://www.cnblogs.com/DTinsight/archive/2023/05/24/szbky
同時,歡迎對大資料開源專案有興趣的同學加入「袋鼠云開源框架釘釘技術qun」,交流最新開源技術資訊,qun號碼:30537511,專案地址:https://github.com/DTStack
轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/553362.html
標籤:其他
上一篇:Elasticsearch與Clickhouse資料存盤對比
下一篇:返回列表
