全部產品
Search
文件中心

Realtime Compute for Apache Flink:Paimon connector

更新時間:Aug 13, 2026

本主題說明如何將 Paimon connector 用於串流資料湖倉庫。為獲得最佳效果,建議將 connector 與 Paimon Catalog 搭配使用。

背景資訊

Apache Paimon 是一種統一的串流與批次資料湖儲存格式,支援高吞吐量寫入與低延遲查詢。Paimon 可與 Alibaba Cloud E-MapReduce 中的主流計算引擎(如 Flink、Spark、Hive 和 Trino)良好整合。您可以使用 Apache Paimon 在 HDFS 或 OSS 上快速建立資料湖,並連接這些計算引擎進行資料湖分析。如需詳細資訊,請參見 Apache Paimon.

類別

描述

支援的類型

來源表、維度表、結果表,以及資料接入目標

執行模式

串流模式與批次模式

資料格式

不支援

監控指標

否ne

API 類型

SQL 與 YAML(用於資料接入)

結果表的更新與刪除

核心功能

Apache Paimon 提供以下核心能力:

  • 在 HDFS 或物件儲存上建立輕量、低成本的資料湖。

  • 以串流和批次模式讀取及寫入大規模資料集。

  • 執行批次和 OLAP 查詢,資料新鮮度可達分鐘乃至秒級。

  • 接入並產生增量資料,可作為傳統離線和現代串流資料倉儲的儲存層。

  • 對資料進行預聚合,以降低儲存成本及下游計算負載。

  • 存取資料的歷史版本。

  • 高效過濾資料。

  • 支援 schema 演進。

