Připojení k Lakebase

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á se catalog o 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říklad id nebo user_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át projects/<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 je databricks_postgres. Viz Správa databází.
  • <schema>.<table>: Cílová tabulka ve schema.table formátu. Pokud schéma vynecháte, jímka použije public sché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říklad id nebo user_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átu schema.table . Pokud schéma vynecháte, jímka použije public schéma.
  • <primary-key-columns>: Volitelné. Čárkami oddělený seznam všech sloupců v primárním klíči cílové tabulky, například id nebo user_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-full s 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 je 1000.
  • Stáří vyrovnávací paměti překračuje batchinterval, což je výchozí hodnota 100 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.