如何在EMR on ECS Spark環境中通過Paimon REST訪問DLF Catalog。
前提條件
-
已建立EMR叢集且版本 >= 5.12.0,組件選擇Spark3, Paimon。如有其他版本訴求,請加入DingTalk群(106575000021)聯絡DLF研發人員。
-
已建立快速使用DLF。
-
EMR與DLF在同一地區,且添加EMR叢集所在的VPC到DLF的白名單中。
建立DLF Catalog
詳情請參見 DLF 快速入門 。
授予角色DLF許可權
-
授予AliyunECSInstanceForEMRRole角色RAM許可權(EMR產品化整合後可以省略該步驟)。
-
使用阿里雲帳號或Resource Access Management員登入RAM控制台。
-
單擊,查詢AliyunECSInstanceForEMRRole角色。
-
單擊操作列的新增授權,進入新增授權頁面。
-
在權限原則中,查詢並勾選AliyunDLFFullAccess,單擊確認新增授權。
-
-
授予AliyunECSInstanceForEMRRole角色DLF許可權。
登入資料湖構建控制台。
-
在Catalogs列表頁面,單擊Catalog名稱,進入Catalog詳情頁。
-
如果要授予整個 catalog 許可權,則直接單擊許可權。否則點擊進入對應資料庫或者表,再點擊許可權目錄,進行許可權授予。
-
在授權頁面,配置以下資訊,單擊確定。
-
使用者/角色:選擇RAM使用者/RAM角色。
-
選擇授權對象:在下拉式清單中選擇AliyunECSInstanceForEMRRole。
說明如果使用者下拉式清單中未找到AliyunECSInstanceForEMRRole,可以在使用者管理頁面單擊同步。
-
預置權限類別型:可自訂選擇讀取許可權,或者直接使用 Data Reader/Data Editor。
-
升級EMR叢集Paimon依賴
請到Maven倉庫下載1.1+ 版本的兩個JAR:paimon-jindo-*.jar, paimon-spark-3.x-*.jar,並根據EMR叢集的Spark版本選取對應依賴。
-
匯入Paimon依賴。
-
將依賴的兩個JAR包
paimon-jindo-*.jar,paimon-spark-3.x-*.jar上傳至OSS,並設定檔案讀寫權限為公用讀取。請參見簡單上傳。 -
將以下指令碼修改後上傳至OSS。
#!/bin/bash echo 'clean up paimon-dlf-2.5 exists file' rm -rf /opt/apps/PAIMON/paimon-dlf-2.5 rm -rf /opt/apps/PAIMON/paimon-dlf-2.5.tar.gz.* cd /opt/apps/PAIMON/paimon-current/lib/spark3 mkdir -p /opt/apps/PAIMON/paimon-dlf-2.5/lib/spark3 cd /opt/apps/PAIMON/paimon-dlf-2.5/lib/spark3 wget ${paimon-jindo-1.1.0.jar} wget ${paimon-spark-3.x-1.1.0.jar} echo 'link paimon-current to paimon-dlf-2.5' rm -f /opt/apps/PAIMON/paimon-current ln -sf /opt/apps/PAIMON/paimon-dlf-2.5 /opt/apps/PAIMON/paimon-current重要指令碼中的預留位置
${paimon-jindo-1.1.0.jar}和${paimon-spark-3.x-1.1.0.jar}需替換為對應的OSS可下載路徑。EMR on ECS預設開通的叢集不具備訪問公網的能力。-
內網:
https://{bucket}.oss-cn-hangzhou-internal.aliyuncs.com/jars/paimon-jindo-1.1.0.jar。 -
公網:
https://{bucket}.oss-cn-hangzhou.aliyuncs.com/jars/paimon-jindo-1.1.0.jar。
-
-
-
通過EMR叢集引導指令碼執行。詳情請參見手動執行指令碼。
-
在EMR叢集中,選擇頁簽,單擊建立並執行。
-
在彈出的對話方塊中,配置以下資訊,單擊確定。
-
名稱:自訂指令碼名稱。
-
指令碼位置:選擇上傳到OSS的升級指令碼。指令碼路徑格式必須是oss://**/*.sh格式。
-
執行範圍:選擇叢集。
-
-
-
執行完成後,需重啟Spark服務以生效。
Spark讀寫資料
串連Paimon Catalog
在Terminal中執行以下spark-sql命令。
需替換命令中的${regionID}為實際region,如cn-hangzhou;${catalog}替換成在DLF中建立好的catalog名稱。
spark-sql --master yarn \
--conf spark.driver.memory=5g \
--conf spark.sql.defaultCatalog=paimon \
--conf spark.sql.catalog.paimon=org.apache.paimon.spark.SparkCatalog \
--conf spark.sql.catalog.paimon.metastore=rest \
--conf spark.sql.extensions=org.apache.paimon.spark.extensions.PaimonSparkSessionExtensions \
--conf spark.sql.catalog.paimon.uri=http://${regionID}-vpc.dlf.aliyuncs.com \
--conf spark.sql.catalog.paimon.warehouse=${catalog} \
--conf spark.sql.catalog.paimon.token.provider=dlf \
--conf spark.sql.catalog.paimon.dlf.token-loader=ecs
建立資料表
執行以下SQL,建立資料表。
CREATE TABLE user_samples
(
user_id INT,
age INT,
gender_code STRING,
clk BOOLEAN
);
CREATE TABLE user_samples_di (
user_id INT,
age INT,
gender_code STRING,
clk BOOLEAN
)
USING CSV
OPTIONS(
'path'='oss://${bucket}/user/user_samples_di'
);
-
不指定資料庫時,建立資料表會預設建在Catalog下的
default資料庫中,也可建立並指定其他資料庫。 -
目錄
/user/user_samples需要提前在OSS中建立。指定path時,為建立外部表格。表的中繼資料會被儲存在DLF中管理,實際的資料檔案仍然儲存在指定的OSS路徑下。刪除該表只刪除中繼資料,而不會刪除OSS上的未經處理資料檔案。
插入資料
運行以下SQL,插入資料。
INSERT INTO user_samples VALUES
(1, 25, 'M', true),
(2, 18, 'F', false);
INSERT INTO user_samples_di VALUES
(1, 25, 'M', true),
(2, 18, 'F', true),
(3, 35, 'M', true);
查詢資料
運行以下SQL,查詢資料。
SELECT * FROM user_samples;
SELECT * FROM user_samples_di;
結果如下。
spark-sql (default)> SELECT * FROM user_samples;
1 25 M true
2 18 F false
spark-sql (default)> SELECT * FROM user_samples_di;
1 25 M true
2 18 F true
3 35 M true
合并資料
使用上述user_samples_di表merge into到user_samples:
MERGE INTO user_samples
USING user_samples_di
ON user_samples.user_id = user_samples_di.user_id
WHEN MATCHED THEN
UPDATE SET
age = user_samples_di.age,
gender_code = user_samples_di.gender_code,
clk = user_samples_di.clk
WHEN NOT MATCHED THEN
INSERT (user_id, age, gender_code, clk)
VALUES (user_samples_di.user_id, user_samples_di.age, user_samples_di.gender_code, user_samples_di.clk);
此時,users_samples 表中的資料將被 user_samples_di 表中具有相同 user_id 的資料覆蓋,同時會插入 user_samples 表中沒有 user_id 的新資料。
spark-sql (default)> SELECT * FROM user_samples;
1 25 M true
2 18 F true
3 35 M true