適用於 Apache Iceberg 的 Dataflow 受管理 I/O

代管 I/O 支援 Apache Iceberg 的下列功能:

目錄
讀取功能 批次讀取
寫入功能

如果是 Apache Iceberg 專用 BigQuery 資料表,請搭配使用 BigQueryIO 連接器和 BigQuery Storage API。資料表必須已存在,不支援動態建立資料表。

需求條件

下列 SDK 支援 Apache Iceberg 的代管 I/O:

  • Apache Beam SDK for Java 2.58.0 以上版本
  • Apache Beam SDK for Python 2.61.0 以上版本

設定

Apache Iceberg 的代管 I/O 支援下列設定參數:

ICEBERG 閱讀

設定 類型 說明
資料表 str Iceberg 資料表的 ID。
catalog_name str 資料表所在目錄的名稱。
catalog_properties map[str, str] 用於設定 Iceberg 目錄的屬性。
config_properties map[str, str] 傳遞至 Hadoop 設定的屬性。
drop list[str] 要排除讀取的資料欄名稱子集。如果為空值或空白,系統會讀取所有資料欄。
篩選器 str 類似 SQL 的述詞,可在掃描時篩選資料。例如:「id > 5 AND status = 'ACTIVE'」。使用 Apache Calcite 語法:https://calcite.apache.org/docs/reference.html
保留 list[str] 要讀取的資料欄名稱子集。如果為空值或空白,系統會讀取所有資料欄。

ICEBERG 寫入

')。

設定 類型 說明
資料表 str 完整資料表 ID。您也可以提供範本,將資料寫入多個動態目的地,例如:`dataset.my_{col1}_{col2.nested}_table`。
allowed_lateness_seconds int32 延遲記錄在完全捨棄前,可能落後浮水印的時間長度,而不是傳送至 dead_letter 輸出內容。預設值為 21600 (6 小時)。目前僅支援「讀取時合併」模式。
自動分片 boolean 啟用動態分片功能,根據資料量自動調整平行寫入器數量。這項功能會將分區進一步細分為多個分片,藉此處理資料偏斜問題,避免高處理量寫入期間發生瓶頸。僅適用於「雜湊」發布模式。
catalog_name str 資料表所在目錄的名稱。
catalog_properties map[str, str] 用於設定 Iceberg 目錄的屬性。
change_type_column str 僅限讀取時合併。代表資料列變更類型的選用資料欄名稱 (INSERT、UPDATE_BEFORE、UPDATE_AFTER 或 DELETE)。在寫入 Iceberg 前,系統會從資料列中移除這個資料欄。如未設定,接收器會使用元素的原生 ValueKind
change_type_map map[str, str] 僅限讀取時合併。選用對應,可將 change_type_column 值對應至標準變更類型名稱 (請參閱上文)。
config_properties map[str, str] 傳遞至 Hadoop 設定的屬性。
direct_write_byte_limit int32 針對串流管道,設定資料組合改用直接寫入路徑時的資料大小上限。
distribution_mode str 定義寫入資料的分配情形。支援的分配方式: - 無:不隨機排序資料列 (預設) - 雜湊:在寫入資料前,依據分割區鍵隨機排序資料列
drop list[str] 寫入前要從輸入記錄捨棄的欄位名稱清單,與「keep」和「only」互斥。在讀取時合併模式中,控制項資料欄一律會捨棄。
equality_columns list[str] 定義列 ID 的欄 (等號刪除欄位)。預設為目標資料表的 ID (主鍵) 欄位。如果資料表尚未建立,則為必要欄位。目前僅支援「讀取時合併」模式。
保留 list[str] 要保留於輸入記錄的欄位名稱清單。寫入之前,其他欄位都會捨棄。與「drop」和「only」互斥。在讀取時合併模式中,除非列於此處,否則控制項資料欄會遭到捨棄。
maximum_table_cache_size int32 如果是批次管道,則設定要在記憶體中快取的資料表 Metadata 規格數量上限。如果資料表超過這個限制,系統會改為載入工作人員本機目錄。
模式 str 控制資料列的寫入方式。「append」(預設) 會將每個資料列附加為新資料。「merge-on-read」會將每個資料列視為套用至資料表 (透過主鍵) 的變更 (INSERT、UPDATE_BEFORE、UPDATE_AFTER 或 DELETE)。
num_shards int32 每個目的地確定性主鍵雜湊分片的數量,也就是每個目的地的寫入平行處理量上限。如果值太低,可能會造成寫入作業瓶頸;如果值太高,可能會產生更多檔案。預設值為 16。目前僅支援「讀取時合併」模式。
僅限 str 要寫入的單一記錄欄位名稱,與「keep」和「drop」互斥。
partition_fields list[str] 用於建立分區規格的欄位,該規格會在建立資料表時套用。以欄位「foo」來說,可用的分區轉換作業包括:
  • foo
  • truncate(foo, N)
  • bucket(foo, N)
  • hour(foo)
  • day(foo)
  • month(foo)
  • year(foo)
  • void(foo)

