Nota
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare ad accedere o modificare le directory.
L'accesso a questa pagina richiede l'autorizzazione. È possibile provare a modificare le directory.
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 esempioidouser_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 formatoprojects/<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 formatoschema.table. Se si omette lo schema, il sink usa lopublicschema . 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 esempioidouser_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 formatoschema.table. Se si omette lo schema, il sink usa lopublicschema . -
<primary-key-columns>: facoltativo. Un elenco separato da virgole di tutte le colonne nella chiave primaria della tabella di destinazione, ad esempioidouser_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-fullquando il certificato Trust Server non è selezionato. Se non fornisci un certificato, la connessione viene utilizzatasslmode=verify-fullcon 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
batchsizerighe, 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
batchintervalper 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
batchsizeper 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 |
Sì | No |
ProcessingTime |
Sì | No |
AvailableNow |
Sì | Sì |
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 |
Sì |
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.