Verbinding maken met Lakebase

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. Dit catalog is 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, bijvoorbeeld id of user_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 indeling projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>heeft. Zie Compute-identificatoren.
  • <database>: optioneel. De naam van de doel-PostgreSQL-database. Wordt standaard ingesteld op databricks_postgres. Zie Databases beheren.
  • <schema>.<table>: De doeltabel in de indeling van schema.table. Als u het schema weglaat, gebruikt de sink het public schema. 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, bijvoorbeeld id of user_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 formaat schema.table . Als u het schema weglaat, gebruikt de sink het public schema.
  • <primary-key-columns>: optioneel. 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.
  • /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-full wanneer het Trust-servercertificaat niet is geselecteerd. Als je geen certificaat opgeeft, maakt de verbinding gebruik van sslmode=verify-full met 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 batchsize rijen, met als standaardwaarde 1000.
  • De bufferleeftijd overschrijdt batchinterval, die standaard is ingesteld op 100 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 batchinterval zodat 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 batchsize om 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.