本文示範了使用Realtime ComputeFlink版和EMR Serverless Spark構建Paimon資料湖分析流程。該流程包括將資料寫入OSS、進行互動式查詢以及執行離線資料Compact操作。EMR Serverless Spark完全相容Paimon,通過內建的DLF中繼資料與其他雲產品(例如,Realtime ComputeFlink版)實現中繼資料互連,形成完整的流批一體化解決方案。它支援靈活的任務運行方式和參數配置,滿足即時分析和生產調度的多種需求。
背景資訊
Realtime ComputeFlink版
阿里雲Realtime ComputeFlink版是一種全託管Serverless的Flink雲端服務,是一站式開發營運管理平台,開箱即用,計費靈活。具備作業開發、資料調試、運行與監控、自動調優、智能診斷等全生命週期能力。更多資訊,請參見什麼是阿里雲Realtime ComputeFlink版。
Apache Paimon
Apache Paimon是一種統一的資料湖儲存格式,結合Flink和Spark構建了流批處理的即時湖倉一體架構。Paimon創新地將湖格式與LSM(Log-structured merge-tree)技術結合,使資料湖具備了即時資料流更新和完整的流處理能力。更多資訊,請參見Apache Paimon。
操作流程
步驟一:通過Realtime ComputeFlink建立Paimon Catalog
Paimon Catalog可以方便地管理同一個warehouse目錄下的所有Paimon表,並與其它阿里雲產品連通。建立並使用Paimon Catalog,詳情請參見管理Paimon Catalog。
-
單擊目標工作空間操作列下的控制台。
-
建立Paimon Catalog。
通過UI方法添加DLF類型的Paimon Catalog,無需手動填寫AccessKey資訊。具體操作,請參見管理Paimon Catalog。
-
建立Paimon表。
在查询脚本文本編輯地區輸入如下命令後,選中代碼後單擊运行。
CREATE TABLE IF NOT EXISTS `paimon`.`test_paimon_db`.`test_append_tbl` ( id STRING, data STRING, category INT, ts STRING, dt STRING, hh STRING ) PARTITIONED BY (dt, hh) WITH ( 'write-only' = 'true' ); -
建立流作業。
-
新增作業。
-
在左側導覽列,選擇数据开发 > ETL。
-
建立流作業,在新增作業草稿對話方塊中,填寫作業配置資訊。
作業參數
說明
文件名称
作業的名稱。
說明作業名稱在當前專案中必須保持唯一。
引擎版本
當前作業使用的Flink的引擎版本。引擎版本號碼含義、版本對應關係和生命週期重要時間點詳情請參見引擎版本介紹。
-
單擊创建。
-
-
編寫代碼。
在建立的作業草稿中,輸入以下代碼,通過datagen源源不斷產生資料寫入Paimon表中。
CREATE TEMPORARY TABLE datagen ( id string, data string, category int ) WITH ( 'connector' = 'datagen', 'rows-per-second' = '100', 'fields.category.kind' = 'random', 'fields.category.min' = '1', 'fields.category.max' = '10' ); INSERT INTO `paimon`.`test_paimon_db`.`test_append_tbl` SELECT id, data, category, cast(LOCALTIMESTAMP as string) as ts, cast(CURRENT_DATE as string) as dt, cast(hour(LOCALTIMESTAMP) as string) as hh FROM datagen; -
單擊部署,即可將資料發布至生產環境。
-
您可以在作業營運頁面啟動作業進入運行階段,詳情請參見作業啟動。
-
步驟二:通過EMR Serverless Spark建立SQL會話
建立的SQL會話用於SQL開發和查詢。有關會話的詳細介紹,請參見會話管理。
-
進入會話管理頁面。
-
在左側導覽列,選擇。
-
在Spark頁面,單擊目標工作空間名稱。
-
在EMR Serverless Spark頁面,單擊左側導覽列中的會話管理。
-
在Serverless Spark工作空間的資料目錄中,添加步驟一中建立的DLF Catalog。
-
建立SQL會話。
-
在SQL 会话頁簽,單擊创建 SQL 会话。
-
在建立SQL會話頁面,配置以下資訊,其餘參數無需配置,然後單擊创建。
參數
說明
名称
自訂SQL會話的名稱。例如,paimon_compute。
-
單擊操作列的启动。
-
步驟三:通過EMR Serverless Spark進行互動式查詢或任務調度
EMR Serverless Spark提供了互動式查詢和任務調度兩種操作模式,以滿足不同的使用需求。互動式查詢適用於快速查詢和調試,而任務調度則支援任務的開發、發布和營運,實現完整的生命週期管理。
在資料寫入過程中,我們可以隨時使用EMR Serverless Spark對Paimon表進行互動式查詢,以便即時擷取資料狀態和執行快速分析。此外,通過發布開發好的任務並建立工作流程,可以編排各項任務並完成工作流程的發布。您可以配置調度策略,實現任務的定期調度,從而保證資料處理和分析的自動化與高效性。
互動式查詢
-
建立SQL開發。
-
在EMR Serverless Spark頁面,單擊左側導覽列中的数据开发。
-
在开发目录頁簽下,單擊新建。
-
在彈出的對話方塊中,輸入名称(例如,paimon_compact),类型選擇為SparkSQL,然後單擊确定。
-
在右上方選擇資料目錄、資料庫和前一步驟中啟動的SQL會話。
-
在建立的任務編輯器中輸入SQL語句。
-
樣本1:查詢
test_append_tbl表中前10行的資料。SELECT * FROM paimon.test_paimon_db.test_append_tbl limit 10;說明請根據實際情況替換SQL中的Catalog名稱、資料庫名稱和表名稱。
返回結果包含 id、data、category、ts、dt、hh 等欄位,其中 id 和 data 列為雜湊字串,category 為數字分類值,ts 為完整時間戳記(如
2024-06-24 19:00:00.446),dt 為日期,hh 為小時。 -
樣本2:統計
test_append_tbl表中滿足特定條件的行數。SELECT COUNT(*) FROM paimon.test_paimon_db.test_append_tbl WHERE dt = '2024-06-24' AND hh = '19';說明請根據實際情況替換SQL中的Catalog名稱、資料庫名稱、表名稱和分區條件(dt、hh)。
查詢返回一行結果,
count(1)值為360000。
-
-
-
運行並發布任務。
-
單擊运行。
返回結果資訊可以在下方的运行结果中查看。如果有異常,則可以在运行问题中查看。
-
確認運行無誤後,單擊右上方的发布。
-
在发布對話方塊中,可以輸入發布資訊,然後單擊确定。
-
任務調度
-
查詢Compact前檔案資訊。
在数据开发頁面,建立SQL開發,查詢Paimon的files系統資料表,快速地得到Compact前檔案的資料。建立SQL開發的具體操作,請參見SparkSQL開發。
SELECT file_path, record_count, file_size_in_bytes FROM paimon.test_paimon_db.`test_append_tbl$files`;說明請根據實際情況替換SQL中的Catalog名稱和資料庫名稱。
查詢結果顯示 Compact 前共有 14 個 ORC 資料檔案,各檔案
record_count大多為 18100(部分為 18200,首行為 8700),file_size_in_bytes範圍約為 1752386~3664183。 -
在上述SQL開發(paimon_compact)中輸入以下Compact SQL,然後直接發布。
CALL paimon.sys.compact ( table => 'test_paimon_db.test_append_tbl', partitions => 'dt=\"2024-06-24\",hh=\"19\"', order_strategy => 'zorder', order_by => 'category' ); -
建立工作流程。
-
在EMR Serverless Spark頁面,單擊左側導覽列中的任务编排。
-
在任务编排頁面,單擊创建工作流。
-
在创建工作流面板中,輸入工作流名称(例如,paimon_workflow_task),然後單擊下一步。
其他设置地區的參數,請根據您的實際情況配置,更多參數資訊請參見管理工作流程。
-
在建立的節點畫布中,單擊添加节点。
-
在来源文件路径下拉式清單中選擇發行的SQL開發(paimon_compact),填寫Spark 配置參數,然後單擊保存。
參數
說明
名称
自訂SQL會話的名稱。例如,paimon_compute。
-
在建立的節點畫布中,單擊发布工作流,然後單擊确定。
-
-
運行工作流程。
-
在任务编排頁面,單擊建立工作流程(例如,paimon_workflow_task)的工作流名称。
-
在工作流实例列表頁面,單擊手动运行。
-
在触发运行對話方塊中,單擊确定。
-
-
驗證Compact效果。
工作流程調度執行成功後,再次執行與開始相同的SQL查詢,對比Compact前後檔案的數量、記錄數和大小,以驗證Compact操作的效果。
SELECT file_path, record_count, file_size_in_bytes FROM paimon.test_paimon_db.`test_append_tbl$files`;Compact 後查詢結果顯示錶中僅剩 3 個 ORC 檔案:
data-c971b2a7-…-0.orc(record_count 144591,file_size_in_bytes 29122164)、data-17d0b4c6-…-0.orc(record_count 179571,file_size_in_bytes 36159348)、data-a8f27c6b-…-0.orc(record_count 35838,file_size_in_bytes 7243404),說明小檔案已被合并。