Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
Important
Den här funktionen finns som allmänt tillgänglig förhandsversion.
Använd Structured Streaming för att skriva till Lakebase med inbyggd batchbearbetning, automatiska återförsök och arbetsytehanterad autentisering.
När du ska använda Lakebase sink
Använd Lakebase-sinken för strömmande skrivningar med låg latens till Lakebase. Den här sinken kräver inte att du implementerar egna foreachBatch-funktioner för att hantera batchning, anslutningshantering och felhantering.
Vanliga användningsfall är:
- Uppdatera programdatabaser i realtid för operativa instrumentpaneler eller kundinriktade funktioner.
- Synkronisera kontinuerligt föränderliga data, till exempel aggregerade eller filtrerade direktuppspelningsresultat, till en transaktionsdatabas.
- Skriv utdata från en fråga för strukturerad direktuppspelning till en Lakebase-tabell med svarstid under sekund med realtidsläge.
Om du vill synkronisera data från Lakebase till Delta Lake-tabeller i Lakehouse, den omvända riktningen, se Lakebase Change Data Feed.
krav
- Databricks Runtime 18 och senare
- Klassisk beräkning med dedikerade eller standardåtkomstlägen.
- En Lakebase-databas
Anslut till en databas
Lakebase-mottagaren har stöd för följande anslutningsmetoder:
Lakebase-tabeller registrerade med Unity Catalog
För Lakebase-tabeller som registrerats med Unity Catalog hanterar anslutningsappen automatiskt autentiseringsuppgifterna och använder identiteten för användaren eller tjänstens huvudnamn som kör frågan. Om tabellen inte finns skapar anslutningsappen tabellen.
Information om hur du registrerar en Lakebase-databas med Unity Catalog finns i Registrera en Lakebase-databas i Unity Catalog.
Om du vill skriva till en Lakebase-tabell använder du .toTable() metoden med ett fullständigt kvalificerat tabellnamn, catalog.schema.table. I följande exempel visas de alternativ som krävs, plus det valfria upsertkey alternativet:
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>")
Ersätt följande platshållare:
-
<catalog>.<schema>.<table>: Måltabellens fullständiga kvalificerade namn.catalogär den Unity Catalog-katalog som du skapade när du registrerade Lakebase-databasen. Mer information finns i Registrera en Lakebase-databas i Unity Catalog. Om tabellen inte finns skapar anslutningsappen den. -
<primary-key-column>: Valfritt. En kommaavgränsad lista över kolumnerna som utgör upsert-nyckeln, till exempelidelleruser_id,event_type. Om du utelämnarupsertkeyhärleder mottagaren nyckeln från måltabellens primära nyckel. Se Upsert-beteende. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: En volymsökväg i Unity Catalog där frågan lagrar sin kontrollpunkt. Du kan också använda en lagrings-URI för molnobjekt. Platsen måste vara lagring som du kan skriva till, inte lokal disk, och måste vara unik för varje strömmande fråga. Detta är oberoende av måltabellen. Se Kontrollpunkter för strukturerad strömning.
Valfria konfigurationer, till exempel batchsize och batchinterval, finns i Konfigurationsalternativ.
Lakebase-tabeller som inte har registrerats med Unity Catalog
För Lakebase-tabeller som inte har registrerats med Unity Catalog hanterar anslutningsappen automatiskt autentiseringsuppgifterna och använder identiteten för användaren eller tjänstens huvudnamn som kör frågan. Om tabellen inte finns skapar anslutningsappen tabellen.
Om du vill skriva till en Lakebase-tabell använder du endpoint alternativen och dbtable . I följande exempel ingår även valfria database alternativ och upsertkey alternativ:
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()
Ersätt följande platshållare:
-
<project-id>.<branch-id>.<endpoint-id>: Din Lakebase-slutpunkt. Hitta alla tre värdena i resursnamnet på menyn Hämta ID på fliken Beräkningar , som har formatetprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Se Beräkningsidentifierare. -
<database>: Valfritt. Namnet på postgres-måldatabasen. Standardinställningen ärdatabricks_postgres. Se Hantera databaser. -
<schema>.<table>: Måltabellen i formatetschema.table. Om du utelämnar schemat använder mottagaren schematpublic. Använd enkla identifierare som börjar med en bokstav eller ett understreck och som endast innehåller bokstäver, siffror och understreck. Citerade identifierare och specialtecken, till exempel bindestreck, stöds inte. -
<primary-key-column>: Valfritt. En kommaavgränsad lista över kolumnerna som utgör upsert-nyckeln, till exempelidelleruser_id,event_type. Om du utelämnarupsertkeyhärleder mottagaren nyckeln från måltabellens primära nyckel. Se Upsert-beteende. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: En volymsökväg i Unity Catalog där frågan lagrar sin kontrollpunkt. Du kan också använda en lagrings-URI för molnobjekt. Platsen måste vara lagring som du kan skriva till, inte lokal disk, och måste vara unik för varje strömmande fråga. Detta är oberoende av måltabellen. Se Kontrollpunkter för strukturerad strömning.
Valfria konfigurationer, till exempel batchsize och batchinterval, finns i Konfigurationsalternativ.
Konfigurationsalternativ
Sinken ger ett fel för okända alternativ, JDBC_STREAMING_SINK_INVALID_OPTIONS.
Följande alternativ gäller för alla anslutningsmetoder:
| Key | Förinställning | Description |
|---|---|---|
batchinterval |
100 milliseconds |
Optional. Den maximala tiden för att behålla rader i bufferten innan bufferten töms. Till exempel "50 milliseconds". |
batchsize |
1000 |
Optional. Det maximala antalet rader för varje databastransaktion. |
checkpointLocation |
None | Required. Sökväg till en kontrollpunktskatalog, till exempel en Unity Catalog-volym (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Måste vara unikt för varje fråga. Se Kontrollpunkter för strukturerad strömning. |
upsertkey |
None | Optional. En kommaavgränsad lista med kolumnnamn som utgör upsert-nyckeln. Exempel: "id" eller "user_id,event_type". Om du anger upsertkeymåste kolumnerna matcha tabellens primära nyckel, annars misslyckas frågan. Om du utelämnar det, använder sinken primärnyckeln automatiskt. Mer information finns i Upsert-beteende. |
Lakebase-tabeller som inte har registrerats med Unity Catalog
Följande alternativ gäller när du ansluter till en Lakebase-tabell som inte är registrerad i Unity Catalog:
| Key | Förinställning | Description |
|---|---|---|
database |
databricks_postgres |
Optional. PostgreSQL-måldatabasnamnet. |
dbtable |
None | Required. Namnet på måltabellen i formatet schema.table. Om du inte anger något schema är publicstandardschemavärdet . Använd enkla identifierare som börjar med en bokstav eller ett understreck och som endast innehåller bokstäver, siffror och understreck. Citera inte tabell- eller schemanamn. Citerade identifierare och namn med specialtecken, till exempel bindestreck, stöds inte. |
endpoint |
None | Required. Lakebase-slutpunkten, i project_id.branch_id eller project_id.branch_id.endpoint_id format.
endpoint_id är valfritt; om du utelämnar det och grenen har en enda läs- och skrivbar slutpunkt väljer sinken den slutpunkten som standard. |
Beteende för upsert
När upsert-nycklar finns, antingen angivna med upsertkey eller härledda av sinken från tabellens primärnycklar, gör sinken en upsert i tabellen med PostgreSQLs syntax INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ....
När det inte finns några upsert-nycklar gör sänkmålet infogningar. En frågas utdataläge har ingen effekt på upsert- eller insert-beteendet.
Kolumnerna upsertkey måste:
- Vara en icke-tom delmängd av DataFrame-kolumnerna.
- Matcha måltabellens
PRIMARY KEYexakt. Om de kolumner du anger inte matchar den primära nyckeln misslyckas frågan. - Vara jämförbara typer, till exempel numeriska typer eller strängtyper. För att förhindra dödlägen i databasen vid samtidiga skrivningar sorterar sinken rader efter upsertnyckel i varje batch. Upsert-nycklar stöder inte komplexa eller struct-typer.
Kolumnnamn omges automatiskt av PostgreSQLs standard, dubbla citationstecken ", vilket hanterar reserverade nyckelord och namn med blandade versaler och gemener.
Tabell- och schemanamn måste använda enkla identifierare som börjar med en bokstav eller ett understreck och som endast innehåller bokstäver, siffror och understreck. Mottagaren stöder inte citerade identifierare eller specialtecken, till exempel bindestreck, i tabell- eller schemanamn.
Prestandaoptimering
Batchbearbetning och ryggtryck
En tömning utlöses när något av villkoren uppfylls:
- Bufferten når upp till
batchsizerader, där standardvärdet är1000. - Buffertåldern överskrider
batchinterval, vilket är standardvärdet100 milliseconds.
När databasen inte kan hänga med i den inkommande datahastigheten sprider mottagaren tillbakatryck uppströms till källan.
Vägledning för svarstid och dataflöde:
- För arbetsbelastningar med låg latens med realtidsläge minskar du
batchintervalför att garantera en kortare maximal tid före tömning. Se begrepp i realtidsläge för koncept och exempel på realtidsläge för ett kodexempel. - För arbetsbelastningar med hög genomströmning, öka
batchsizeför att minska overheaden för varje transaktion.
Anslutningsbeteende
Sinken använder anslutningspoolning på executorer. Som standard använder varje aktivitet en databasanslutning.
Databricks rekommenderar att du använder standardvärdet 1 för aktiviteten för varje anslutning. Om du ökar antalet uppgifter för varje anslutning kan du orsaka anslutningskonkurrens och öka svarstiderna för anslutningar med högt dataflöde.
Om du vill konfigurera förhållandet mellan aktiviteter och anslutningar anger du spark.databricks.sql.streaming.jdbc.tasksPerConnection Spark-konfigurationen. Om måldatabasen har en låg anslutningsgräns kan du minska antalet shuffle-partitioner eller öka spark.databricks.sql.streaming.jdbc.tasksPerConnection.
Sinken försöker automatiskt igen vid tillfälliga fel i JDBC, inklusive anslutningsfel, dödlägen och frekvensbegränsning. Om sinken har förbrukat alla återförsök, misslyckas frågan.
Utlösare och utdatalägen som stöds
Triggers
Den här tabellen visar stöd för utlösartyper för strukturerad direktuppspelning:
| Trigger | Supported |
|---|---|
realTime |
Yes |
ProcessingTime |
Yes |
AvailableNow |
Yes |
Once |
Yes |
Utdatalägen
Den här tabellen visar stöd för utdatalägen för strukturerad direktuppspelning:
| Utdataläge | Supported |
|---|---|
update |
Yes |
append |
Yes. Beteendet är identiskt med update. Frågan utför en upsert när måltabellen har en primärnyckel, annars infogar den. Se Upsert-beteende. |
complete |
No |
Begränsningar
- Serverlös beräkning och Lakeflow-pipelines stöds inte.
- Endast Lakebase stöds som skrivmål. Externa PostgreSQL-kompatibla databaser stöds inte.