Connettersi a Lakebase

Usa Structured Streaming per scrivere su Lakebase o su un database esterno PostgreSQL con batching integrato, tentativi automatici e autenticazione gestita dall'area di lavoro.

Quando usare il sink Lakebase

Usa il sink di Lakebase per scritture in streaming a bassa latenza su Lakebase o su un database esterno PostgreSQL. Questo sink non richiede l'implementazione di funzioni personalizzate foreach per gestire l'invio in batch, la gestione delle connessioni e la gestione degli errori.

I casi d'uso comuni includono:

  • Aggiornare i database dell'applicazione in tempo reale per dashboard operativi o funzionalità rivolte ai clienti.
  • Sincronizzare in modo continuo i dati, ad esempio i risultati di streaming aggregati o filtrati, in un database transazionale.
  • Scrivi l'output di una query di Structured Streaming in una tabella Lakebase con latenza inferiore al secondo usando la modalità in tempo reale.

Per sincronizzare i dati da Lakebase verso le tabelle Delta Lake nel Lakehouse, nella direzione inversa, vedere Lakebase Change Data Feed.

Requisiti

  • Databricks Runtime 18 LTS e versioni successive.
    • Le connessioni PostgreSQL esterne richiedono l'utilizzo di Databricks Runtime 19 o versioni successive e l'adesione all'anteprima di Custom JDBC on UC Compute.
    • I tipi di dati di intervallo richiedono Databricks Runtime 19 o versione successiva.
  • Calcolo classico con modalità di accesso dedicate o standard, oppure calcolo serverless per notebook o lavori. Sul calcolo serverless, si usa Trigger.AvailableNow(). Consulta Streaming sull'elaborazione serverless.
  • Un database Lakebase, o una connessione Unity Catalog a un database esterno PostgreSQL.

Requisiti per l'identificazione

Per tutti i target, Databricks raccomanda di utilizzare nomi di schema, tabelle, colonne e colonne di chiave primaria che iniziano con una lettera o sottolineamento e contengono solo lettere, numeri e sottolineature. Il sink fa rispettare questi requisiti quando crea automaticamente una tabella Lakebase. Per utilizzare identificatori che non soddisfano questi requisiti, crea la tabella target prima di iniziare la query.

Connettersi a un database

Il sink Lakebase supporta i metodi di connessione seguenti:

Tabelle Lakebase registrate con Unity Catalog

Per le tabelle Lakebase registrate con Unity Catalog, il connettore gestisce automaticamente le credenziali e utilizza l'identità dell'utente o del service principal che esegue la query. Se la tabella non esiste, il connettore crea la tabella.

Per registrare un database Lakebase con Il catalogo Unity, vedere Registrare un database Lakebase nel catalogo unity.

Per scrivere su una tabella Lakebase, usa il .toTable() metodo con un nome di tabella completamente qualificato, 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>")

Sostituire i segnaposto seguenti:

  • <catalog>.<schema>.<table>: Il nome completo della tabella di destinazione. catalog è il catalogo di Unity Catalog che hai creato quando hai registrato il database Lakebase; vedi Registrare un database Lakebase in Unity Catalog. Se la tabella non esiste, il connettore lo crea.
  • <primary-key-columns>: facoltativo. Un elenco separato da virgole di tutte le colonne nella chiave primaria della tabella di destinazione, ad esempio id o user_id,event_type. Vedi comportamento di Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: percorso di un volume di Unity Catalog in cui la query archivia il checkpoint. È anche possibile usare un URI di archiviazione di oggetti cloud. La posizione deve essere un'unità di archiviazione su cui sia possibile scrivere, non un disco locale, e deve essere univoca per ogni query di streaming. Questo è indipendente dalla tabella di destinazione. Consulta Checkpoint di Structured Streaming.

Per le configurazioni facoltative, ad esempio batchsize e batchinterval, vedere Opzioni di configurazione.

Tabelle Lakebase non registrate con Unity Catalog

Per le tabelle Lakebase non registrate in Unity Catalog, il connettore gestisce automaticamente le credenziali e usa l'identità dell'utente o dell'entità servizio che esegue la query. Se la tabella non esiste, il connettore crea la tabella.

Per scrivere su una tabella Lakebase, usa le opzioni endpoint e 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()

Sostituire i segnaposto seguenti:

  • <project-id>.<branch-id>.<endpoint-id>: Il tuo endpoint Lakebase. Trovare tutti e tre i valori nel nome della risorsa nel menu Recupera ID della scheda Calcolo , che ha il formato projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>. Vedere Identificatori di calcolo.
  • <database>: facoltativo. Il nome del database postgreSQL di destinazione. Di default è databricks_postgres. Vedere Gestire i database.
  • <schema>.<table>: La tabella di destinazione in formato schema.table. Se si omette lo schema, il sink usa lo public schema . Per la creazione automatica di tabelle, usa identificatori che iniziano con una lettera o un sottolineamento e contengono solo lettere, numeri e sottolineature.
  • <primary-key-columns>: facoltativo. Un elenco separato da virgole di tutte le colonne nella chiave primaria della tabella di destinazione, ad esempio id o user_id,event_type. Vedi comportamento di Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: percorso di un volume di Unity Catalog in cui la query archivia il checkpoint. È anche possibile usare un URI di archiviazione di oggetti cloud. La posizione deve essere un'unità di archiviazione su cui sia possibile scrivere, non un disco locale, e deve essere univoca per ogni query di streaming. Questo è indipendente dalla tabella di destinazione. Consulta Checkpoint di Structured Streaming.

