Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
Gebruik Structured Streaming om te schrijven naar Lakebase of een externe PostgreSQL-database met ingebouwde batching, automatische herkansingen en door de werkruimte beheerde authenticatie.
Wanneer gebruikt u de Lakebase-sink?
Gebruik de Lakebase-sink voor stream-schrijfopdrachten met lage latentie naar Lakebase of een externe PostgreSQL-database. Deze sink vereist niet dat u aangepaste foreach-functies implementeert om batchverwerking, verbindingsbeheer en foutafhandeling af te handelen.
Veelvoorkomende gebruiksvoorbeelden zijn:
- Toepassingsdatabases in realtime bijwerken voor operationele dashboards of klantgerichte functies.
- Synchroniseer continu veranderende gegevens, zoals geaggregeerde of gefilterde streamingresultaten, in een transactionele database.
- Schrijf de uitvoer van een Structured Streaming-query naar een Lakebase-tabel met een latentie van minder dan een seconde met behulp van real-time-modus.
Als u gegevens vanuit Lakebase wilt synchroniseren met Delta Lake-tabellen in Lakehouse, raadpleegt u Lakebase Change Data Feed.
vereisten voor
-
Databricks Runtime 18 LTS en hoger.
- Externe PostgreSQL-verbindingen vereisen dat je Databricks Runtime 19 en hoger gebruikt en je aanmeldt voor de Custom JDBC op UC Compute preview.
- Gegevenstypen voor intervallen vereisen dat u Databricks Runtime 19 of hoger gebruikt.
- Klassieke rekenkracht met dedicatede of standaardtoegangsmodi, of serverless rekenkracht voor notebooks of taken. Gebruik
Trigger.AvailableNow()voor serverloze rekenkracht. Zie Streaming op serverloze rekenkracht. - Een Lakebase-database, of een Unity Catalog-verbinding met een externe PostgreSQL-database.
Identificatievereisten
Voor alle doelen raadt Databricks aan om kolomnamen voor schema's, tabels, kolommen en primaire sleutels te gebruiken die beginnen met een letter of onderstreep en alleen letters, cijfers en onderscores bevatten. De put handhaaft deze vereisten wanneer deze automatisch een Lakebase-tabel aanmaakt. Om identifiers te gebruiken die niet aan deze eisen voldoen, maak je de doel-tabel aan voordat je de query start.
Verbinding maken met een database
De Lakebase-sink ondersteunt de volgende verbindingsmethoden:
Lakebase-tabellen die zijn geregistreerd bij Unity Catalog
Voor Lakebase-tabellen die zijn geregistreerd bij Unity Catalog, beheert de connector automatisch de referenties en gebruikt de identiteit van de gebruiker of service-principal die de query uitvoert. Als de tabel niet bestaat, maakt de connector de tabel.
Als u een Lakebase-database wilt registreren bij Unity Catalog, raadpleegt u Een Lakebase-database registreren in Unity Catalog.
Om naar een Lakebase-tabel te schrijven, gebruik je de .toTable() methode met een volledig gekwalificeerde tabelnaam, 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>")
Vervang de volgende tijdelijke aanduidingen:
-
<catalog>.<schema>.<table>: De volledig gekwalificeerde naam van de doeltabel. Ditcatalogis de Unity Catalog-catalogus die u hebt gemaakt toen u de Lakebase-database registreerde. Zie Een Lakebase-database registreren in Unity Catalog. Als de tabel niet bestaat, wordt deze gemaakt door de connector. -
<primary-key-columns>: optioneel. Een komma-gescheiden lijst van alle kolommen in de primaire sleutel van de doeltabel, bijvoorbeeldidofuser_id,event_type. Zie gedrag van upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Een Unity Catalog-volumepad waar de query zijn checkpoint opslaat. U kunt ook een opslag-URI voor cloudobjecten gebruiken. De locatie moet opslag zijn waarnaar u kunt schrijven, niet naar de lokale schijf en moet uniek zijn voor elke streamingquery. Dit is onafhankelijk van de doeltabel. Zie Controlepunten voor gestructureerd streamen.
Zie batchsize voor optionele configuraties, zoals batchinterval en.
Lakebase-tabellen die niet zijn geregistreerd bij Unity Catalog
Voor Lakebase-tabellen die niet zijn geregistreerd bij Unity Catalog, beheert de connector automatisch de referenties en gebruikt de identiteit van de gebruiker of service-principal die de query uitvoert. Als de tabel niet bestaat, maakt de connector de tabel.
Om naar een Lakebase-tabel te schrijven, gebruik je de endpoint en-opties 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()
Vervang de volgende tijdelijke aanduidingen:
-
<project-id>.<branch-id>.<endpoint-id>: uw Lakebase-eindpunt. Zoek alle drie de waarden in de resourcenaam in het menu Id ophalen van het tabblad Berekeningen , die de indelingprojects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>heeft. Zie Compute-identificatoren. -
<database>: optioneel. De naam van de doel-PostgreSQL-database. Wordt standaard ingesteld opdatabricks_postgres. Zie Databases beheren. -
<schema>.<table>: De doeltabel in de indeling vanschema.table. Als u het schema weglaat, gebruikt de sink hetpublicschema. Voor automatische tabelcreatie gebruik je identificaties die beginnen met een letter of onderstreep en alleen letters, cijfers en onderstreepjes bevatten. -
<primary-key-columns>: optioneel. Een komma-gescheiden lijst van alle kolommen in de primaire sleutel van de doeltabel, bijvoorbeeldidofuser_id,event_type. Zie gedrag van upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Een Unity Catalog-volumepad waar de query zijn checkpoint opslaat. U kunt ook een opslag-URI voor cloudobjecten gebruiken. De locatie moet opslag zijn waarnaar u kunt schrijven, niet naar de lokale schijf en moet uniek zijn voor elke streamingquery. Dit is onafhankelijk van de doeltabel. Zie Controlepunten voor gestructureerd streamen.
Zie batchsize voor optionele configuraties, zoals batchinterval en.
Externe PostgreSQL met Unity Catalog-inloggegevens
Important
Deze functie bevindt zich in openbare preview-versie. Workspace-beheerders kunnen de toegang tot Custom JDBC op UC Compute beheren via de Previews-pagina . Zie Azure Databricks previews beheren.
Gebruik een Unity Catalog-verbinding om te authenticeren bij een externe PostgreSQL-database zonder inloggegevens in je code op te slaan. De doel-tabel moet al bestaan.
Maak een verbinding van type POSTGRESQL, zie Maak een verbinding aan. De gebruiker of service-principal die de query uitvoert, moet beschikken over USE CONNECTION voor de verbinding.
Om naar de PostgreSQL-tabel te schrijven, gebruik je de databricks.connection, database, en dbtable opties:
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()
Vervang de volgende tijdelijke aanduidingen:
-
<connection-name>: De naam van de Unity Catalog-verbinding. -
<database>: De naam van de doel-PostgreSQL-database. -
<schema>.<table>: De bestaande doeltabel in formaatschema.table. Als u het schema weglaat, gebruikt de sink hetpublicschema. -
<primary-key-columns>: optioneel. Een komma-gescheiden lijst van alle kolommen in de primaire sleutel van de doeltabel, bijvoorbeeldidofuser_id,event_type. Zie gedrag van upsert. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: Een Unity Catalog-volumepad waar de query zijn checkpoint opslaat. U kunt ook een opslag-URI voor cloudobjecten gebruiken. De locatie moet opslag zijn waarnaar u kunt schrijven, niet naar de lokale schijf en moet uniek zijn voor elke streamingquery. Dit is onafhankelijk van de doeltabel. Zie Controlepunten voor gestructureerd streamen.
PostgreSQL-verbindingen gebruiken altijd TLS. Certificaatverificatie volgt de instellingen op de Unity Catalog-verbinding, die je kiest wanneer je de verbinding aanmaakt:
-
Trust server-certificaat: Wanneer geselecteerd, gebruikt
sslmode=requirede verbinding , dat de verbinding versleutelt zonder het servercertificaat te verifiëren. - Door de gebruiker verstrekt servercertificaat: Geef een PEM-gecodeerd servercertificaat om te gebruiken
sslmode=verify-fullwanneer het Trust-servercertificaat niet is geselecteerd. Als je geen certificaat opgeeft, maakt de verbinding gebruik vansslmode=verify-fullmet de standaard truststore van de JVM.
Configuratieopties
De sink meldt een fout voor niet-herkende opties, JDBC_STREAMING_SINK_INVALID_OPTIONS.
De volgende opties zijn van toepassing op alle verbindingsmethoden:
| Key | Default | Description |
|---|---|---|
batchinterval |
100 milliseconds |
Optional. De maximale tijd om rijen in de buffer te houden voordat deze wordt geleegd. Bijvoorbeeld: "50 milliseconds". |
batchsize |
1000 |
Optional. Het maximum aantal rijen voor elke databasetransactie. |
checkpointLocation |
Geen | Required. Pad naar een controlepuntmap, zoals een Unity Catalog-volume (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Moet uniek zijn voor elke query. Zie Controlepunten voor gestructureerd streamen. |
upsertkey |
Geen | Optional. Een komma-gescheiden lijst van alle kolommen in de primaire sleutel van de doeltabel, bijvoorbeeld, "id" of "user_id,event_type". Zie gedrag van upsert. |
Lakebase-tabellen die niet zijn geregistreerd bij Unity Catalog
De volgende opties zijn van toepassing wanneer u verbinding maakt met een Lakebase-tabel die niet is geregistreerd bij Unity Catalog:
| Key | Default | Description |
|---|---|---|
database |
databricks_postgres |
Optional. De naam van de PostgreSQL-doeldatabase. |
dbtable |
Geen | Required. De naam van de doeltabel in de indeling schema.table. Als u geen schema opgeeft, is publicde standaardschemawaarde . Voor automatische tabelcreatie gebruik je identificaties die beginnen met een letter of onderstreep en alleen letters, cijfers en onderstreepjes bevatten. |
endpoint |
Geen | Required. Het Lakebase-eindpunt, in project_id.branch_id of project_id.branch_id.endpoint_id indeling. De endpoint_id optie is optioneel. Als je het weglaat en de tak één lees-schrijf eindpunt heeft, selecteert de sink standaard dat eindpunt. |
Externe PostgreSQL met Unity Catalog-inloggegevens
De volgende opties gelden wanneer je verbinding maakt met een externe PostgreSQL-database met Unity Catalog-inloggegevens:
| Key | Default | Description |
|---|---|---|
database |
Geen | Required. De naam van de PostgreSQL-doeldatabase. |
databricks.connection |
Geen | Required. De Unity Catalog-verbindingsnaam voor Unity Catalog-beheerde authenticatie naar externe PostgreSQL. |
dbtable |
Geen | Required. De bestaande naam van de doeltabel, in schema.table formaat. Als u geen schema opgeeft, is publicde standaardschemawaarde . |
Gegevenstypetoewijzingen
De sink controleert of elke DataFrame-kolom compatibel is met de bijbehorende doelkolom voordat deze naar een bestaande Lakebase- of externe PostgreSQL-tabel wordt geschreven.
De volgende tabel bevat typen die worden ondersteund in Databricks Runtime 18 LTS en hoger:
| Sparktype | Automatisch aangemaakt Lakebase-tabeltype | Compatibele types in bestaande PostgreSQL-tabellen |
|---|---|---|
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 |
De volgende tabel bevat typen die worden ondersteund in Databricks Runtime 19 en hoger:
| Sparktype | Automatisch aangemaakt Lakebase-tabeltype | Compatibele types in bestaande PostgreSQL-tabellen |
|---|---|---|
DayTimeIntervalType, YearMonthIntervalType |
interval |
interval |
Upsert gedrag
De upsertkey optie identificeert de primaire sleutelkolommen van de doeltabel. Voor een bestaande tabel moeten de kolommen in upsertkey exact overeenkomen met de primaire sleutel van de tabel. Als je de optie weglaat, leest de sink de primaire sleutel uit de tabel. Voor een Lakebase-tabel die de sink maakt, definieert upsertkey de primaire sleutel. Als u de optie weglaat, maakt de sink de tabel zonder primaire sleutel.
Wanneer de doeltabel een primaire sleutel heeft, voert de sink een upsert uit met de PostgreSQL-syntaxis INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET .... Wanneer de doeltabel geen primaire sleutel heeft, voert de sink invoegingen uit. De uitvoermodus van een query heeft geen effect op dit gedrag.
Alle primaire sleutelkolommen moeten aanwezig zijn in het DataFrame en vergelijkbare typen gebruiken, zoals numerieke of stringtypes.
Prestatieafstemming
Batchverwerking en tegendruk
Het leegmaken wordt uitgevoerd wanneer aan een van beide voorwaarden wordt voldaan:
- De buffer gaat tot
batchsizerijen, met als standaardwaarde1000. - De bufferleeftijd overschrijdt
batchinterval, die standaard is ingesteld op100 milliseconds.
Wanneer de database de snelheid van de binnenkomende gegevensstroom niet kan bijhouden, geeft de sink stroomopwaarts tegendruk door aan de bron.
Richtlijnen voor latentie en doorvoer:
- Voor workloads met lage latentie in realtime-modus verlaagt u
batchintervalzodat de maximale tijd vóór het wegschrijven korter wordt. Zie real-time mode concepten voor concepten en Real-time mode voorbeelden voor een codevoorbeeld. - Voor workloads met hoge doorvoer verhoogt u
batchsizeom de overhead voor elke transactie te verlagen.
Verbindingsgedrag
De sink gebruikt verbindingspooling bij uitvoerders. Elke taak maakt standaard gebruik van één databaseverbinding.
Databricks raadt u aan om voor elke verbinding de standaardwaarde van de taak 1 te gebruiken. Als u het aantal taken voor elke verbinding verhoogt, kan dit leiden tot conflicten tussen verbindingen en hogere latenties voor verbindingen met hoge doorvoer.
Als u de verhouding tussen taken en verbindingen wilt configureren, stelt u de spark.databricks.sql.streaming.jdbc.tasksPerConnection Spark-configuratie in. Als de doeldatabase een lage verbindingslimiet heeft, verminder dan het aantal shuffle-partities of verhoog spark.databricks.sql.streaming.jdbc.tasksPerConnection.
De sink probeert automatisch opnieuw bij tijdelijke JDBC-fouten, waaronder verbindingsfouten, deadlocks en rate limiting. Als de sink alle nieuwe pogingen uitput, mislukt de query.
Ondersteunde triggers en uitvoermodi
Triggers
Deze tabel toont ondersteuning voor Structured Streaming triggertypes op klassieke en serverless compute:
| Trigger | Klassieke rekenkracht | Serverless rekenkracht (notebooks en taken) |
|---|---|---|
RealTime |
Yes | No |
ProcessingTime |
Yes | No |
AvailableNow |
Yes | Yes |
Once |
Ja. Deprecated. Gebruik AvailableNow. |
Ja. Deprecated. Gebruik AvailableNow. |
Uitvoermodi
In deze tabel ziet u ondersteuning voor modi voor structured streaming-uitvoer:
| Uitvoermodus | Supported |
|---|---|
update |
Yes |
append |
Ja. Gedrag is identiek aan update. De query voert een upsert uit wanneer de doeltabel een primaire sleutel heeft; anders voert de query een insert uit. Zie gedrag van upsert. |
complete |
No |
Beperkingen
- Voor een externe PostgreSQL-database die via een Unity Catalog-verbinding is verbonden, moet de doeltabel al bestaan. De sink maakt ontbrekende tabellen alleen automatisch aan in Lakebase.
- Lakeflow-pijpleidingen worden niet ondersteund.