限制與建議

  • Paimon connector 需要 Flink 計算引擎 VVR 6.0.6 或更高版本。

  • 下表列出了 Paimon 與 VVR 之間的版本相容性。

    Apache Paimon 版本

    VVR

    1.3.1

    11.5, 11.6, 11.7, 11.8

    1.3

    11.4

    1.2

    11.2, 11.3

    1.1

    11.1

    1.0

    8.0.11

  • 並發寫入的儲存建議

    當多個作業並發寫入同一 Paimon 表時,使用標準 OSS 儲存(oss://)偶爾會因原子檔案操作的限制而導致提交衝突或作業失敗。

    為實現穩定一致的寫入,請使用提供強原子保證的元資料或儲存服務。首選方案是 Data Lake Formation(DLF),它提供 Paimon 元資料和儲存的統一管理。或者,您也可以使用 OSS-HDFS 或 HDFS。

  • 配置變更如何生效

    Paimon 表的配置參數變更只有在重啟相關作業後才會生效。正在執行的作業不會動態載入這些變更。

  • 已刪除分區的延遲物理回收

    當您執行 DROP PARTITION 操作時,系統 不會立即刪除 底層的物理資料檔案。
    此操作執行的是邏輯刪除。Paimon 僅從最新快照中移除目標分區的元資料。由於 Paimon 支援時間旅行(time travel)功能,歷史快照仍然引用該分區的資料檔案。只有當引用該分區的所有歷史快照達到其保留限制並被快照過期機制清理後,物理資料檔案才會被永久刪除。

SQL

在 SQL 作業中使用 Paimon connector 作為來源表或結果表。

語法

  • 如果您在 Paimon Catalog中建立 Paimon 表,則無需指定 connector 參數。語法如下:

    CREATE TABLE `<YOUR-PAIMON-CATALOG>`.`<YOUR-DB>`.paimon_table (
      id BIGINT,
      data STRING,
      PRIMARY KEY (id) NOT ENFORCED
    ) WITH (
      ...
    );
    說明

    如果您已在 Paimon Catalog 中建立了 Paimon 表,可直接使用。

  • 如果您在其他 Catalog 中建立 Paimon 臨時表,則必須指定 'connector' 和 'path' 參數。語法如下:

    CREATE TEMPORARY TABLE paimon_table (
      id BIGINT,
      data STRING,
      PRIMARY KEY (id) NOT ENFORCED
    ) WITH (
      'connector' = 'paimon',
      'path' = '<path-to-paimon-table-files>',
      'auto-create' = 'true', -- If Paimon table data files do not exist at the specified path, they are automatically created.
      ...
    );
    說明
    • Path example: 'path' = 'oss://<bucket>/test/order.db/orders'. Do not omit the .db 後綴。Paimon 依賴此後綴來識別資料庫。

    • 寫入同一個表的多個作業必須共用相同的路徑配置。

    • 如果兩個路徑配置不同,Paimon 不會將其識別為同一個表。即使物理路徑相同,不一致的 Catalog 配置也可能導致並發寫入衝突、compaction 操作失敗和資料遺失。例如,Paimon 認為 oss://b/test and oss://b/test/ 是不同的表(因為尾部斜線的差異),儘管它們可能指向相同的物理位置。

WITH 參數

參數

描述

Type

是否必填

預設值

備註

connector

指定表的 connector。

String

否ne

  • 如果您在 Paimon catalog 中建立 Paimon 表,則不需要此參數。

  • 如果您在其他 catalog 中建立 Paimon 臨時表,此參數必須設定為 paimon.

path

表的儲存路徑。

String

否ne

  • 如果您在 Paimon catalog 中建立 Paimon 表,則不需要此參數。

  • 如果您在其他 catalog 中建立 Paimon 臨時表,此參數指定表在 HDFS 或 OSS 中的儲存目錄。

auto-create

指定當指定路徑下不存在表檔案時,是否自動建立。

Boolean

false

有效值:

  • false(預設值):如果指定路徑下不存在 Paimon 表檔案,作業將失敗。

  • true:如果指定路徑不存在,Flink 將自動建立 Paimon 表檔案。

file.format

資料檔案的格式。

String

parquet

有效值:

  • orc

  • parquet

  • avro

  • lance(Realtime Compute for Apache Flink 11.6 及更高版本支援)

bucket

每個分區的 bucket 數量。

Integer

1

Paimon 根據 bucket-key.

說明

建議每個 bucket 包含的資料少於 5 GB。

bucket-key

用作 bucket key 的欄位。

String

否ne

指定用於將資料分配到 bucket 的欄位。

多個欄位名稱用逗號(,)分隔。例如, 'bucket-key' = 'order_id,cust_id' 根據 order_id and cust_id 欄位分配資料。

說明
  • If this parameter is not specified, Paimon 根據 primary key.

  • If the table has no primary key, Paimon 根據 values of all 欄位分配資料。

changelog-producer

changelog 的生產機制。

String

none

Paimon 可以為任何輸入串流生成完整的 changelog,這意味著每條 update_after 記錄都有對應的 update_before 記錄。這簡化了下游消費。有效值:

  • none(預設值):不生成額外的 changelog。下游消費者仍然可以以串流模式讀取 Paimon 表,但 changelog 不完整(只包含 update_after 記錄,沒有對應的 update_before 記錄)。

  • input:將輸入串流寫入 changelog 檔案,作為完整的 changelog。

  • full-compaction:在每次完整 compaction 時生成完整的 changelog。

  • lookup:在每個 snapshot 提交前生成完整的 changelog。

有關如何選擇 changelog 生產者的更多資訊,請參見 Changelog 生產.

full-compaction.delta-commits

兩次連續完整 compaction 之間的最大提交次數。

Integer

否ne

指定在觸發完整 compaction 之前允許的最大 snapshot 提交次數。

lookup.cache-max-memory-size

Paimon 維度表的記憶體快取大小。

String

256 MB

此參數控制維度表查找和 lookup changelog 生產者的快取大小。

merge-engine

合併具有相同主鍵的記錄的機制。

String

deduplicate

有效值:

  • deduplicate:僅保留最新記錄。

  • partial-update:使用最新記錄的非空值更新現有記錄。其他欄位保持不變。

  • aggregation:使用指定的聚合函數進行預聚合。

有關合併引擎的詳細分析,請參見 合併引擎.

partial-update.ignore-delete

指定是否忽略刪除(-D)訊息。

Boolean

false

有效值:

  • true:忽略刪除訊息。

  • false:處理刪除訊息。您必須使用如 sequence.field 等參數配置刪除處理策略,以防止潛在的 IllegalStateException or IllegalArgumentException 錯誤。

說明
  • 在 Realtime Compute for Apache Flink 8.0.6 及更早版本中,此參數僅在 merge-engine = 'partial-update'.

  • 在 Realtime Compute for Apache Flink 8.0.7 及更高版本中,此參數也相容於非 partial-update 場景,功能與 ignore-delete 參數相同。建議您使用 ignore-delete 代替。

  • 根據您的業務需求和是否預期收到刪除訊息來決定是否啟用此參數。如果刪除訊息不符合作業的預期語義,快速失敗通常是更好的選擇。

ignore-delete

指定是否忽略刪除(-D)訊息。

Boolean

false

有效值與 partial-update.ignore-delete.

說明
  • 此參數僅在 Realtime Compute for Apache Flink 8.0.7 及更高版本中受支援。

  • 此參數與 partial-update.ignore-delete功能相同。建議您使用 ignore-delete ,並避免同時配置這兩個參數。

partition.default-name

預設分區名稱。

String

__DEFAULT_PARTITION__

當分區欄位的值為 null 或空字串時使用的分區名稱。

partition.expiration-check-interval

系統檢查過期分區的頻率。

String

1h

For details為前綴的參數的資訊,請參見 如何配置自動分區過期.

partition.expiration-time

分區過期的時長。

String

否ne

當分區的存在時間超過此值時,分區將過期。預設情況下,分區永不過期。

系統根據分區值計算分區的存在時間。詳情請參見 如何配置自動分區過期.

partition.timestamp-formatter

將時間字串轉換為時間戳記的格式字串。

String

否ne

指定從分區值中提取分區存在時間的格式。詳情請參見 如何配置自動分區過期.

partition.timestamp-pattern

將分區值轉換為時間字串的格式字串。

String

否ne

指定從分區值中提取時間字串的模式。詳情請參見 如何配置自動分區過期.

scan.bounded.watermark

表示掃描結束的 watermark 值。當來源表的 watermark 超過此值時,停止產出資料。

Long

否ne

N/A

scan.mode

指定 Paimon 來源表的消費位置。

String

default

For details為前綴的參數的資訊,請參見 如何設定 Paimon 來源表的消費位置.

scan.snapshot-id

指定 Paimon 來源表開始消費的 snapshot。

Integer

否ne

For details為前綴的參數的資訊,請參見 如何設定 Paimon 來源表的消費位置.

scan.timestamp-millis

指定 Paimon 來源表開始消費的時間點。

Integer

否ne

For details為前綴的參數的資訊,請參見 如何設定 Paimon 來源表的消費位置.

snapshot.num-retained.max

保留最新 snapshot 的最大數量。

Integer

2147483647

如果此條件或 snapshot.time-retained 條件滿足,且 snapshot.num-retained.min 條件也滿足,則觸發 snapshot 過期。

snapshot.num-retained.min

保留最新 snapshot 的最小數量。

Integer

10

N/A

snapshot.time-retained

snapshot 的保留期。

String

1h

如果此條件或 snapshot.num-retained.max 條件滿足,且 snapshot.num-retained.min 條件也滿足,則觸發 snapshot 過期。

write-mode

Paimon 表的寫入模式。

String

change-log

有效值:

  • change-log:Paimon 表支援基於主鍵的插入、刪除和更新操作。

  • append-only:Paimon 表僅接受插入操作,不支援主鍵。此模式比 change-log 模式更高效。

有關寫入模式的更多資訊,請參見 寫入模式.

scan.infer-parallelism

指定是否自動推斷 Paimon 來源表的並行度。

Boolean

true

有效值:

  • true:根據 bucket 數量自動推斷 Paimon 來源表的並行度。

  • false: Uses the default parallelism configured in Realtime Compute for Apache Flink. If expert mode is enabled, the job uses your configured parallelism 代替。

scan.parallelism

Paimon 來源表的並行度。

Integer

否ne

說明

如果在作業的 Resource Mode is set to Mode on the job's Configuration > Resources 標籤頁,則忽略此參數。

sink.parallelism

Paimon sink 表的並行度。

Integer

否ne

說明

如果在作業的 Resource Mode is set to Mode on the job's Configuration > Resources 標籤頁,則忽略此參數。

sink.clustering.by-columns

指定寫入 Paimon sink 表時的聚簇欄位。

String

否ne

For Paimon append-only 表(無主鍵表), this parameter enables clustered writing in batch jobs. This process improves query performance by grouping data on the specified 欄位分配資料。

多個欄位名稱用逗號(,)分隔,例如 'col1,col2'.

有關聚簇的更多資訊,請參見 Apache Paimon 官方文件.

sink.delete-strategy

指定驗證策略,以確保系統正確處理回縮訊息(-D/-U)。

​​

Enum

NONE

有效值及 sink 算子處理回縮訊息時的預期行為:

  • NONE(預設值):不執行驗證。

  • IGNORE_DELETE:sink 算子必須忽略 -U 和 -D 訊息,不發生回縮。

  • NON_PK_FIELD_TO_NULL:sink 算子必須忽略 -U 訊息。對於 -D 訊息,保留主鍵並將所有其他非主鍵欄位設為 null。

    這主要用於多個 sink 寫入同一表時的部分更新場景。

  • DELETE_ROW_ON_PK:sink 算子必須忽略 -U 訊息,但對於 -D 訊息,刪除主鍵對應的行。

  • CHANGELOG_STANDARD:sink 算子對於 -U 和 -D 訊息都必須刪除主鍵對應的行。

說明
  • 此參數僅在 Realtime Compute for Apache Flink 8.0.8 及更高版本中受支援。

  • 回縮的實際 sink 行為由其他參數決定,例如 ignore-delete and merge-engine。此參數僅驗證實際行為是否與所選策略匹配。如果不匹配,作業將失敗並顯示建議更正的錯誤訊息。

blob-as-descriptor

指定讀取 Blob 欄位時是否輸出 Blob 描述符(descriptor)位元組。

Boolean

false

將此參數設定為 true 時,查詢將返回 BlobDescriptor 的序列化位元組,而不是實際的 Blob 內容。將此參數與 Blob 預簽名 URL 函數配合使用。您無需在建立表時設定此參數,可以在讀取時通過 SQL hint 動態設定。此參數在 VVR 11.9-preview1 及更高版本中受支援。

說明

僅在 VVR 11.9-preview1 及更高版本中受支援。

sink.existing-table.schema-check.enabled

指定向預先建立的表寫入資料時,是否檢查表 schema 的主鍵、欄位數量和欄位類型的相容性。

Boolean

true

僅在 VVR 11.8 及更高版本中受支援。

說明

有關配置選項的更多資訊,請參見 Apache Paimon 官方文件.

向量表參數

重要

以下參數僅在 VVR 11.8 及更高版本中受支援。

參數

描述

資料類型

是否必填

預設值

備註

ivf.nprobe

搜尋時探查的 IVF 叢集數量。

Integer

16

較高的值通常會提升召回率,但也會增加延遲。

ivf.refine_factor

檢索 top_k × refine_factor 個 IVF 候選結果,並使用 Paimon 表中儲存的原始向量重新排序。

Float

停用

所有 IVF 變體預設禁用。最適合在召回率優先於延遲的場景中使用壓縮索引(如 ivf-pq 和 ivf-hnsw-sq)。

hnsw.ef_search

搜尋時的 HNSW 搜尋寬度。

Integer

0

較高的值通常會提升召回率,但也會增加延遲。 0 indicates that the native library default is used.

diskann.search.list_size

Lumina DiskANN 搜尋清單大小。

Integer

max(1.5 × top_k, 16)

較高的值通常會提升召回率,但也會增加延遲。

diskann.search.beam_width

Lumina DiskANN 搜尋束寬度。

Integer

4

search.parallel_number

並行 Lumina 搜尋的數量。

Integer

5

功能詳情

資料新鮮度與一致性

Paimon sink 表使用兩階段提交協議在每次 Flink 作業 checkpoint 時提交資料。因此,資料新鮮度由 Flink 作業的 checkpoint 間隔決定。每次提交最多生成兩個 snapshot。

當兩個 Flink 作業並發寫入同一個 Paimon 表時,如果作業寫入不同的 bucket,它們可以實現可序列化一致性。如果作業寫入同一個 bucket,它們只能實現 snapshot 隔離。這意味著表的資料可能是兩個作業結果的混合,但不會發生資料遺失。

Merge engine

當 Paimon sink 表接收到具有相同主鍵的多條記錄時,它會將它們合併為一條記錄以維護唯一性。您可以通過設定 merge-engine 參數來控制此行為。下表描述了可用的合併引擎。

Merge engine

描述

Deduplicate

去重引擎是預設選項。對於具有相同主鍵的多條記錄,Paimon sink 表只保留最新記錄並丟棄其他記錄。

說明

如果最新記錄是刪除訊息,則丟棄所有具有該主鍵的記錄。

Partial Update

partial update 引擎允許您通過多條訊息的增量更新來建立完整記錄。當具有相同主鍵的新記錄到達時,其非空值覆蓋現有記錄中對應的欄位。引擎忽略新記錄中為空的欄位,並保留現有值。

例如,假設 Paimon sink 表按順序接收以下三條記錄:

  • <1, 23.0, 10, NULL>

  • <1, NULL, NULL, 'This is a book'>

  • <1, 25.2, NULL, NULL>

如果第一欄是主鍵,最終合併的記錄為 <1, 25.2, 10, 'This is a book'>。

說明
  • 要以串流模式讀取 partial update 的結果,您必須將 changelog-producer parameter to lookup or full-compaction.

  • partial update 引擎無法處理刪除訊息。您可以將 partial-update.ignore-delete parameter to true 以忽略刪除訊息。

Aggregation

在某些使用場景中,您可能只需要記錄的聚合值。聚合引擎使用您指定的聚合函數組合具有相同主鍵的記錄。對於每個非主鍵欄位,您必須使用 fields.<field-name>.aggregate-function 選項指定聚合函數。否則,該欄位預設使用 last_non_null_value 聚合函數。例如,考慮以下 Paimon 表定義。

CREATE TABLE MyTable (
  product_id BIGINT,
  price DOUBLE,
  sales BIGINT,
  PRIMARY KEY (product_id) NOT ENFORCED
) WITH (
  'merge-engine' = 'aggregation',
  'fields.price.aggregate-function' = 'max',
  'fields.sales.aggregate-function' = 'sum'
);

The price 欄位使用 max 函數聚合, sales 欄位使用 sum 函數聚合。給定兩條輸入記錄 <1, 23.0, 15> 和 <1, 30.2, 20>,最終結果為 <1, 30.2, 35>。支援的聚合函數及其對應的資料類型為:

  • sum:支援 DECIMAL、TINYINT、SMALLINT、INTEGER、BIGINT、FLOAT 和 DOUBLE。

  • min 和 max:支援 DECIMAL、TINYINT、SMALLINT、INTEGER、BIGINT、FLOAT、DOUBLE、DATE、TIME、TIMESTAMP 和 TIMESTAMP_LTZ。

  • last_value 和 last_non_null_value:支援所有資料類型。

  • listagg:支援 STRING。

  • bool_and 和 bool_or:支援 BOOLEAN。

說明
  • 只有 sum 函數支援回縮和刪除;其他聚合函數不支援。如果您需要某些欄位忽略回縮和刪除訊息,可以設定 'fields.${field_name}.ignore-retract'='true'.

  • 要以串流模式讀取聚合結果,您必須將 changelog-producer parameter to lookup or full-compaction.

Changelog 生產者

Set the changelog-producer 參數,將 Paimon 配置為為任何輸入串流生成完整的 changelog(其中每條 update_after 記錄都有對應的 update_before 記錄)。下表描述了可用的 changelog 生產者。更多詳情,請參見 Apache Paimon 官方文件.

生產者

描述

否ne

When you set changelog-producer to none (預設值)時,下游 Paimon 來源表只能看到給定主鍵的最新資料狀態。這種不完整的 changelog 使消費者難以進行正確的計算,因為他們無法確定資料的先前狀態,只能知道它是否被刪除或最新狀態是什麼。

例如,如果下游消費者需要計算一個欄位的總和,但只看到最新值 5,它無法確定如何更新總計。如果之前的值是 4,總和應該增加 1;如果之前的值是 6,總和應該減少 1。對 update_before 記錄敏感的消費者不應使用 none 生產者,但其他 changelog 生產者會帶來效能開銷。

說明

如果您的下游消費者(例如資料庫)對 update_before 資料不敏感,您可以使用 none 生產者。請根據您的具體需求配置 changelog 生產者。

Input

When you set changelog-producer to input時,sink 表會將輸入串流雙寫到 changelog 檔案中。

僅當輸入串流本身已經是完整的 changelog 時才使用此生產者,例如來自變更資料擷取(CDC)的資料。

Lookup

When you set changelog-producer to lookup時,sink 表使用點查詢機制(類似於維度表查找)在 snapshot 提交前為當前 snapshot 生成完整的 changelog。此生產者可以從任何輸入串流生成完整的 changelog。

full-compaction producer, the lookup 生產者提供更好的 changelog 時效性,但總體消耗更多資源。

此選項適用於需要高資料新鮮度(例如分鐘級)的使用場景。

Full Compaction

When you set changelog-producer to full-compaction時,sink 表在每次完整 compaction 時生成完整的 changelog。此生產者可以從任何輸入串流生成完整的 changelog。完整 compaction 的間隔由 full-compaction.delta-commits 參數指定。

lookup producer, the full-compaction 生產者的延遲較高,但利用了現有的完整 compaction 流程而無需添加額外計算。這導致整體資源消耗較低。

此選項適用於對資料新鮮度要求較低的使用場景(例如按小時)。

Write mode

Paimon 表支援以下寫入模式。

Mode

描述

Change-log

The change-log 寫入模式是 Paimon 表的預設模式。此模式支援基於主鍵的插入、刪除和更新操作。您也可以在此模式中使用合併引擎和 changelog 生產者。

Append-only

The append-only 寫入模式僅支援資料插入,不使用主鍵。此模式比 change-log 模式更高效,在資料新鮮度要求適中(例如分鐘級新鮮度)的場景中可用作訊息佇列的替代方案。

有關 append-only 寫入模式的詳細描述,請參見 Apache Paimon 官方文件。使用此模式時,請注意以下事項:

  • Set the bucket-key 參數。否則,Paimon 表將根據所有欄位的值分配 bucket,這在計算上效率低下。

  • The append-only 寫入模式在一定程度上可以保證資料的輸出順序。具體的輸出順序如下:

    1. 對於來自不同分區的記錄:如果設定了 scan.plan-sort-partition 參數,值較小的分區的記錄先輸出。否則,較早建立的分區的記錄先輸出。

    2. 對於來自同一分區和 bucket 的記錄,較早寫入的記錄先輸出。

    3. 對於來自同一分區但不同 bucket 的記錄,輸出順序不能保證,因為不同的 bucket 由不同的並行任務處理。

CTAS 與 CDAS 的目標

Paimon 表支援單個表或整個資料庫的即時資料同步。上游表的 schema 變更也會即時同步到 Paimon 表。詳情請參見 管理 Paimon 表 and 管理 Paimon Catalog.

Variant 讀取剪枝

僅讀取查詢引用的 Variant 欄位,減少 I/O 和記憶體開銷。

如何啟用

參數

預設值

描述

variant.read.pushdown.enabled

false

設定為 true 以啟用 Variant 讀取剪枝。

前提條件

  • Variant 資料只能由 VVR 11.6 或更高版本寫入。由早期版本或其他引擎寫入的 Variant 資料不支援讀取剪枝。

  • 要使上游資料生效,您必須啟用 variant.inferShreddingSchema 選項以啟用自動 Variant shredding 推斷,或明確配置 variant.shreddingSchema 及相關選項以指定哪些 Variant 成員適合 shredding。更多資訊請參見 Coreoptions.

支援的場景

當 SQL 通過字串鍵存取 Variant 欄位時,讀取剪枝生效,例如:

SELECT v['a']            FROM t;   -- Single level
SELECT v['a']['b']       FROM t;   -- Nested
SELECT v['a'], v['b']    FROM t;   -- Multiple fields
SELECT id, v['a']        FROM t;   -- Mixed with regular columns
SELECT v['a'] + 1        FROM t;   -- Referenced in expressions

不支援的場景

  • 直接引用整個 Variant,例如 SELECT v FROM t.

  • 將直接引用與欄位存取組合,例如 SELECT v, v['a'] FROM t.

  • 陣列索引存取,例如 v[0].

  • 表 schema 包含巢狀的 Row 欄位。在這種情況下,讀取剪枝不適用於表中的任何 Variant 欄位。

Blob 預簽名 URL

Paimon 可以為儲存在 OSS 中的 Blob 欄位資料生成短期公開 HTTPS GET 預簽名 URL。外部服務(例如多模態模型服務)可以使用這些 URL 直接擷取檔案內容,而無需在 Flink 作業內傳輸檔案位元組。

限制

  • 僅 Flink 計算引擎 VVR 11.9-preview1 及更高版本支援此功能。

  • 僅支援儲存在 OSS 中的 Paimon 表。元資料儲存類型不限,Filesystem、DLF 及其他元資料儲存類型均受支援。

  • 生成的 URL 是公開的 HTTPS URL,必須由具有公網存取權限的外部服務消費。如果 catalog 使用內部端點,建議您將其更改為公開的 HTTPS 端點,以便外部服務可以存取生成的 URL。

  • 要讀取常規 Blob 欄位,您必須先通過 blob-as-descriptor 選項獲取 BlobDescriptor 位元組,然後將位元組傳遞給函數。

語法

descriptor_to_presigned_url(source_table, descriptor, validity)
try_descriptor_to_presigned_url(source_table, descriptor, validity)

參數s

參數

Type

描述

source_table

STRING literal

包含 Blob 欄位的 Paimon 表,格式為 database.table。如果表的 catalog 與函數所屬的 catalog 相同,您也可以使用 catalog.database.table。值必須是非空字面量。不支援表欄位、動態欄位和 CASE 表達式。

descriptor

BYTES

BlobDescriptor 位元組,即啟用 blob-as-descriptor = 'true' 後讀取 Blob 欄位的結果。

validity

INTERVAL

URL 的有效期。值必須為正數秒,例如 INTERVAL '5' MINUTE.

回傳值

Type

描述

STRING

短期公開的 HTTPS GET 預簽名 URL。

範例

-- Use an SQL hint to read the Blob column as a descriptor and generate a presigned URL
SELECT sys.descriptor_to_presigned_url(
    'default.image_table',
    image,
    INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;

-- Fault-tolerant mode: returns NULL on row-level errors, with other behavior unchanged
SELECT sys.try_descriptor_to_presigned_url(
    'default.image_table',
    image,
    INTERVAL '5' MINUTE)
FROM image_table /*+ OPTIONS('blob-as-descriptor'='true') */;
說明
  • descriptor_to_presigned_url 在行級錯誤時拋出異常。 try_descriptor_to_presigned_url 在行級錯誤時返回 NULL。

  • 生成的 URL 不包含副檔名,URL 指向的位元組與 Blob 資料相同。下游服務必須根據返回的位元組而不是 URL 後綴來識別檔案格式。

  • 預簽名 URL 等同於臨時存取憑證。預簽名 URL 生成後,請立即提交給下游服務。不要將其寫入日誌或結果表,也不要持久化儲存。

將 Paimon 用作維度表

Paimon 表可用作維度表。有關 JOIN 語法,請參見 維度表 JOIN 語句.

預設情況下,lookup 會在每個並行實例中載入所有資料。此方法僅適用於小型維度表。對於大型維度表,請使用下面描述的 Shuffle Lookup 方案。

分區維度表

如果您的維度表是分區的,且您只需要最新一到兩個分區的資料,可以使用動態分區載入功能:

SELECT * FROM T
JOIN DIM /*+ OPTIONS('lookup.dynamic-partition'='max_pt()', 'lookup.dynamic-partition.refresh-interval'='1 h') */
FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col = D.col;

參數

資料類型

預設值

描述

lookup.dynamic-partition

String

N/A

max_pt():僅載入最新分區。 max_two_pt():僅載入最新的兩個分區。

lookup.dynamic-partition.refresh-interval

Duration

1 h

系統檢查維度表中分區更新的間隔。

大型維度表:固定 bucket 表

僅在 VVR 8.0.8 及更高版本中受支援。對於固定 bucket 表(bucket > 0),您可以使用 Shuffle Lookup 按 bucket key 將資料分配到各並行實例,每個實例只載入其分配的 bucket 中的資料:

SELECT /*+ LOOKUP('table'='D', 'shuffle'='true') */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;
說明
  • join key 必須是 bucket key。bucket key 預設為主鍵。

  • 僅固定 bucket 表(bucket > 0)支援此功能。

大型維度表:非固定 bucket 表

僅在 VVR 8.0.10 及更高版本中受支援。對於動態 bucket 表或 append 表,您可以使用 SHUFFLE_HASH 或 REPLICATED_SHUFFLE_HASH,讓每個並行實例讀取所有資料,但只保留其所需的部分:

-- Shuffle Hash
SELECT /*+ SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;

-- Replicated Shuffle Hash
SELECT /*+ REPLICATED_SHUFFLE_HASH(D) */ T.col1, D.col2
FROM T
JOIN DIM FOR SYSTEM_TIME AS OF T.proc_time AS D
ON T.col1 = D.col1;

有關 SHUFFLE_HASH 和 REPLICATED_SHUFFLE_HASH 的更多資訊,請參見 維度表 JOIN 語句.

資料接入

您可以在資料接入 YAML 作業中,將 Paimon connector 用作 sink。

語法

sink:
  type: paimon
  name: Paimon Sink
  catalog.properties.metastore: filesystem
  catalog.properties.warehouse: /path/warehouse

參數s

參數

描述

是否必填

Type

預設值

否tes

type

connector 的類型。

STRING

否ne

值必須為 paimon.

name

sink 的名稱。

STRING

否ne

catalog.properties.metastore

Paimon catalog 的類型。

STRING

filesystem

有效值:

  • filesystem(預設值)

  • rest(僅支援 Data Lake Formation(DLF),不支援 DLF-Legacy)

catalog.properties.*

建立 Paimon catalog 的參數。

STRING

否ne

For more information為前綴的參數的資訊,請參見 管理 Paimon Catalog.

table.properties.*

建立 Paimon 表的參數。

STRING

否ne

For more information為前綴的參數的資訊,請參見 Paimon 表選項.

catalog.properties.warehouse

檔案儲存的根目錄。

STRING

否ne

This parameter applies only when catalog.properties.metastore is set to filesystem.

commit.user-prefix

提交資料檔案的使用者名稱前綴。

STRING

否ne

說明

建議為不同作業設定不同的使用者名稱。這樣更容易識別導致提交衝突的作業。

partition.key

分區表的分區鍵。

STRING

否ne

不同表用 ;分隔,不同欄位用 ,分隔,表和欄位用 :分隔。例如: testdb.table1:id1,id2;testdb.table2:name.

sink.cross-partition-upsert.tables

列出需要跨分區 upsert 的表,其中主鍵不包含所有分區鍵。

STRING

否ne

適用於需要跨分區更新的表。

  • 格式:使用分號 ; 分隔表名。

  • 效能建議:此操作資源消耗大。請為這些表建立單獨的作業。

重要
  • 您必須列出所有符合條件的表。省略表名將導致資料重複。

sink.commit.parallelism

指定 Commit 算子的並行度。

INTEGER

否ne

如果 Commit 算子是瓶頸,請使用此參數增加其並行度以提升效能。

此參數僅在 Realtime Compute for Apache Flink 11.6 及更高版本中受支援。

說明

設定此參數會更改算子並行度。重啟有狀態作業時,您必須指定 AllowNonRestoredState ,使作業可以忽略部分算子狀態。

重複使用現有 catalog

從 Realtime Compute for Apache Flink 11.5 開始,您可以在 Flink CDC 資料接入作業中,直接從「資料管理」頁面引用內建的 Paimon catalog。這減少了手動配置。

sink:
  type: paimon
  using.built-in-catalog: paimon_dlf_catalog
  catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com

資料接入作業可以自動重複使用 Paimon catalog 中的所有參數。這等同於在 YAML 作業中手動配置以 catalog.properties. 為前綴的參數。

要覆蓋自動重複使用的參數,請在 YAML 作業中明確設定。明確的 YAML 配置具有更高優先順序。例如,在上述範例中, fs.oss.endpoint 參數使用 YAML 作業中的值,覆蓋 paimon_dlf_catalog.

範例

將 Paimon 用作資料接入 sink 時,請參考以下範例根據您的 Paimon Catalog 類型配置作業。

  • 將資料寫入 Object Storage Service(OSS)並使用 filesystem Paimon catalog 的配置範例:

    source:
      type: mysql
      name: MySQL Source
      hostname: ${secret_values.mysql.hostname}
      port: ${mysql.port}
      username: ${secret_values.mysql.username}
      password: ${secret_values.mysql.password}
      tables: ${mysql.source.table}
      server-id: 8601-8604
    
    sink:
      type: paimon
      name: Paimon Sink
      catalog.properties.metastore: filesystem
      catalog.properties.warehouse: oss://default/test
      catalog.properties.fs.oss.endpoint: oss-cn-beijing-internal.aliyuncs.com
      catalog.properties.fs.oss.accessKeyId: xxxxxxxx
      catalog.properties.fs.oss.accessKeySecret: xxxxxxxx

    有關以 catalog.properties為前綴的參數的資訊,請參見 建立 Paimon Filesystem Catalog.

  • 將資料寫入 Data Lake Formation(DLF) 並使用 rest Paimon catalog 的配置範例:

    source:
      type: mysql
      name: MySQL Source
      hostname: ${secret_values.mysql.hostname}
      port: ${mysql.port}
      username: ${secret_values.mysql.username}
      password: ${secret_values.mysql.password}
      tables: ${mysql.source.table}
      server-id: 8601-8604
    
    sink:
      type: paimon
      name: Paimon Sink
      catalog.properties.metastore: rest
      catalog.properties.uri: dlf_uri
      catalog.properties.warehouse: your_warehouse
      catalog.properties.token.provider: dlf
      # (Optional) Enable deletion vectors to improve read performance.
      table.properties.deletion-vectors.enabled: true

    有關以 catalog.properties為前綴的參數的資訊,請參見 Flink CDC Catalog 配置參數.

Schema 變更

當用作資料接入 sink 時,Paimon 支援以下 schema 變更事件:

  • CREATE TABLE 事件

  • ADD COLUMN 事件

  • ALTER COLUMN TYPE 事件(不支援更改主鍵欄位的資料類型)

  • RENAME COLUMN 事件

  • DROP COLUMN 事件

  • TRUNCATE TABLE 事件

  • DROP TABLE 事件

說明

如果下游 Paimon 表已存在,作業將寫入現有 schema,不會嘗試再次建立表。

常見問題