使用 Structured Streaming 寫入 Lakebase 或具備內建批次處理、自動重試及工作區管理認證的外部 PostgreSQL 資料庫。
何時使用 Lakebase 水槽
使用 Lakebase sink 進行低延遲串流寫入 Lakebase 或外部 PostgreSQL 資料庫。 這個接收端不需要你實作自訂的 foreach 函式來處理批次處理、連線管理和錯誤處理。
常見的使用案例包括:
- 即時更新應用程式資料庫,以支援營運儀表板或面向客戶的功能。
- 將持續變動的資料,例如彙總或過濾的串流結果,同步到交易式資料庫。
- 使用 即時模式,以亞秒級延遲將結構化串流查詢的輸出寫入 Lakebase 資料表。
若要以相反方向將資料從 Lakebase 同步至 Lakehouse 中的 Delta Lake 資料表,請參閱 Lakebase 變更資料摘要。
需求規格
-
Databricks Runtime 18 LTS 及更新版本。
- 外部 PostgreSQL 連線需要你使用 Databricks Runtime 19 及以上版本,並且選擇加入 UC Compute 上的自訂 JDBC 預覽。
- 區間資料型態需要使用 Databricks Runtime 19 及以上版本。
- 經典運算模式,配備專用或標準存取模式,或是用於筆記本或作業的無伺服器運算。 在無伺服器運算中,請使用
Trigger.AvailableNow()。 請參考 無伺服器運算上的串流。 - 一個 Lakebase 資料庫,或 Unity Catalog 與外部 PostgreSQL 資料庫的連線。
識別碼要求
對於所有目標,Databricks 建議使用以字母或底線開頭且僅包含字母、數字與底線的結構、表格、欄位及主鍵欄位名稱。 當匯流器自動建立 Lakebase 表格時,會強制執行這些要求。 若要使用不符合這些需求的識別碼,請在開始查詢前建立目標資料表。
連線至資料庫
湖底排水槽支援以下連接方式:
已在 Unity Catalog 註冊的 Lakebase 資料表
對於註冊於 Unity Catalog 的 Lakebase 資料表,連接器會自動管理憑證,並使用執行查詢的使用者或服務主體的身份。 如果表格不存在,連接器就會建立表格。
若要在 Unity 目錄中註冊 Lakebase 資料庫,請參閱 在 Unity 目錄中註冊 Lakebase 資料庫。
若要寫入 Lakebase 資料表,請使用 .toTable() 方法,並搭配完全限定的資料表名稱 catalog.schema.table:
Python
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
Scala
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
替換下列佔位符:
-
<catalog>.<schema>.<table>: 目標資料表的完全限定名稱。 這是catalog你註冊 Lakebase 資料庫時建立的 Unity 目錄,請參見 「在 Unity 目錄中註冊 Lakebase 資料庫」。 如果表格不存在,連接器會產生它。 -
<primary-key-columns>:選擇性。 目標資料表主索引鍵中所有欄位的逗號分隔清單,例如id或user_id,event_type。 請參考 Upsert的行為。 -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>:查詢將其檢查點儲存在此處的 Unity Catalog 磁碟區路徑。 你也可以使用雲端物件儲存 URI。 位置必須是你可以寫入的儲存空間,而非本地磁碟,且每個串流查詢必須是獨一無二的。 這與目標表無關。 請參閱結構化串流檢查點。
關於可選配置,如 batchsize 和 batchinterval,請參見 PostgreSQL 的匯入選項。
未在 Unity Catalog 中註冊的 Lakebase 資料表
對於未在 Unity Catalog 註冊的 Lakebase 資料表,連接器會自動管理憑證,並使用執行查詢的使用者或服務主體的身份。 如果表格不存在,連接器就會建立表格。
若要寫入 Lakebase 資料表,請使用 dbtable 和 endpoint 選項:
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-columns>") # Optional
.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-columns>") // Optional
.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>:選擇性。 目標 PostgreSQL 資料庫的名稱。 預設為databricks_postgres。 請參閱 管理資料庫。 -
<schema>.<table>:採用schema.table格式的目標資料表。 如果你省略了 schema,sink 就會使用該publicschema。 自動建立表格時,請使用以字母或底線開頭,且僅包含字母、數字和底線的識別碼。 -
<primary-key-columns>:選擇性。 目標資料表主索引鍵中所有欄位的逗號分隔清單,例如id或user_id,event_type。 請參考 Upsert的行為。 -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>:查詢將其檢查點儲存在此處的 Unity Catalog 磁碟區路徑。 你也可以使用雲端物件儲存 URI。 位置必須是你可以寫入的儲存空間,而非本地磁碟,且每個串流查詢必須是獨一無二的。 這與目標表無關。 請參閱結構化串流檢查點。
關於可選配置,如 batchsize 和 batchinterval,請參見 PostgreSQL 的匯入選項。
外部 PostgreSQL 與 Unity 目錄憑證
Important
這項功能目前處於 公開預覽版。 Workspace 管理員可以從預覽頁面控制 UC Compute 上自訂 JDBC 的存取權限。 請參閱 管理 Azure Databricks 預覽。
使用 Unity Catalog 連線來驗證外部 PostgreSQL 資料庫,且不在程式碼中儲存憑證。 目標資料表必須已存在。
建立一種連結 POSTGRESQL,請參見 建立連結。 執行查詢的使用者或服務主體必須具有連線的 USE CONNECTION 權限。
要寫入 PostgreSQL 資料表,請使用 databricks.connection、 database和 dbtable 選項:
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") # Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("databricks.connection", "<connection-name>")
.option("database", "<database>")
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-columns>") // Optional
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
替換下列佔位符:
-
<connection-name>:Unity 目錄連結的名稱。 -
<database>: 目標 PostgreSQL 資料庫名稱。 -
<schema>.<table>:採用schema.table格式的現有目標資料表。 如果你省略了 schema,sink 就會使用該publicschema。 -
<primary-key-columns>:選擇性。 目標資料表主索引鍵中所有欄位的逗號分隔清單,例如id或user_id,event_type。 請參考 Upsert的行為。 -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>:查詢將其檢查點儲存在此處的 Unity Catalog 磁碟區路徑。 你也可以使用雲端物件儲存 URI。 位置必須是你可以寫入的儲存空間,而非本地磁碟,且每個串流查詢必須是獨一無二的。 這與目標表無關。 請參閱結構化串流檢查點。
PostgreSQL 連線總是使用 TLS。 憑證驗證會依照 Unity Catalog 連線上的設定,而這些設定是在你建立連線時選擇的:
-
信任伺服器憑證:選擇時,連線使用
sslmode=require,該連線在不驗證伺服器憑證的情況下加密連線。 -
使用者提供的伺服器憑證:提供一份 PEM 編碼的伺服器憑證,當未選擇信任伺服器憑證時使用
sslmode=verify-full。 如果你沒有提供憑證,連線會使用sslmode=verify-full,並搭配 JVM 預設的信任儲存庫。
設定選項
水槽因未識別選項而產生錯誤訊息。 JDBC_STREAMING_SINK_INVALID_OPTIONS
關於匯組設定選項,包括常見選項及每種連線方法的選項,請參見 PostgreSQL 匯入選項。
資料類型對應
匯入器會檢查每個 DataFrame 欄位是否與其對應的目標欄位相容,然後再寫入現有的 Lakebase 或外部 PostgreSQL 資料表。
下表包含 Databricks Runtime 18 LTS 及以上版本所支援的類型:
| Spark 類型 | 自動建立的 Lakebase 表格類型 | 現有 PostgreSQL 資料表中的相容型別 |
|---|---|---|
ByteType、ShortType |
smallint |
smallint |
IntegerType |
integer |
integer |
LongType |
bigint |
bigint |
FloatType |
real |
real |
DoubleType |
double precision |
double precision |
DecimalType |
numeric |
numeric |
StringType |
text |
varchar、text |
VarcharType(n) |
varchar(n) |
varchar、text |
CharType(n) |
char(n) |
char |
BinaryType |
bytea |
bytea |
BooleanType |
boolean |
boolean |
TimestampType |
timestamptz |
timestamptz |
TimestampNTZType |
timestamp |
timestamp |
DateType |
date |
date |
ArrayType、、 MapType、 StructType、 VariantType、 NullType |
jsonb |
json、jsonb |
下表列出 Databricks Runtime 19 及以上版本所支援的類型:
| Spark 類型 | 自動建立的 Lakebase 表格類型 | 現有 PostgreSQL 資料表中的相容型別 |
|---|---|---|
DayTimeIntervalType、YearMonthIntervalType |
interval |
interval |
Upsert 的行為
該 upsertkey 選項會識別目標資料表的主鍵欄位。 對於現有的資料表,欄位 upsertkey 必須與該資料表的主鍵完全一致。 如果你省略這個選項,匯入器會從表格讀取主金鑰。 對於 sink 建立的 Lakebase 表格, upsertkey 定義了主鍵。 如果你省略了這個選項,水槽會建立沒有主鍵的表格。
當目標資料表具有主鍵時,接收端會使用 PostgreSQL 的 INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... 語法執行 upsert。 當目標資料表沒有主鍵時,匯入器會執行插入。 查詢的輸出模式不會影響此行為。
所有主鍵欄位必須存在於資料框架中,並使用可比較的類型,例如數字型或字串型別。
性能調校
批次處理與背壓
當任一條件成立時,將觸發刷新:
- 緩衝區可達到
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 |
是的 | No |
ProcessingTime |
是的 | No |
AvailableNow |
是的 | 是的 |
Once |
Yes. Deprecated. 請使用 AvailableNow。 |
Yes. Deprecated. 請使用 AvailableNow。 |
輸出模式
此表顯示對結構化串流輸出模式的支援:
| 輸出模式 | 支援 |
|---|---|
update |
是的 |
append |
Yes. 行為與 update相同。 當目標資料表有主鍵時,查詢會上溢出,否則查詢會插入。 請參考 Upsert的行為。 |
complete |
No |
限制
- 若透過 Unity 目錄連線連接外部 PostgreSQL 資料庫,目標資料表必須已存在。 Sink 只會自動在 Lakebase 中產生缺失的資料表。
- 不支援湖流量管線。