本主題說明如何將 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/testandoss://b/test/是不同的表(因為尾部斜線的差異),儘管它們可能指向相同的物理位置。
-
WITH 參數
|
參數 |
描述 |
Type |
是否必填 |
預設值 |
備註 |
|
connector |
指定表的 connector。 |
String |
否 |
否ne |
|
|
path |
表的儲存路徑。 |
String |
否 |
否ne |
|
|
auto-create |
指定當指定路徑下不存在表檔案時,是否自動建立。 |
Boolean |
否 |
false |
有效值:
|
|
file.format |
資料檔案的格式。 |
String |
否 |
parquet |
有效值:
|
|
bucket |
每個分區的 bucket 數量。 |
Integer |
否 |
1 |
Paimon 根據 說明
建議每個 bucket 包含的資料少於 5 GB。 |
|
bucket-key |
用作 bucket key 的欄位。 |
String |
否 |
否ne |
指定用於將資料分配到 bucket 的欄位。 多個欄位名稱用逗號(,)分隔。例如, 說明
|
|
changelog-producer |
changelog 的生產機制。 |
String |
否 |
none |
Paimon 可以為任何輸入串流生成完整的 changelog,這意味著每條
有關如何選擇 changelog 生產者的更多資訊,請參見 Changelog 生產. |
|
full-compaction.delta-commits |
兩次連續完整 compaction 之間的最大提交次數。 |
Integer |
否 |
否ne |
指定在觸發完整 compaction 之前允許的最大 snapshot 提交次數。 |
|
lookup.cache-max-memory-size |
Paimon 維度表的記憶體快取大小。 |
String |
否 |
256 MB |
此參數控制維度表查找和 |
|
merge-engine |
合併具有相同主鍵的記錄的機制。 |
String |
否 |
deduplicate |
有效值:
有關合併引擎的詳細分析,請參見 合併引擎. |
|
partial-update.ignore-delete |
指定是否忽略刪除(-D)訊息。 |
Boolean |
否 |
false |
有效值:
說明
|
|
ignore-delete |
指定是否忽略刪除(-D)訊息。 |
Boolean |
否 |
false |
有效值與 partial-update.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.num-retained.min |
保留最新 snapshot 的最小數量。 |
Integer |
否 |
10 |
N/A |
|
snapshot.time-retained |
snapshot 的保留期。 |
String |
否 |
1h |
如果此條件或 |
|
write-mode |
Paimon 表的寫入模式。 |
String |
否 |
change-log |
有效值:
有關寫入模式的更多資訊,請參見 寫入模式. |
|
scan.infer-parallelism |
指定是否自動推斷 Paimon 來源表的並行度。 |
Boolean |
否 |
true |
有效值:
|
|
scan.parallelism |
Paimon 來源表的並行度。 |
Integer |
否 |
否ne |
說明
如果在作業的 Resource Mode is set to Mode on the job's 標籤頁,則忽略此參數。 |
|
sink.parallelism |
Paimon sink 表的並行度。 |
Integer |
否 |
否ne |
說明
如果在作業的 Resource Mode is set to Mode on the job's 標籤頁,則忽略此參數。 |
|
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 欄位分配資料。 多個欄位名稱用逗號(,)分隔,例如 有關聚簇的更多資訊,請參見 Apache Paimon 官方文件. |
|
sink.delete-strategy |
指定驗證策略,以確保系統正確處理回縮訊息(-D/-U)。 |
Enum |
否 |
NONE |
有效值及 sink 算子處理回縮訊息時的預期行為:
說明
|
|
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 叢集數量。 |
Integer |
否 |
16 |
較高的值通常會提升召回率,但也會增加延遲。 |
|
|
檢索 top_k × refine_factor 個 IVF 候選結果,並使用 Paimon 表中儲存的原始向量重新排序。 |
Float |
否 |
停用 |
所有 IVF 變體預設禁用。最適合在召回率優先於延遲的場景中使用壓縮索引(如 ivf-pq 和 ivf-hnsw-sq)。 |
|
|
搜尋時的 HNSW 搜尋寬度。 |
Integer |
否 |
0 |
較高的值通常會提升召回率,但也會增加延遲。 0 indicates that the native library default is used. |
|
|
Lumina DiskANN 搜尋清單大小。 |
Integer |
否 |
max(1.5 × top_k, 16) |
較高的值通常會提升召回率,但也會增加延遲。 |
|
|
Lumina DiskANN 搜尋束寬度。 |
Integer |
否 |
4 |
— |
|
|
並行 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, 25.2, 10, 'This is a book'>。 說明
|
|
Aggregation |
在某些使用場景中,您可能只需要記錄的聚合值。聚合引擎使用您指定的聚合函數組合具有相同主鍵的記錄。對於每個非主鍵欄位,您必須使用
The
說明
|
Changelog 生產者
Set the changelog-producer 參數,將 Paimon 配置為為任何輸入串流生成完整的 changelog(其中每條 update_after 記錄都有對應的 update_before 記錄)。下表描述了可用的 changelog 生產者。更多詳情,請參見 Apache Paimon 官方文件.
|
生產者 |
描述 |
|
否ne |
When you set 例如,如果下游消費者需要計算一個欄位的總和,但只看到最新值 5,它無法確定如何更新總計。如果之前的值是 4,總和應該增加 1;如果之前的值是 6,總和應該減少 1。對 說明
如果您的下游消費者(例如資料庫)對 |
|
Input |
When you set 僅當輸入串流本身已經是完整的 changelog 時才使用此生產者,例如來自變更資料擷取(CDC)的資料。 |
|
Lookup |
When you set 與 此選項適用於需要高資料新鮮度(例如分鐘級)的使用場景。 |
|
Full Compaction |
When you set 與 此選項適用於對資料新鮮度要求較低的使用場景(例如按小時)。 |
Write mode
Paimon 表支援以下寫入模式。
|
Mode |
描述 |
|
Change-log |
The |
|
Append-only |
The 有關
|
CTAS 與 CDAS 的目標
Paimon 表支援單個表或整個資料庫的即時資料同步。上游表的 schema 變更也會即時同步到 Paimon 表。詳情請參見 管理 Paimon 表 and 管理 Paimon Catalog.
Variant 讀取剪枝
僅讀取查詢引用的 Variant 欄位,減少 I/O 和記憶體開銷。
如何啟用
|
參數 |
預設值 |
描述 |
|
|
false |
設定為 |
前提條件
-
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 |
描述 |
|
|
STRING literal |
包含 Blob 欄位的 Paimon 表,格式為 |
|
|
BYTES |
BlobDescriptor 位元組,即啟用 |
|
|
INTERVAL |
URL 的有效期。值必須為正數秒,例如 |
回傳值
|
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 |
|
|
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 |
值必須為 |
|
name |
sink 的名稱。 |
否 |
STRING |
否ne |
|
|
catalog.properties.metastore |
Paimon catalog 的類型。 |
否 |
STRING |
filesystem |
有效值:
|
|
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 |
|
commit.user-prefix |
提交資料檔案的使用者名稱前綴。 |
否 |
STRING |
否ne |
說明
建議為不同作業設定不同的使用者名稱。這樣更容易識別導致提交衝突的作業。 |
|
partition.key |
分區表的分區鍵。 |
否 |
STRING |
否ne |
不同表用 |
|
sink.cross-partition-upsert.tables |
列出需要跨分區 upsert 的表,其中主鍵不包含所有分區鍵。 |
否 |
STRING |
否ne |
適用於需要跨分區更新的表。
重要
|
|
sink.commit.parallelism |
指定 Commit 算子的並行度。 |
否 |
INTEGER |
否ne |
如果 Commit 算子是瓶頸,請使用此參數增加其並行度以提升效能。 此參數僅在 Realtime Compute for Apache Flink 11.6 及更高版本中受支援。 說明
設定此參數會更改算子並行度。重啟有狀態作業時,您必須指定 |
重複使用現有 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)並使用
filesystemPaimon 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) 並使用
restPaimon 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,不會嘗試再次建立表。