Per le configurazioni facoltative, ad esempio batchsize e batchinterval, vedere Opzioni di configurazione.

PostgreSQL esterno con credenziali del Catalogo Unity

Important

Questa funzionalità è in Anteprima Pubblica. Gli amministratori dello spazio di lavoro possono controllare l'accesso a Custom JDBC su UC Compute dalla pagina delle anteprime . Vedere Gestire le anteprime di Azure Databricks.

Usa una connessione Unity Catalog per autenticarti a un database PostgreSQL esterno senza memorizzare credenziali nel tuo codice. La tabella di destinazione deve già esistere.

Crea una connessione di tipo POSTGRESQL, vedi Crea una connessione. L'utente o il principale del servizio che esegue la query deve avere USE CONNECTION sulla connessione.

Per scrivere nella tabella PostgreSQL, usa le databricks.connectionopzioni , database, e dbtable :

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()

Sostituire i segnaposto seguenti:

  • <connection-name>: Il nome della connessione di Unity Catalog.
  • <database>: Il nome del database postgreSQL di destinazione.
  • <schema>.<table>: La tabella di destinazione esistente in formato schema.table. Se si omette lo schema, il sink usa lo public schema .
  • <primary-key-columns>: facoltativo. Un elenco separato da virgole di tutte le colonne nella chiave primaria della tabella di destinazione, ad esempio id o user_id,event_type. Vedi comportamento di Upsert.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: percorso di un volume di Unity Catalog in cui la query archivia il checkpoint. È anche possibile usare un URI di archiviazione di oggetti cloud. La posizione deve essere un'unità di archiviazione su cui sia possibile scrivere, non un disco locale, e deve essere univoca per ogni query di streaming. Questo è indipendente dalla tabella di destinazione. Consulta Checkpoint di Structured Streaming.

Le connessioni PostgreSQL usano sempre TLS. La verifica del certificato segue le impostazioni sulla connessione Unity Catalog, che scegli quando crei la connessione:

  • Certificato server fiduciario: Quando selezionato, la connessione utilizza sslmode=require, che cripta la connessione senza verificare il certificato server.
  • Certificato server fornito dall'utente: Fornire un certificato server codificato in PEM da utilizzare sslmode=verify-full quando il certificato Trust Server non è selezionato. Se non fornisci un certificato, la connessione viene utilizzata sslmode=verify-full con il trust store predefinito della JVM.

Opzioni di configurazione

Il sink segnala un errore in presenza di opzioni non riconosciute, JDBC_STREAMING_SINK_INVALID_OPTIONS.

Le opzioni seguenti si applicano a tutti i metodi di connessione:

Key Default Description
batchinterval 100 milliseconds Optional. Tempo massimo di memorizzazione delle righe nel buffer prima dello scaricamento. Ad esempio: "50 milliseconds".
batchsize 1000 Optional. Numero massimo di righe per ogni transazione di database.
checkpointLocation Nessuno Required. Percorso di una directory di checkpoint, ad esempio un volume di Unity Catalog (/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>). Deve essere univoco per ogni query. Consulta Checkpoint di Structured Streaming.
upsertkey Nessuno Optional. Una lista separata da virgole di tutte le colonne nella chiave primaria della tabella di destinazione, ad esempio, "id" o "user_id,event_type". Vedi comportamento di Upsert.

Tabelle Lakebase non registrate con Unity Catalog

Quando ci si connette a una tabella Lakebase non registrata in Unity Catalog, si applicano le opzioni seguenti:

Key Default Description
database databricks_postgres Optional. Nome del database PostgreSQL di destinazione.
dbtable Nessuno Required. Nome della tabella di destinazione in schema.table formato. Se non si specifica uno schema, il valore dello schema predefinito è public. Per la creazione automatica di tabelle, usa identificatori che iniziano con una lettera o un sottolineamento e contengono solo lettere, numeri e sottolineature.
endpoint Nessuno Required. L'endpoint Lakebase, nel formato project_id.branch_id o project_id.branch_id.endpoint_id. Il endpoint_id è facoltativo. Se lo ometti e il branch ha un unico endpoint di lettura e scrittura, il sink seleziona automaticamente quell'endpoint per impostazione predefinita.

PostgreSQL esterno con credenziali del Catalogo Unity

Le seguenti opzioni si applicano quando si connette a un database PostgreSQL esterno con credenziali del Catalogo Unity:

