全部產品
Search
文件中心

Data Lake Formation:EMR on ECS Spark訪問DLF

更新時間:Jun 03, 2026

如何在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許可權

  1. 授予AliyunECSInstanceForEMRRole角色RAM許可權(EMR產品化整合後可以省略該步驟)。

    1. 使用阿里雲帳號或Resource Access Management員登入RAM控制台

    2. 單擊身份管理 > 角色,查詢AliyunECSInstanceForEMRRole角色。

    3. 單擊操作列的新增授權,進入新增授權頁面。

    4. 權限原則中,查詢並勾選AliyunDLFFullAccess,單擊確認新增授權

  2. 授予AliyunECSInstanceForEMRRole角色DLF許可權。

    1. 登入資料湖構建控制台

    2. Catalogs列表頁面,單擊Catalog名稱,進入Catalog詳情頁。

    3. 如果要授予整個 catalog 許可權,則直接單擊許可權。否則點擊進入對應資料庫或者表,再點擊許可權目錄,進行許可權授予。

    4. 在授權頁面,配置以下資訊,單擊確定

      • 使用者/角色:選擇RAM使用者/RAM角色

      • 選擇授權對象:在下拉式清單中選擇AliyunECSInstanceForEMRRole

        說明

        如果使用者下拉式清單中未找到AliyunECSInstanceForEMRRole,可以在使用者管理頁面單擊同步。

      • 預置權限類別型:可自訂選擇讀取許可權,或者直接使用 Data Reader/Data Editor。

升級EMR叢集Paimon依賴

請到Maven倉庫下載1.1+ 版本的兩個JAR:paimon-jindo-*.jar, paimon-spark-3.x-*.jar,並根據EMR叢集的Spark版本選取對應依賴。

  1. 匯入Paimon依賴。

    1. 將依賴的兩個JAR包paimon-jindo-*.jar, paimon-spark-3.x-*.jar上傳至OSS,並設定檔案讀寫權限為公用讀取。請參見簡單上傳

    2. 將以下指令碼修改後上傳至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

  2. 通過EMR叢集引導指令碼執行。詳情請參見手動執行指令碼

    1. 在EMR叢集中,選擇指令碼操作 > 手動執行頁簽,單擊建立並執行

    2. 在彈出的對話方塊中,配置以下資訊,單擊確定

      • 名稱:自訂指令碼名稱。

      • 指令碼位置:選擇上傳到OSS的升級指令碼。指令碼路徑格式必須是oss://**/*.sh格式。

      • 執行範圍:選擇叢集

  3. 執行完成後,需重啟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