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.
Gli SDK di Zerobus Ingest offrono diversi metodi per acquisire un record, che comportano un compromesso tra il throughput e il livello di conferma della durabilità restituito. Questa pagina spiega ogni metodo e quando bloccare la durabilità. Per reagire agli acknowledgment in modo asincrono invece che bloccare, vedi Acknowledgment callbacks.
Gli esempi in questa pagina utilizzano l'SDK Python. Per le opzioni esatte di timeout e configurazione accettate da ogni metodo (inclusi i loro valori predefiniti e unità), consulta il repository Zerobus SDK. Gli altri SDK linguistici espongono opzioni equivalenti.
Cos'è un offset?
A ogni record che acquisisci viene assegnato un offset: la sua posizione nello stream. Il contrasto è come ti riferisci a un documento specifico quando vuoi confermare che sia stato scritto in modo duraturo. Zerobus Ingest offre garanzie di consegna almeno una volta, e attendere un offset è il modo con cui un client conferma tale garanzia per un determinato record.
Confermare un offset significa che il record è durabile, non che sia ancora consultabile nella tabella Delta. Zerobus Ingest scrive i record persistenti nella tabella poco dopo. Per i dati di latenza, vedi Latenza.
Metodi di inserimento
Gli SDK offrono due modi per ingerire un record. (I nomi dei metodi qui sotto provengono dall'SDK Python. Altri SDK espongono metodi equivalenti.)
| metodo | Restituzioni | Usarlo quando |
|---|---|---|
Basato sugli offset, ingest_record_offset() |
Il record offset, dopo che il record è stato messo in coda nello streaming. | Impostazione predefinita consigliata. Vuoi accodare i record in ordine e, facoltativamente, confermare la durabilità in un secondo momento attendendo un offset. |
Basata sul futuro, ingest_record() |
Un RecordAcknowledgment su cui puoi contare. |
Deprecated. Preferire l'approccio basato sull'offset per prestazioni migliori. |
Basato su offset (consigliato)
ingest_record_offset() invia il record e restituisce il suo offset una volta che il record è in coda sul stream. La chiamata viene eseguita sul tuo thread di chiamata, quindi i record vengono messi in coda nell'ordine in cui chiami il metodo, e l'offset restituito ti permette di confermare la durabilità successivamente con wait_for_offset(). Questo è il metodo predefinito raccomandato per la maggior parte dei produttori, ed è il metodo usato negli esempi di Usa Zerobus Ingest .
Orientato al futuro (obsoleto)
ingest_record() restituisce un oggetto RecordAcknowledgment su cui puoi attendere per garantirne la durabilità. Viene deprecato a favore del metodo basato su offset, che funziona meglio. Usalo solo per codice esistente che non è ancora stato migrato.
Registrazione per registrazione vs. ingestione batch
Ogni metodo di ingestione ha una variante batch (ad esempio, ingest_records_offset()) che invia una lista di record in una singola chiamata. Il batching è più efficiente rispetto alle singole chiamate per l'ingestione in massa.
Per JSON e Protocol Buffer (protobuf), un batch effettua un commit atomico: o ogni record del batch viene accettato e reso durabile, oppure l'intero batch viene rifiutato. Zerobus Ingest non esegue caricamenti parziali né confermi parziali per questi formati, quindi la tua tabella non contiene mai un batch parziale. Un batch che fallisce la validazione (ad esempio, un disadattamento dello schema) fallisce rapidamente, prima di toccare la tabella, invece di ottenere alcuni record e perderne altri.
Poiché un batch JSON o protobuf viene inviato come un singolo messaggio, la dimensione massima di 10 MB del messaggio si applica sia a un singolo record che a un intero lotto: tutti i record in un batch insieme devono stare entro 10 MB. Dimensiona i tuoi lotti per restare sotto quel limite. Vedi Dimensione del record.
I lotti di Arrow Flight sono l'eccezione
L'ingestione di Apache Arrow Flight non segue il modello tutto o niente, a messaggio singolo sopra descritto. Un batch Arrow può essere molto più grande di un batch JSON o protobuf, e il percorso Arrow Flight divide un grande lotto in messaggi di trasporto più piccoli che vengono inviati e riconosciuti individualmente invece che come un'unica unità atomica. Di conseguenza:
- Il limite di 10 MB per messaggio che si applica ai lotti JSON e protobuf non si applica allo stesso modo a un batch Arrow. Un batch Arrow di grandi dimensioni viene suddiviso in messaggi di trasporto invece di essere rifiutato perché troppo grande.
- La durabilità è confermata alla granularità del messaggio di trasporto, quindi un lotto logico molto grande può essere parzialmente resistente se si verifica un guasto a metà strada, invece di impegnare tutto o niente.
ingest_batch() restituisce comunque un singolo offset logico per il batch che hai inviato, e wait_for_offset() su quel offset si completa solo dopo che ogni messaggio di trasporto che compone il batch è stato confermato. Per il modello completo di Arrow Flight, la guida in batching e il recupero di dati non riconosciuti, vedi Usa Arrow Flight con Zerobus Ingest.
Quando dovresti bloccare un messaggio?
Attendere un offset sacrifica il throughput in cambio di una garanzia di durabilità più forte per ogni record nel codice cliente. Scegli in base al carico di lavoro:
- Non bloccare: il predefinito giusto per lo streaming ad alto volume, dove tieni a un throughput sostenuto e puoi confermare la durabilità in totale (ad esempio, alla chiusura dello stream o tramite un acknowledgment callback). La maggior parte dei produttori dovrebbe iniziare da qui.
-
Blocco su un offset: prendilo in considerazione quando la tua applicazione deve avere la certezza che un record specifico sia stato reso permanente prima di eseguire un'altra azione. Per esempio:
- Stai per eliminare o confermare la fonte dei dati (un messaggio in coda, un file, un cursore a monte) e non devi perderla se l'ingestione fallisce.
- Stai acquisendo dati in corrispondenza di checkpoint o limiti transazionali e hai bisogno che ogni checkpoint sia persistente prima di procedere.
- Stai facendo scritture di basso volume e alto valore dove la conferma per record conta più del throughput.
Non bloccare ogni record in un loop ad alta produttività. Questo serializza il produttore in un viaggio di andata e ritorno al server per ogni record e riduce drasticamente la velocità di lavoro. Invece, Azure Databricks consiglia di acquisire un ampio blocco di record e poi confermare la persistenza una sola volta per l'intero blocco. Hai due modi per farlo: attendere l'offset più recente, o svuotare il flusso. Il blocco per singolo record dovrebbe essere riservato ai casi specifici sopra in cui un singolo record deve essere confermato prima dell'azione successiva.
Aspetta un offset
wait_for_offset() rimane bloccato finché Zerobus Ingest non conferma che il record a tale offset sia stato scritto in modo permanente, oppure finché non scade il timeout. Usalo per confermare un punto specifico del flusso, in genere l'ultimo record di un segmento. Acquisisci il chunk, mantieni l'offset finale restituito dal loop e attendi solo quell'offset invece di aspettare dopo ogni record:
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties
sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)
table_properties = TableProperties("main.default.air_quality")
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)
try:
last_offset = 0
for row in records:
last_offset = stream.ingest_record_offset(row)
# Block until everything up to the last record of the chunk is durable
stream.wait_for_offset(last_offset)
print("Chunk durably written.")
finally:
stream.close()
Svuotare il flusso
flush() si blocca finché tutti i record acquisiti fino a quel momento non vengono scritti in modo permanente, quindi restituisce il controllo. A differenza di wait_for_offset(), non si tiene traccia di un offset: l'operazione di flush attende tutte le operazioni in sospeso nello stream. Non chiude il getto, quindi puoi continuare a ingerire dopo.
try:
for row in records:
stream.ingest_record_offset(row)
# Block until every pending record is durable
stream.flush()
print("All ingested records durably written.")
finally:
stream.close()
wait_for_offset vs. colore
Entrambi confermano la durata per un po' di tempo. Scegli in base a ciò che stai confermando:
- Usa
wait_for_offset(offset)quando vuoi confermare fino a un determinato record, ad esempio il limite di un checkpoint, mentre altri record potrebbero essere ancora in transito dopo di esso. - Usa
flush()quando vuoi confermare che tutti i record in sospeso siano persistenti prima di procedere, ad esempio alla fine di un batch, prima di far avanzare un cursore upstream o prima di arrestare il sistema.flush()è governato da un timeout di flush configurabile.
close() svuota e chiude lo stream, quindi i record vengono sempre resi persistenti in caso di arresto controllato. Chiamalo sempre in un blocco finally.
Reagire alle conferme di ricezione in modo asincrono
Se invece di bloccare vuoi reagire alle conferme di durabilità e agli errori man mano che arrivano, mentre il tuo produttore continua a spingere a tutta velocità, registra un acknowledgement callback sullo stream. I callback sono una funzione separata dalle chiamate di blocco presenti in questa pagina. Vedi richiami ai riconoscimenti.
Related
- Richiami di conferma: Reagire in modo asincrono agli acknowledgment e agli errori.
- Usa Zerobus Ingest: Scrivi un client.
- Flussi: Flussi e offset.
- Gestione degli errori di Zerobus Ingest: Gestione degli errori.