Key Default Description
database Nessuno Required. Nome del database PostgreSQL di destinazione.
databricks.connection Nessuno Required. Il nome della connessione del Catalogo Unity per l'autenticazione gestita dal Catalogo Unity verso PostgreSQL esterno.
dbtable Nessuno Required. Il nome della tabella target esistente in schema.table formato. Se non si specifica uno schema, il valore dello schema predefinito è public.

Mappatura dei tipi di dati

Il sink verifica che ogni colonna DataFrame sia compatibile con la colonna target corrispondente prima di scrivere su una Lakebase esistente o su una tabella esterna PostgreSQL.

La tabella seguente contiene i tipi supportati in Databricks Runtime 18 LTS e superiori:

Tipo Spark Tipo di tabella Lakebase creato automaticamente Tipi compatibili nelle tabelle PostgreSQL esistenti
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

La seguente tabella contiene i tipi supportati in Databricks Runtime 19 e superiori:

Tipo Spark Tipo di tabella Lakebase creato automaticamente Tipi compatibili nelle tabelle PostgreSQL esistenti
DayTimeIntervalType, YearMonthIntervalType interval interval

Comportamento da Upsert

L'opzione upsertkey identifica le colonne chiave principali della tabella di destinazione. Per una tabella esistente, le colonne in upsertkey devono corrispondere esattamente alla chiave primaria della tabella. Se ometti questa opzione, il sink legge la chiave primaria dalla tabella. Per una tabella Lakebase creata dal sink, upsertkey definisce la chiave primaria. Se ometti questa opzione, il sink crea la tabella senza chiave primaria.

Quando la tabella target ha una chiave primaria, il sink esegue un upsert utilizzando la sintassi INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... di PostgreSQL. Quando la tabella di destinazione non ha una chiave primaria, il sink esegue inserimenti. La modalità di output di una query non ha alcun effetto su questo comportamento.

Tutte le colonne chiave primarie devono essere presenti nel DataFrame e utilizzare tipi comparabili, come tipi numerici o di stringhe.

Messa a punto delle prestazioni

Invio in batch e backpressione

Lo svuotamento viene eseguito quando è soddisfatta una delle due condizioni:

  • Il buffer arriva a batchsize righe, il cui valore predefinito è 1000.
  • L'età del buffer supera batchinterval, che per impostazione predefinita è 100 milliseconds.

Quando il database non riesce a tenere il passo con il tasso di dati in ingresso, il sink propaga la backpressure a monte, fino alla sorgente.

Indicazioni sulla latenza e sulla velocità effettiva:

  • Per i carichi di lavoro a bassa latenza in modalità in tempo reale, diminuire batchinterval per garantire un intervallo massimo più breve prima dello svuotamento. Vedi concetti di modalità in tempo reale per i concetti e esempi di modalità in tempo reale per un esempio di codice.
  • Per i carichi di lavoro a throughput elevato, aumentare batchsize per ridurre il sovraccarico di ogni transazione.

Comportamento della connessione

Il "sink" utilizza un pool di connessioni sugli executor. Per impostazione predefinita, ogni attività usa una connessione al database.

Databricks consiglia di usare il valore predefinito dell'attività 1 per ogni connessione. Se si aumenta il numero di attività per ogni connessione, è possibile che si verifichino conflitti di connessione e si aumentino le latenze per le connessioni a velocità effettiva elevata.

Per configurare il rapporto tra attività e connessioni, impostare la spark.databricks.sql.streaming.jdbc.tasksPerConnection configurazione di Spark. Se il database di destinazione ha un limite di connessioni basso, ridurre il numero di partizioni di shuffle o aumentare spark.databricks.sql.streaming.jdbc.tasksPerConnection.

Il sink ritenta automaticamente gli errori JDBC temporanei, inclusi errori di connessione, deadlock e limitazione della frequenza. Se il sink esaurisce tutti i nuovi tentativi, la query ha esito negativo.

Trigger e modalità di output supportati

Triggers

Questa tabella mostra il supporto per i tipi di trigger di Structured Streaming sia nel calcolo classico che serverless:

Attivatore Calcolo classico Calcolo serverless (notebook e lavori)
RealTime No
ProcessingTime No
AvailableNow
Once Yes. Deprecated. Utilizzare il AvailableNow. Yes. Deprecated. Utilizzare il AvailableNow.

Modalità di output

Questa tabella mostra il supporto per le modalità di output di Structured Streaming:

Modalità output Supported
update
append Yes. Il comportamento è identico a update. La query esegue un upsert quando la tabella di destinazione ha una chiave primaria; altrimenti esegue un inserimento. Vedi comportamento di Upsert.
complete No

Limitazioni

  • Per un database PostgreSQL esterno collegato tramite una connessione Unity Catalog, la tabella target deve già esistere. Il sink crea automaticamente tabelle mancanti solo in Lakebase.
  • I condotti a flusso lacustre non sono supportati.