如要進一步瞭解分割區轉換,請前往 https://iceberg.apache.org/spec/#partition-transforms。

sequence_number_column str 僅限讀取時合併。必要資料欄名稱,代表用於排序單一金鑰變更的單調遞增序號。預設值為「_commit_snapshot_sequence_number」。在寫入 Iceberg 之前,系統會從資料列中移除這個資料欄。
shards_per_partition int32 單一資料分割的資料列可占用的分片數量上限。值越小,每次提交寫入的檔案就越少,但也會降低每個分割區的寫入平行處理能力。值為 1 時,每個分割區會固定指派給一個寫入器。未分區資料表會忽略這項設定。必須介於 1 至 num_shards 之間;預設為 num_shards。目前僅支援「讀取時合併」模式。
sink_id str 這個接收器的穩定 ID,用於為寫入每個提交內容 Iceberg 快照摘要的等冪權杖設定命名空間。預設值為每個寫入作業的專屬 UUID。針對特定串流寫入作業的重新啟動,明確設定 (並在重新啟動時保持穩定),確保只會提交一次。使用穩定 sink_id 的批次載入只會提交一次 (之後使用相同 sink_id 的批次載入會略過)。目前僅支援「讀取時合併」模式。
snapshot_properties map[str, str] 要新增至每個提交的 Iceberg 快照摘要的額外鍵/值屬性。前置字串為「beam.cdc.」的鍵為保留項目,系統會拒絕這類鍵。目前僅支援「讀取時合併」模式。
sort_fields list[str] 用於設定資料表排序順序的欄位,會在建立資料表時套用。每個項目的格式為 <term> [asc|desc] [nulls first|nulls last],其中 <term> 是欄位名稱或其中一個分割區轉換 (例如 bucket(col, 4)、day(ts))。方向預設為遞增;空值順序預設為遞增時空值優先,遞減時空值在後。注意:這會將資料表的宣告排序順序設為中繼資料,不會導致 Beam 在寫入前實際排序記錄。 如要進一步瞭解排序順序,請前往 https://iceberg.apache.org/spec/#sort-orders。
sorter_memory_mb int32 預先寫入排序的記憶體內緩衝區空間 (MB);大於此大小的群組會溢出至磁碟。必須大於或等於 1。預設值為 100。目前僅支援「讀取時合併」模式。
table_cache_polling_buckets int32 設定用於在重新整理期間查詢 Iceberg 目錄的平行 bucket/worker 數量。預設值為 1。
table_cache_refresh_interval_seconds int32 設定串流管道從目錄重新整理資料表中繼資料的時間間隔 (以秒為單位)。
table_properties map[str, str] 建立資料表時要設定的 Iceberg 資料表屬性。 如要進一步瞭解資料表屬性,請前往 https://iceberg.apache.org/docs/latest/configuration/#table-properties。
token_heartbeat_seconds int32 僅限串流播放。如果已設定,接收器會在閒置時定期發出空白權杖重新整理提交內容,因此其 sink_id 蓋章快照的執行緒會保持最新狀態,較不容易因 expire_snapshots 而遺失。預設為停用。目前僅支援「讀取時合併」模式。
triggering_frequency_seconds int32 設定串流管道的快照產生頻率。
upsert boolean 僅限讀取時合併。如為 true,系統只會套用每個變更的後續圖片 (INSERT/UPDATE_AFTER),做為 upsert;系統會捨棄 UPDATE_BEFORE 記錄。預設值:false。
use_side_input_table_cache boolean 在各個工作人員之間啟用 Iceberg 資料表的中繼資料快取功能,可減少目錄負載。
write_properties map[str, str] 套用至基礎檔案寫入器的屬性 (例如 Parquet 寫入屬性,如「write.parquet.bloom-filter-enabled.column.

後續步驟

如需更多資訊和程式碼範例,請參閱下列主題: