主頁 > 資料庫 > Apache Hudi 在袋鼠云資料湖平臺的設計與實踐

Apache Hudi 在袋鼠云資料湖平臺的設計與實踐

2023-05-25 09:41:05 資料庫

在大資料處理中,實時資料分析是一個重要的需求,隨著資料量的不斷增長,對于實時分析的挑戰也在不斷加大,傳統的批處理方式已經不能滿足實時資料處理的需求,需要一種更加高效的技術來解決這個問題,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 等,

file

Hudi 的使用場景

Apache Hudi 可以幫助企業和組織實作實時資料處理和分析,實時資料處理需要快速地處理和查詢資料,同時還需要保證資料的一致性和可靠性,

Apache Hudi 的增量資料處理、ACID 事務性保證、寫時合并等技術特性可以幫助企業更好地實作實時資料處理和分析,基于 Hudi 的特性可以在一定程度上在實時數倉的構建程序中承擔上下游資料鏈路的對接(類似 Kafka 的角色),既能實作增量的資料處理,也能為批流一體的處理提供存盤基礎,

Hudi 的優勢和劣勢

● 優勢

· 高效處理大規模資料集;

· 支持實時資料更新和查詢;

· 實作了增量寫入機制,提高了資料訪問效率;

· Hudi 可以與流處理管道集成;

· Hudi 提供了時間旅行功能,允許回溯資料的歷史版本,

● 劣勢

· 在讀寫資料時需要付出額外的代價;

· 操作比較復雜,需要使用專業的編程語言和工具,

Hudi 在袋鼠云資料湖平臺上的實踐

Hudi 在袋鼠云資料湖的技術架構

Hudi 在袋鼠云的資料湖平臺上主要對資料湖管理提供助力:

· 元資料的接入,讓用戶可以快速的對表進行管理;

· 資料快速接入,包括對符合條件的原有表資料進行轉換,快速搭建資料湖能力;

· 湖表的管理,監控小檔案定期進行合并,提升表的查詢性能,內在豐富的表操作功能,包括 time travel ,孤兒檔案清理,過期快照清理等;

· 索引構建,提供多種索引包括 bloom filter,zorder 等,提升計算引擎的查詢性能,

file

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資料存盤對比

下一篇:返回列表

標籤雲
其他(159663) Python(38169) JavaScript(25450) Java(18123) C(15231) 區塊鏈(8268) C#(7972) AI(7469) 爪哇(7425) MySQL(7211) html(6777) 基礎類(6313) sql(6102) 熊猫(6058) PHP(5873) 数组(5741) R(5409) Linux(5340) 反应(5209) 腳本語言(PerlPython)(5129) 非技術區(4971) Android(4576) 数据框(4311) css(4259) 节点.js(4032) C語言(3288) json(3245) 列表(3129) 扑(3119) C++語言(3117) 安卓(2998) 打字稿(2995) VBA(2789) Java相關(2746) 疑難問題(2699) 细绳(2522) 單片機工控(2479) iOS(2433) ASP.NET(2403) MongoDB(2323) 麻木的(2285) 正则表达式(2254) 字典(2211) 循环(2198) 迅速(2185) 擅长(2169) 镖(2155) .NET技术(1976) 功能(1967) Web開發(1951) HtmlCss(1944) C++(1922) python-3.x(1918) 弹簧靴(1913) xml(1889) PostgreSQL(1878) .NETCore(1861) 谷歌表格(1846) Unity3D(1843) for循环(1842)

