Poznámka:
Přístup k této stránce vyžaduje autorizaci. Můžete se zkusit přihlásit nebo změnit adresáře.
Přístup k této stránce vyžaduje autorizaci. Můžete zkusit změnit adresáře.
Použijte Structured Streaming k zápisu do Lakebase nebo externí databáze PostgreSQL s vestavěným dávkováním, automatickými opakovaními a autentizací spravovanou pracovním prostorem.
Kdy použít jímku Lakebase
Používejte jímku Lakebase pro streamované zápisy s nízkou latencí do Lakebase nebo externí databáze PostgreSQL. Tento cíl nevyžaduje, abyste implementovali vlastní funkce foreach pro dávkové zpracování, správu připojení a zpracování chyb.
Mezi běžné případy použití patří:
- Aktualizujte aplikační databáze v reálném čase pro provozní řídicí panely nebo funkce určené pro zákazníky.
- Synchronizovat průběžně se měnící data, jako jsou agregované nebo filtrované výsledky streamování, do transakční databáze.
- Zapište výstup dotazu strukturovaného streamování do tabulky Lakebase s latencí podsekundy pomocí režimu v reálném čase.
Pro synchronizaci dat z Lakebase do tabulek Delta Lake v Lakehouse, tedy v opačném směru, viz Kanál změn dat Lakebase.
Požadavky
-
Databricks Runtime 18 LTS a vyšší.
- Externí připojení PostgreSQL vyžadují použití Databricks Runtime 19 a novější a přihlášení k Custom JDBC na UC Compute preview.
- Intervalové datové typy vyžadují použití Databricks Runtime 19 a vyšší.
- Klasické výpočetní prostředky s vyhrazenými nebo standardními režimy přístupu, nebo bezserverové výpočetní prostředky pro poznámkové bloky či úlohy. Pro serverless výpočty použijte
Trigger.AvailableNow(). Viz Streamování na výpočetních prostředcích bez serveru. - Databáze Lakebase nebo připojení Unity Catalog k externí databázi PostgreSQL.
Požadavky na identifikátor
Pro všechny cílové systémy Databricks doporučuje používat názvy schémat, tabulek, sloupců a sloupců primárních klíčů, které začínají písmenem nebo podtržítkem a obsahují pouze písmena, číslice a podtržítka. Sink tyto požadavky vynucuje, když automaticky vytváří tabulku Lakebase. Pro použití identifikátorů, které tyto požadavky nespĺnějí, vytvořte cílovou tabulku před zahájením dotazu.
Připojit k databázi
Jímka Lakebase podporuje následující metody připojení:
Tabulky Lakebase zaregistrované v katalogu Unity
V případě tabulek Lakebase zaregistrovaných v Unity Catalog konektor automaticky spravuje přihlašovací údaje a používá identitu uživatele nebo instančního objektu služby, který dotaz spouští. Pokud tabulka neexistuje, konektor vytvoří tabulku.
Pokud chcete zaregistrovat databázi Lakebase v katalogu Unity, přečtěte si téma Registrace databáze Lakebase v katalogu Unity.
Pro zápis do tabulky Lakebase použijte metodu .toTable() s plně kvalifikovaným názvem tabulky: 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>")
Nahraďte následující zástupné symboly:
-
<catalog>.<schema>.<table>: Plně kvalifikovaný název cílové tabulky. Jedná secatalogo katalog Unity, který jste vytvořili při registraci databáze Lakebase, viz Registrace databáze Lakebase v Katalogu Unity. Pokud tabulka neexistuje, konektor ji vytvoří. -
<primary-key-columns>: Volitelné. Čárkami oddělený seznam všech sloupců v primárním klíči cílové tabulky, napříkladidnebouser_id,event_type. Viz chování Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Cesta ke svazku katalogu Unity, kde dotaz ukládá svůj kontrolní bod. Můžete také použít identifikátor URI cloudového úložiště objektů. Umístění musí být úložiště, do kterého můžete zapisovat, nikoli na místní disk, a musí být jedinečné pro každý dotaz streamování. To je nezávislé na cílové tabulce. Viz kontrolní body strukturovaného streamování.
Volitelné konfigurace, například batchsize a batchinterval, viz Možnosti konfigurace.
Tabulky Lakebase nejsou zaregistrované v katalogu Unity
U tabulek Lakebase, které nejsou zaregistrované ve službě Unity Catalog, konektor automaticky spravuje přihlašovací údaje a používá identitu uživatele nebo instančního objektu, který dotaz spouští. Pokud tabulka neexistuje, konektor vytvoří tabulku.
Pro zápis do tabulky Lakebase použijte možnosti endpoint a dbtable:
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()
Nahraďte následující zástupné symboly:
-
<project-id>.<branch-id>.<endpoint-id>: Váš koncový bod Lakebase. Všechny tři hodnoty v názvu prostředku najdete v nabídce Získat ID na kartě Výpočty , která má formátprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Viz Identifikátory výpočetních prostředků. -
<database>: Volitelné. Název cílové databáze PostgreSQL. Výchozí hodnota jedatabricks_postgres. Viz Správa databází. -
<schema>.<table>: Cílová tabulka veschema.tableformátu. Pokud schéma vynecháte, jímka použijepublicschéma. Pro automatické vytváření tabulek používejte identifikátory začínající písmenem nebo podtržítkem a obsahující pouze písmena, čísla a podtržítka. -
<primary-key-columns>: Volitelné. Čárkami oddělený seznam všech sloupců v primárním klíči cílové tabulky, napříkladidnebouser_id,event_type. Viz chování Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Cesta ke svazku katalogu Unity, kde dotaz ukládá svůj kontrolní bod. Můžete také použít identifikátor URI cloudového úložiště objektů. Umístění musí být úložiště, do kterého můžete zapisovat, nikoli na místní disk, a musí být jedinečné pro každý dotaz streamování. To je nezávislé na cílové tabulce. Viz kontrolní body strukturovaného streamování.
Volitelné konfigurace, například batchsize a batchinterval, viz Možnosti konfigurace.
Externí PostgreSQL s přihlašovacími údaji Unity Catalog
Important
Tato funkce je ve verzi Public Preview. Správci pracovního prostoru mohou ovládat přístup k Custom JDBC na UC Compute ze stránky Previews . Viz Manage Azure Databricks preview.
Použijte připojení Unity Catalog k autentizaci k externí databázi PostgreSQL bez ukládání přihlašovacích údajů do kódu. Cílová tabulka už musí existovat.
Vytvořte spojení typu POSTGRESQL, viz Vytvořit spojení. Uživatel nebo princip služby, který dotaz provádí, musí mít USE CONNECTION na spojení.
Pro zápis do tabulky PostgreSQL použijte databricks.connection, database, a dbtable možnosti:
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()
Nahraďte následující zástupné symboly:
-
<connection-name>: Název připojení ke službě Unity Catalog. -
<database>: Název cílové databáze PostgreSQL. -
<schema>.<table>: Stávající cílová tabulka ve formátuschema.table. Pokud schéma vynecháte, jímka použijepublicschéma. -
<primary-key-columns>: Volitelné. Čárkami oddělený seznam všech sloupců v primárním klíči cílové tabulky, napříkladidnebouser_id,event_type. Viz chování Upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Cesta ke svazku katalogu Unity, kde dotaz ukládá svůj kontrolní bod. Můžete také použít identifikátor URI cloudového úložiště objektů. Umístění musí být úložiště, do kterého můžete zapisovat, nikoli na místní disk, a musí být jedinečné pro každý dotaz streamování. To je nezávislé na cílové tabulce. Viz kontrolní body strukturovaného streamování.
PostgreSQL připojení vždy používá TLS. Ověření certifikátu probíhá podle nastavení připojení v Unity Catalog, která zvolíte při vytváření připojení:
-
Trust server certificate: Po výběru spojení použije
sslmode=require, což zašifruje spojení bez ověření serverového certifikátu. -
Uživatelem poskytnutý serverový certifikát: Poskytnout PEM kódovaný serverový certifikát pro použití
sslmode=verify-full, když není vybrán certifikát Trust serveru . Pokud neposkytnete certifikát, připojení používásslmode=verify-fulls výchozím úložištěm důvěry JVM.
Možnosti konfigurace
Sink vrátí chybu při nerozpoznaných volbách, JDBC_STREAMING_SINK_INVALID_OPTIONS.
Následující možnosti platí pro všechny metody připojení:
| Key | Výchozí | Description |
|---|---|---|
batchinterval |
100 milliseconds |
Optional. Maximální doba uložení řádků do vyrovnávací paměti před vyprázdněním. Například: "50 milliseconds". |
batchsize |
1000 |
Optional. Maximální počet řádků pro každou transakci databáze. |
checkpointLocation |
None | Required. Cesta k adresáři kontrolních bodů, například ke svazku katalogu Unity (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Každý dotaz musí být jedinečný. Viz kontrolní body strukturovaného streamování. |
upsertkey |
None | Optional. Seznam všech sloupců v primárním klíči cílové tabulky oddělený čárkami, "id" například ."user_id,event_type" Viz chování Upsert. |
Tabulky Lakebase nejsou zaregistrované v katalogu Unity
Následující možnosti platí pro připojení k tabulce Lakebase, která není zaregistrovaná v katalogu Unity:
| Key | Výchozí | Description |
|---|---|---|
database |
databricks_postgres |
Optional. Název cílové databáze PostgreSQL. |
dbtable |
None | Required. Název cílové tabulky ve schema.table formátu. Pokud nezadáte schéma, výchozí hodnota schématu je public. Pro automatické vytváření tabulek používejte identifikátory začínající písmenem nebo podtržítkem a obsahující pouze písmena, čísla a podtržítka. |
endpoint |
None | Required. Koncový bod Lakebase ve formátu project_id.branch_id nebo project_id.branch_id.endpoint_id.
endpoint_id je nepovinný. Pokud ji vynecháte a větev má jeden read-write endpoint, sink tento endpoint ve výchozím nastavení vybere. |
Externí PostgreSQL s přihlašovacími údaji Unity Catalog
Následující možnosti platí při připojení k externí databázi PostgreSQL s přihlašovacími údaji Unity Catalog:
| Key | Výchozí | Description |
|---|---|---|
database |
None | Required. Název cílové databáze PostgreSQL. |
databricks.connection |
None | Required. Název spojení v Unity Catalog pro autentizaci spravovanou Unity Catalog k externímu PostgreSQL. |
dbtable |
None | Required. Stávající název cílové tabulky je ve formátu schema.table . Pokud nezadáte schéma, výchozí hodnota schématu je public. |
Mapování datových typů
Jímka před zápisem do existující tabulky Lakebase nebo externí tabulky PostgreSQL kontroluje, zda je každý sloupec objektu DataFrame kompatibilní s odpovídajícím cílovým sloupcem.
Následující tabulka obsahuje typy podporované v Databricks Runtime 18 LTS a vyšších:
| Typ Spark | Automaticky vytvářený typ tabulky Lakebase | Kompatibilní typy v existujících tabulkách 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, , VariantTypeNullType |
jsonb |
json, jsonb |
Následující tabulka obsahuje typy podporované v Databricks Runtime 19 a výše:
| Typ Spark | Automaticky vytvářený typ tabulky Lakebase | Kompatibilní typy v existujících tabulkách PostgreSQL |
|---|---|---|
DayTimeIntervalType, YearMonthIntervalType |
interval |
interval |
Chování operace upsert
Možnost upsertkey identifikuje sloupce primárního klíče cílové tabulky. U existující tabulky musí sloupce v upsertkey přesně odpovídat primárnímu klíči tabulky. Pokud tuto možnost vynecháte, dřez přečte hlavní klíč ze stolu. Pro tabulku Lakebase, kterou sink vytváří, definuje upsertkey primární klíč. Pokud tuto možnost vynecháte, dřez vytvoří stůl bez primárního klíče.
Když má cílová tabulka primární klíč, sink se vyrovná se syntaxí INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... PostgreSQL. Když cílová tabulka nemá primární klíč, dřez provádí vložení. Výstupní režim dotazu nemá na toto chování žádný vliv.
Všechny sloupce primárních klíčů musí být přítomny v DataFrame a používat srovnatelné typy, například číselné nebo řetězcové typy.
Ladění výkonu
Dávkování a zpětný tlak
Při splnění některé z podmínek se aktivuje vyprázdnění:
- Vyrovnávací paměť dosáhne velikosti
batchsizeřádků, přičemž výchozí hodnota je1000. - Stáří vyrovnávací paměti překračuje
batchinterval, což je výchozí hodnota100 milliseconds.
Pokud databáze nedokáže zpracovávat příchozí data požadovanou rychlostí, cílová komponenta propaguje zpětný tlak proti proudu ke zdroji.
Pokyny k latenci a propustnosti:
- U úloh s nízkou latencí v režimu reálného času snižte
batchinterval, abyste zajistili kratší maximální dobu před vyprázdněním. Viz koncepty režimu v reálném čase pro koncepty a příklady režimu v reálném čase pro příklad kódu. - U úloh s vysokou propustností zvyšte
batchsize, abyste snížili režii u každé transakce.
Chování připojení
Jímka používá sdružování připojení u exekutorů. Ve výchozím nastavení používá každý úkol jedno připojení k databázi.
Databricks doporučuje, abyste pro každé připojení použili výchozí hodnotu 1 úkolu. Pokud zvýšíte počet úloh pro každé připojení, můžete způsobit kolize připojení a zvýšit latenci u připojení s vysokou propustností.
Pokud chcete nakonfigurovat poměr úloh k připojením, nastavte konfiguraci Sparku spark.databricks.sql.streaming.jdbc.tasksPerConnection . Pokud má cílová databáze nízký limit připojení, snižte počet oddílů pro přeskupení nebo zvyšte spark.databricks.sql.streaming.jdbc.tasksPerConnection.
Jímka automaticky opakuje přechodné chyby JDBC, včetně selhání připojení, zablokování a omezování rychlosti. Pokud jímka vyčerpá všechny pokusy o opakování, dotaz selže.
Podporované aktivační události a výstupní režimy
Triggers
Tato tabulka ukazuje podporu typů spouštěčů strukturovaného streamování na klasickém a serverless výpočetním systému:
| Trigger | Klasické výpočetní prostředky | Bezserverové výpočty (notebooky a úlohy) |
|---|---|---|
RealTime |
Ano | No |
ProcessingTime |
Ano | No |
AvailableNow |
Ano | Ano |
Once |
Yes. Deprecated. Použijte AvailableNow. |
Yes. Deprecated. Použijte AvailableNow. |
Výstupní režimy
Tato tabulka ukazuje podporu pro režimy výstupu strukturovaného streamování:
| Výstupní režim | Podporováno |
|---|---|
update |
Ano |
append |
Yes. Chování je identické s update. Dotaz provede operaci upsert, pokud má cílová tabulka primární klíč; jinak provede vložení. Viz chování Upsert. |
complete |
No |
omezení
- Pro externí PostgreSQL databázi připojenou přes Unity Catalog musí cílová tabulka již existovat. Dřez automaticky vytváří chybějící tabulky pouze v Lakebase.
- Potrubí Lakeflow není podporováno.