Important
這項功能目前處於 公開預覽版。
使用 Structured Streaming 將資料寫入 Lakebase,具備內建批次處理、自動重試,以及由工作區管理的驗證。
何時使用 Lakebase 水槽
使用 Lakebase 接收器將資料以低延遲串流寫入 Lakebase。 這個接收端不需要你實作自訂的 foreachBatch 函式來處理批次處理、連線管理和錯誤處理。
常見的使用案例包括:
- 即時更新應用程式資料庫,以支援營運儀表板或面向客戶的功能。
- 將持續變動的資料,例如彙總或過濾的串流結果,同步到交易式資料庫。
- 使用 即時模式,以亞秒級延遲將結構化串流查詢的輸出寫入 Lakebase 資料表。
若要以相反方向將資料從 Lakebase 同步至 Lakehouse 中的 Delta Lake 資料表,請參閱 Lakebase 變更資料摘要。
需求規格
- Databricks Runtime 18 及更新版本
- 經典運算模式,具備專用或標準存取模式。
- Lakebase 的資料庫
連線至資料庫
湖底排水槽支援以下連接方式:
已在 Unity Catalog 註冊的 Lakebase 資料表
對於註冊於 Unity Catalog 的 Lakebase 資料表,連接器會自動管理憑證,並使用執行查詢的使用者或服務主體的身份。 如果表格不存在,連接器就會建立表格。
若要在 Unity 目錄中註冊 Lakebase 資料庫,請參閱 在 Unity 目錄中註冊 Lakebase 資料庫。
若要將資料寫入 Lakebase 資料表,請使用 .toTable() 方法,並搭配完整限定資料表名稱 catalog.schema.table。 以下範例展示了所需的選項,以及可選 upsertkey 選項:
Python
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
Scala
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
替換下列佔位符:
-
<catalog>.<schema>.<table>: 目標資料表的完全限定名稱。 這是catalog你註冊 Lakebase 資料庫時建立的 Unity 目錄,請參見 「在 Unity 目錄中註冊 Lakebase 資料庫」。 如果表格不存在,連接器會產生它。 -
<primary-key-column>:選擇性。 以逗號分隔、構成 upsert 索引鍵的欄位清單,例如id或user_id,event_type。 若省略upsertkey,接收端會根據目標資料表的主鍵推斷該鍵值。 請參考 Upsert的行為。 -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>:查詢將其檢查點儲存在此處的 Unity Catalog 磁碟區路徑。 你也可以使用雲端物件儲存 URI。 位置必須是你可以寫入的儲存空間,而非本地磁碟,且每個串流查詢必須是獨一無二的。 這與目標表無關。 請參閱結構化串流檢查點。
關於可選配置,如 batchsize 和 batchinterval,請參見 配置選項。
未在 Unity Catalog 中註冊的 Lakebase 資料表
對於未在 Unity Catalog 註冊的 Lakebase 資料表,連接器會自動管理憑證,並使用執行查詢的使用者或服務主體的身份。 如果表格不存在,連接器就會建立表格。
要寫入 Lakebase 資料表,請使用 endpoint and dbtable 選項。 以下範例亦包含可選 database 與 upsertkey 選項:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") # Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") // Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
替換下列佔位符:
-
<project-id>.<branch-id>.<endpoint-id>:您的 Lakebase 端點。 在 Computes 標籤的 Get ID 選單中,找到資源名稱中的三個值,格式為projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>。 請參見 計算識別碼。 -
<database>:選擇性。 目標 Postgres 資料庫的名稱。 預設為databricks_postgres。 請參閱 管理資料庫。 -
<schema>.<table>:採用schema.table格式的目標資料表。 如果你省略了 schema,sink 就會使用該publicschema。 使用以字母或底線開頭的簡單識別碼,且僅包含字母、數字和底線;不支援引用的識別碼和特殊字元,例如連字號。 -
<primary-key-column>:選擇性。 以逗號分隔、構成 upsert 索引鍵的欄位清單,例如id或user_id,event_type。 若省略upsertkey,接收端會根據目標資料表的主鍵推斷該鍵值。 請參考 Upsert的行為。 -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>:查詢將其檢查點儲存在此處的 Unity Catalog 磁碟區路徑。 你也可以使用雲端物件儲存 URI。 位置必須是你可以寫入的儲存空間,而非本地磁碟,且每個串流查詢必須是獨一無二的。 這與目標表無關。 請參閱結構化串流檢查點。
關於可選配置,如 batchsize 和 batchinterval,請參見 配置選項。
設定選項
水槽因未識別選項而產生錯誤訊息。 JDBC_STREAMING_SINK_INVALID_OPTIONS
以下選項適用於所有連接方式:
| Key | 預設 | Description |
|---|---|---|
batchinterval |
100 milliseconds |
Optional. 沖洗前保持緩衝區行數的最長時間。 例如: "50 milliseconds" 。 |
batchsize |
1000 |
Optional. 每個資料庫交易的最大列數。 |
checkpointLocation |
沒有 | Required. 檢查點目錄的路徑,例如 Unity Catalog 磁碟區(/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>)。 每個查詢必須是獨一無二的。 請參閱結構化串流檢查點。 |
upsertkey |
沒有 | Optional. 一個逗號分隔的欄位名稱清單,構成 upsert 鍵。 例如,"id" 或 "user_id,event_type"。 如果你指定 upsertkey,欄位必須與資料表的主鍵相符,否則查詢會失敗。 如果你省略了,水槽會自動使用主鍵。 欲了解更多資訊,請參閱 Upsert 行為。 |
未在 Unity Catalog 中註冊的 Lakebase 資料表
當你 連接到未註冊於 Unity 目錄的 Lakebase 表格時,以下選項適用:
| Key | 預設 | Description |
|---|---|---|
database |
databricks_postgres |
Optional. 目標 PostgreSQL 資料庫名稱。 |
dbtable |
沒有 | Required.
schema.table 格式的目標資料表名稱。 如果你沒有指定結構,預設結構值為 public。 使用以字母或底線開頭的簡單識別碼,且僅包含字母、數字和底線。 請勿引用表格或結構名稱;不支援帶有特殊字元的引號識別碼及名稱,如連字號。 |
endpoint |
沒有 | Required. Lakebase 端點,格式為 project_id.branch_id 或 project_id.branch_id.endpoint_id。
endpoint_id 是選用的;如果省略它,且該分支只有一個可讀寫端點,接收端會預設選取該端點。 |
Upsert 行為
當 upsert 鍵存在時,無論是以 upsertkey 指定,或由 sink 根據資料表的主鍵推斷,sink 都會使用 PostgreSQL 的 INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ... 語法對資料表執行 upsert。
當沒有 upsert 鍵時,sink 會進行插入。 查詢的輸出模式不會影響 upsert 或 insert 行為。
upsertkey 欄必須:
- 必須為 DataFrame 欄位的非空子集。
- 完全符合目標表格的
PRIMARY KEY。 如果你指定的欄位與主鍵不符,查詢就會失敗。 - 可以是可比較的類型,例如數值型或字串型。 為了避免在同時寫入時發生資料庫死結,sink 會在每個批次中依 upsert 鍵排序資料列。 Upsert 鍵不支援複雜型或結構型。
欄位名稱會自動以 PostgreSQL 預設的雙引號 "引號引號,該格式處理保留關鍵字及混合大小寫名稱。
表格與結構名稱必須使用以字母或底線開頭的簡單識別碼,且僅包含字母、數字和底線。 匯入器不支援引號識別碼或表格或結構名稱中的特殊字元,如連字號。
性能調校
批次處理與背壓
當任一條件成立時,將觸發刷新:
- 緩衝區可達到
batchsize列,預設值為1000。 - 緩衝年齡超過
batchinterval,預設為100 milliseconds。
當資料庫無法跟上資料傳入速率時,sink 會將背壓往上游傳回來源端。
延遲與吞吐量指引:
- 對於使用即時模式的低延遲工作負載,請減少
batchinterval以確保刷新前的最長等待時間更短。 關於概念,請參見 即時模式概念 ,並以程式碼範例參考即時 模式範例 。 - 對於高吞吐量工作負載,請增加
batchsize以減少每筆交易的開銷。
連線行為
匯流器在執行器上使用連線池。 預設情況下,每個任務使用一個資料庫連線。
Databricks 建議您對每個連線使用 1 工作的預設值。 如果你增加每個連線的任務數量,可能會造成連線爭用,並增加高吞吐量連線的延遲。
要設定任務與連線的比例,請設定 spark.databricks.sql.streaming.jdbc.tasksPerConnection Spark 設定。 如果目標資料庫的連線限制很低,請減少洗牌分割區的數量或增加 spark.databricks.sql.streaming.jdbc.tasksPerConnection。
匯入器會自動重試暫時性的 JDBC 錯誤,包括連線失敗、死結及速率限制。 如果接收端耗盡所有重試次數,查詢就會失敗。
支援的觸發器與輸出模式
Triggers
下表顯示對結構化串流觸發器類型的支援:
| 觸發程序 | 支援 |
|---|---|
realTime |
是的 |
ProcessingTime |
是的 |
AvailableNow |
是的 |
Once |
是的 |
輸出模式
此表顯示對結構化串流輸出模式的支援:
| 輸出模式 | 支援 |
|---|---|
update |
是的 |
append |
Yes. 行為與 update相同。 當目標資料表有主鍵時,查詢會上溢出,否則查詢會插入。 請參考 Upsert的行為。 |
complete |
No |
限制
- 不支援無伺服器運算與 Lakeflow 管線。
- 僅支援將 Lakebase 作為寫入目標。 不支援外部 PostgreSQL 相容的資料庫。