熱門瀏覽
  • GPU虛擬機創建時間深度優化

    **?桔妹導讀:**GPU虛擬機實體創建速度慢是公有云面臨的普遍問題,由于通常情況下創建虛擬機屬于低頻操作而未引起業界的重視,實際生產中還是存在對GPU實體創建時間有苛刻要求的業務場景。本文將介紹滴滴云在解決該問題時的思路、方法、并展示最終的優化成果。 從公有云服務商那里購買過虛擬主機的資深用戶,一 ......

    uj5u.com 2020-09-10 06:09:13 more
  • 可編程網卡芯片在滴滴云網路的應用實踐

    **?桔妹導讀:**隨著云規模不斷擴大以及業務層面對延遲、帶寬的要求越來越高,采用DPDK 加速網路報文處理的方式在橫向縱向擴展都出現了局限性。可編程芯片成為業界熱點。本文主要講述了可編程網卡芯片在滴滴云網路中的應用實踐,遇到的問題、帶來的收益以及開源社區貢獻。 #1. 資料中心面臨的問題 隨著滴滴 ......

    uj5u.com 2020-09-10 06:10:21 more
  • 滴滴資料通道服務演進之路

    **?桔妹導讀:**滴滴資料通道引擎承載著全公司的資料同步,為下游實時和離線場景提供了必不可少的源資料。隨著任務量的不斷增加,資料通道的整體架構也隨之發生改變。本文介紹了滴滴資料通道的發展歷程,遇到的問題以及今后的規劃。 #1. 背景 資料,對于任何一家互聯網公司來說都是非常重要的資產,公司的大資料 ......

    uj5u.com 2020-09-10 06:11:05 more
  • 滴滴AI Labs斬獲國際機器翻譯大賽中譯英方向世界第三

    **桔妹導讀:**深耕人工智能領域,致力于探索AI讓出行更美好的滴滴AI Labs再次斬獲國際大獎,這次獲獎的專案是什么呢?一起來看看詳細報道吧! 近日,由國際計算語言學協會ACL(The Association for Computational Linguistics)舉辦的世界最具影響力的機器 ......

    uj5u.com 2020-09-10 06:11:29 more
  • MPP (Massively Parallel Processing)大規模并行處理

    1、什么是mpp? MPP (Massively Parallel Processing),即大規模并行處理,在資料庫非共享集群中,每個節點都有獨立的磁盤存盤系統和記憶體系統,業務資料根據資料庫模型和應用特點劃分到各個節點上,每臺資料節點通過專用網路或者商業通用網路互相連接,彼此協同計算,作為整體提供 ......

    uj5u.com 2020-09-10 06:11:41 more
  • 滴滴資料倉庫指標體系建設實踐

    **桔妹導讀:**指標體系是什么?如何使用OSM模型和AARRR模型搭建指標體系?如何統一流程、規范化、工具化管理指標體系?本文會對建設的方法論結合滴滴資料指標體系建設實踐進行解答分析。 #1. 什么是指標體系 ##1.1 指標體系定義 指標體系是將零散單點的具有相互聯系的指標,系統化的組織起來,通 ......

    uj5u.com 2020-09-10 06:12:52 more
  • 單表千萬行資料庫 LIKE 搜索優化手記

    我們經常在資料庫中使用 LIKE 運算子來完成對資料的模糊搜索,LIKE 運算子用于在 WHERE 子句中搜索列中的指定模式。 如果需要查找客戶表中所有姓氏是“張”的資料,可以使用下面的 SQL 陳述句: SELECT * FROM Customer WHERE Name LIKE '張%' 如果需要 ......

    uj5u.com 2020-09-10 06:13:25 more
  • 滴滴Ceph分布式存盤系統優化之鎖優化

    **桔妹導讀:**Ceph是國際知名的開源分布式存盤系統,在工業界和學術界都有著重要的影響。Ceph的架構和演算法設計發表在國際系統領域頂級會議OSDI、SOSP、SC等上。Ceph社區得到Red Hat、SUSE、Intel等大公司的大力支持。Ceph是國際云計算領域應用最廣泛的開源分布式存盤系統, ......

    uj5u.com 2020-09-10 06:14:51 more
  • es~通過ElasticsearchTemplate進行聚合~嵌套聚合

    之前寫過《es~通過ElasticsearchTemplate進行聚合操作》的文章,這一次主要寫一個嵌套的聚合,例如先對sex集合,再對desc聚合,最后再對age求和,共三層嵌套。 Aggregations的部分特性類似于SQL語言中的group by,avg,sum等函式,Aggregation ......

    uj5u.com 2020-09-10 06:14:59 more
  • 爬蟲日志監控 -- Elastc Stack(ELK)部署

    傻瓜式部署,只需替換IP與用戶 導讀: 現ELK四大組件分別為:Elasticsearch(核心)、logstash(處理)、filebeat(采集)、kibana(可視化) 下載均在https://www.elastic.co/cn/downloads/下tar包,各組件版本最好一致,配合fdm會 ......

    uj5u.com 2020-09-10 06:15:05 more
最新发布
  • Apache Hudi 在袋鼠云資料湖平臺的設計與實踐

    在大資料處理中,[實時資料分析](https://www.dtstack.com/dtengine/easylake?src=https://www.cnblogs.com/DTinsight/archive/2023/05/24/szsm)是一個重要的需求。隨著資料量的不斷增長,對于實時分析的挑戰也在不斷加大,傳統的批處理方式已經不能滿足[實時資料處理](https://www.dtstack.com ......

    uj5u.com 2023-05-25 09:41:05 more
  • Elasticsearch與Clickhouse資料存盤對比

    Elasticsearch的查詢陳述句維護成本較高、在聚合計算場景下出現資料不精確等問題。Clickhouse是列式資料庫,列式型資料庫天然適合OLAP場景,類似SQL語法降低開發和學習成本,采用快速壓縮演算法節省存盤成本,采用向量執行引擎技術大幅縮減計算耗時。所以做此對比,進行Elasticsearc... ......

    uj5u.com 2023-05-25 09:40:45 more
  • 【資料庫】時區及JDBC的時區設定

    JDBC連接時有個TimeZone配置,這玩意到底有用嗎?我是使用Postgresql和Mysql兩個資料庫驗證的。結果如下: 資料庫 部署方式 版本 JDBC連接TimeZone引數 JDBC連接serverTimezone引數 總結 Mysql docker 8.0 沒用 有用,會使用客戶端時區 ......

    uj5u.com 2023-05-25 09:30:15 more
  • es筆記六之聚合操作之指標聚合

    > 本文首發于公眾號:Hunter后端 > 原文鏈接:[es筆記六之聚合操作之指標聚合](https://mp.weixin.qq.com/s/UyiZ2bzFxi7zCGmL1Xf3CQ) 聚合操作,在 es 中的聚合可以分為大概四種聚合: * bucketing(桶聚合) * mertic(指標 ......

    uj5u.com 2023-05-25 09:23:32 more
  • Elasticsearch與Clickhouse資料存盤對比

    Elasticsearch的查詢陳述句維護成本較高、在聚合計算場景下出現資料不精確等問題。Clickhouse是列式資料庫,列式型資料庫天然適合OLAP場景,類似SQL語法降低開發和學習成本,采用快速壓縮演算法節省存盤成本,采用向量執行引擎技術大幅縮減計算耗時。所以做此對比,進行Elasticsearc... ......

    uj5u.com 2023-05-25 09:23:21 more
  • 與世界分享我剛編的mysql http隧道工具-hersql原理與使用

    原文地址:[https://blog.fanscore.cn/a/53/](https://blog.fanscore.cn/a/53/) # 1. 前言 本文是[與世界分享我剛編的轉發ntunnel_mysql.php的工具](https://blog.fanscore.cn/a/47/)的后續, ......

    uj5u.com 2023-05-25 09:22:35 more
  • 01_MySQL基礎架構

    01_MySQL基礎架構 MySQL 45 講Note: 課程專欄名稱:《MySQL實戰45講》課程 筆記參考:MYSQL45 講 01_基礎架構:一條SQL查詢陳述句是如何執行的? 一條SQL查詢是如何執行的 先看一下下面這個圖 ?? 我們首先理解一下 Mysql 的基礎架構,理解如果執行一條簡單的 ......

    uj5u.com 2023-05-25 09:22:24 more
  • 【資料庫】時區及JDBC的時區設定

    JDBC連接時有個TimeZone配置,這玩意到底有用嗎?我是使用Postgresql和Mysql兩個資料庫驗證的。結果如下: 資料庫 部署方式 版本 JDBC連接TimeZone引數 JDBC連接serverTimezone引數 總結 Mysql docker 8.0 沒用 有用,會使用客戶端時區 ......

    uj5u.com 2023-05-25 09:22:15 more
  • 150萬學術名詞中英對照字典ACCESS資料庫

    今天這個資料是一款字典的型別的軟體,專門用來查詢一些學術上面的名詞的中英對照,超過180個學科分類,150多萬條記錄,伴隨您悠游于學海之中,是您做學問、寫論文的好幫手。 主要科目有:電子計算機名詞(107213)、電機工程名詞(100395)、電力工程(68379)、外國地名譯名(64487)、機械 ......

    uj5u.com 2023-05-25 09:21:56 more
  • Apache Hudi 在袋鼠云資料湖平臺的設計與實踐

    在大資料處理中,[實時資料分析](https://www.dtstack.com/dtengine/easylake?src=https://www.cnblogs.com/DTinsight/p/szsm)是一個重要的需求。隨著資料量的不斷增長,對于實時分析的挑戰也在不斷加大,傳統的批處理方式已經不能滿足[實時資料處理](https://www.dtstack.com ......

    uj5u.com 2023-05-25 09:21:46 more