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.
Importante
Questa funzionalità è in versione beta. Gli amministratori dell'area di lavoro possono controllare l'accesso a questa funzionalità dalla pagina Anteprime . Vedere Gestire le anteprime di Azure Databricks.
Zerobus Ingest offre API di produttore compatibili con Kafka che ti permettono di ingerire utilizzando qualsiasi client produttore Apache Kafka senza un SDK Azure Databricks. Punti un produttore Kafka esistente all'endpoint Zerobus e produci su un argomento chiamato in base alla tua tabella target, e i record finiscono in una tabella Delta del Unity Catalog. Le API compatibili con Kafka sono adatte quando hai già un produttore Kafka, un collector che parla il protocollo Kafka, o uno strumento che emette a Kafka, e vuoi instradare quei dati in Delta con modifiche minime al codice.
Le API compatibili con Kafka girano su TCP con SASL_SSL e il OAUTHBEARER meccanismo e implementano il sottoinsieme lato produttore del protocollo Kafka, che copre Produce, Metadata, ApiVersions, e le API di handshake SASL. Le API per i consumatori, di amministrazione e transazionali non sono disponibili. Le API sono di sola scrittura.
Quando utilizzare le API compatibili con Kafka
Le API compatibili con Kafka sono le più adatte nei seguenti scenari:
- Vuoi inviare dati a Delta senza utilizzare un SDK Zerobus e utilizzi già un producer Kafka o un'applicazione, un agente o un collettore che invia dati a Kafka.
- Vuoi riutilizzare la configurazione Kafka di produttore, il batching e gli strumenti operativi esistenti.
- Invii record JSON e non hai bisogno di Protocol Buffer (protobuf) o Apache Arrow.
Se stai creando un nuovo client da zero e vuoi il massimo throughput, conferme per record e il recupero automatico, utilizza un SDK Zerobus basato su gRPC invece delle API compatibili con Kafka. Vedere Scegliere un'interfaccia. Per carichi di lavoro di tipo colonnare o in batch, consulta Usare Arrow Flight con Zerobus Ingest.
Come funziona il modello di ingestione
Le API compatibili con Kafka mappano i concetti Kafka su Zerobus Ingest come segue:
- Topics. Un nome di argomento Kafka è il nome completo della tabella a tre livelli del Catalogo Unity (
catalog.schema.table). La tabella target deve già esistere, perché Zerobus non crea mai argomenti. - Registri Zerobus assorbe solo il valore del record, che deve essere un oggetto JSON codificato in UTF-8 che corrisponde allo schema della tabella Delta di destinazione. Ignora le chiavi di record, le intestazioni, la partizione e il timestamp forniti dal client, e non li fa persistere.
- Autenticazione. Ogni connessione si autentica con un token OAuth di Azure Databricks limitato alla tabella di destinazione, presentato tramite SASL/
OAUTHBEARER. Vedere Autenticazione. - Ringraziamenti. Zerobus invia una
Producerisposta solo dopo aver memorizzato i record in modo durevole. Configura il tuo produttore conacks=all.
Zerobus è progettato senza partizioni. L'endpoint Zerobus è un singolo broker logico e un'unica partizione, quindi ogni record di un topic viene confermato sulla partizione 0. I produttori non devono prevedere questo. Zerobus Ingest scala orizzontalmente per gestire il carico in ingresso.
Inoltre, Zerobus Ingest fornisce almeno una consegna una volta. Le API compatibili con Kafka sono in Beta e hanno una propria quota, e condividono le caratteristiche di latenza, dimensione del record e tabelle partizionate del resto di Zerobus Ingest. Per la quota Beta e tutte le altre quote, vedi quote di Zerobus Ingest.
Authentication
Le API compatibili con Kafka utilizzano SASL/OAUTHBEARER. Il token portatore è un token Azure Databricks OAuth che ottieni con le credenziali client di un principale di servizio con accesso alla tabella target. Il token ha come ambito quella tabella tramite OAuth authorization_details e usa la risorsa zerobusDirectWriteApi, seguendo lo stesso flusso dell'API REST di Zerobus.
I token OAuth scadono dopo un'ora, quindi fornisci il token tramite il callback del provider di token del tuo client Kafka invece che come stringa statica. Il client poi recupera un token nuovo ogni volta che si riconnette. Le connessioni hanno anche una durata limitata lato server. Quando una connessione raggiunge quel limite, Zerobus la chiude e il produttore si riconnette e si riautentica autonomamente. La riautenticazione su una connessione live non è supportata.
Assegna al service principal i privilegi di Unity Catalog richiesti sulla tabella di destinazione prima di connetterti. Vedi Crea un principio di servizio e concede permessi.
Scrivere un client
L'esempio seguente produce la stessa air_quality tabella utilizzata negli esempi di Use Zerobus Ingest. Utilizza kafka-python, ma funziona qualsiasi client producer Kafka che supporti SASL_SSL con il meccanismo OAUTHBEARER. Adatta il pattern del token provider alla tua libreria client.
Il produttore si collega al server bootstrap Zerobus sulla porta 9092. Trova l'ID e la regione del tuo workspace come descritto in Ottieni l'URL del tuo workspace e l'endpoint Zerobus Ingest.
-
Server Bootstrap:
<workspace-id>.zerobus.<region>.azuredatabricks.net:9092
pip install kafka-python requests
Passo 1: Crea un provider di token
Zerobus autentica ciascuna connessione con un token OAuth di Azure Databricks di breve durata, con ambito limitato alla tabella. Poiché il token scade, fornisci un callback che genera un nuovo token su richiesta anziché un token statico.
La funzione fetch_zerobus_token() scambia le credenziali della tua entità servizio con un token con ambito limitato alla tabella di destinazione, e ZerobusTokenProvider lo incapsula nell'interfaccia di callback prevista da kafka-python.
import json
import requests
from kafka.sasl.oauth import AbstractTokenProvider
# See "Get your workspace URL and Zerobus Ingest endpoint" in zerobus-ingest.md.
WORKSPACE_ID = "1234567890123456"
WORKSPACE_URL = "https://dbc-a1b2c3d4-e5f6.cloud.databricks.com"
TABLE_NAME = "main.default.air_quality"
CLIENT_ID = "your-client-id"
CLIENT_SECRET = "your-client-secret"
def fetch_zerobus_token():
catalog, schema, table = TABLE_NAME.split(".")
authorization_details = [
{
"type": "unity_catalog_privileges",
"privileges": ["USE CATALOG"],
"object_type": "CATALOG",
"object_full_path": catalog,
},
{
"type": "unity_catalog_privileges",
"privileges": ["USE SCHEMA"],
"object_type": "SCHEMA",
"object_full_path": f"{catalog}.{schema}",
},
{
"type": "unity_catalog_privileges",
"privileges": ["SELECT", "MODIFY"],
"object_type": "TABLE",
"object_full_path": TABLE_NAME,
},
]
response = requests.post(
f"{WORKSPACE_URL}/oidc/v1/token",
auth=(CLIENT_ID, CLIENT_SECRET),
data={
"grant_type": "client_credentials",
"scope": "all-apis",
"resource": f"api://databricks/workspaces/{WORKSPACE_ID}/zerobusDirectWriteApi",
"authorization_details": json.dumps(authorization_details),
},
timeout=30,
)
response.raise_for_status()
return response.json()["access_token"]
# kafka-python calls token() whenever it needs a fresh OAuth token.
class ZerobusTokenProvider(AbstractTokenProvider):
def token(self):
return fetch_zerobus_token()
Passo 2: Configura il produttore e invia i record
Indirizza il producer al server di bootstrap, configura SASL_SSL utilizzando il meccanismo OAUTHBEARER e passa il provider di token dal passaggio 1. Utilizzare acks="all" in modo che ogni lotto venga riconosciuto solo dopo essere stato persistito duraturamente, e inviare i record non compressi.
from kafka import KafkaProducer
BOOTSTRAP_SERVERS = "1234567890123456.zerobus.us-west-2.cloud.databricks.com:9092"
producer = KafkaProducer(
bootstrap_servers=BOOTSTRAP_SERVERS,
security_protocol="SASL_SSL",
sasl_mechanism="OAUTHBEARER",
sasl_oauth_token_provider=ZerobusTokenProvider(),
# Wait for durable acknowledgment before treating a record as ingested.
acks="all",
# Compression is not supported by the endpoint; send records uncompressed.
compression_type=None,
)
# Each send() returns a future immediately. The topic is the full table name.
futures = [
producer.send(
topic=TABLE_NAME,
value=json.dumps(
{"device_name": f"sensor-{i}", "temp": 20 + i % 15, "humidity": 50 + i % 40}
).encode("utf-8"),
)
for i in range(1000)
]
producer.flush()
# Block on each future to confirm every record was durably acknowledged.
for future in futures:
future.get(timeout=30)
producer.close()
print("All records ingested successfully")
Opzioni di configurazione
Le API compatibili con Kafka implementano il sottoinsieme produttore del protocollo Kafka. Configura il tuo produttore secondo le seguenti opzioni.
| Option | dettagli |
|---|---|
| Formato del record | Solo JSON. Ogni valore di record deve essere un oggetto JSON codificato in UTF-8 che corrisponda allo schema della tabella target. Per inviare protobuf o Avro, usa un SDK Zerobus. |
| Compressione | Non supportato. Invia batch non compressi, ad esempio compression.type=none. Zerobus rifiuta i batch gzip, snappy, lz4 e zstd con il codice di errore UNSUPPORTED_COMPRESSION_TYPE. |
| Campi del record | Solo il valore. Zerobus assorbe il valore del record e non persiste chiavi, intestazioni, assegnazioni di partizioni o timestamp. |
| Supporto dell'API | Solo scrittura. Zerobus accetta Produce richieste, oltre ai metadati e alle richieste SASL necessarie a stabilire una sessione. API consumer, admin e transazionali non sono supportate. |
| Collegamento privato front-end | Non supportato. Connettiti invece tramite l'endpoint pubblico. |
| Schema | Applicato. Zerobus rifiuta i record con campi che non corrispondono allo schema della tabella di destinazione e considera le colonne Delta aggiuntive che accettano valori null come una modifica non incompatibile. Per catturare campi non abbinati invece di rifiutarli, configura una colonna di salvataggio. |
Per i limiti relativi a latenza, quote, dimensione del record e tabelle partizionate, vedere Quote di acquisizione di Zerobus.
Procedure consigliate
Segui queste linee guida per ottenere la migliore prestazione e affidabilità dalle API compatibili con Kafka.
- Riutilizzare un producer a lunga durata per molti record invece di crearne uno per ogni batch, perché la creazione del producer e l'handshake SASL comportano un costo di inizializzazione.
- Lascia che il producer accumuli i record in batch, ad esempio configurando
linger.msebatch.size, invece di svuotare il buffer dopo ogni record. L'elaborazione in batch è la principale leva per il throughput. - Usare
acks=allper ottenere una conferma permanente per ogni batch, in linea con la semantica at-least-once di Zerobus. - Recupera i token OAuth tramite il callback del provider del tuo client in modo che si aggiornino automaticamente alla riconnessione, invece di passare un token statico che scade.
- Esegui il produttore nella stessa regione cloud dell'endpoint Zerobus per ottenere la massima produttività.
Gestione degli errori
Zerobus segnala i fallimenti usando codici di errore Kafka standard sull'argomento e sulla partizione interessati. I codici comuni includono:
| Errore di Kafka | Meaning |
|---|---|
SASL_AUTHENTICATION_FAILED |
Il token OAuth manca, non è valido o non dispone dei privilegi Unity Catalog richiesti sulla tabella. |
UNKNOWN_TOPIC_OR_PARTITION |
La tabella target non esiste, è stata eliminata, oppure il token non è autorizzato a scriverci. |
INVALID_RECORD |
Un record non ha validato lo schema o non poteva essere decodificato come UTF-8 JSON. |
MESSAGE_TOO_LARGE |
Un singolo disco supera il limite di dimensione di 10 MB. Vedi Dimensione del record. |
UNSUPPORTED_COMPRESSION_TYPE |
Il lotto è stato compresso. Invia i record non compressi. |
Dopo una richiesta fallita Produce , Zerobus restituisce i codici di errore per partizione e chiude la connessione. I producer di Kafka si riconnettono automaticamente, ma progetta il client in modo che segnali gli errori di invio, ad esempio controllando il risultato di ogni invio, così che i record non vengano scartati in modo silenzioso.
Risorse aggiuntive
- Usa Zerobus Ingest: Se non hai ancora configurato Zerobus Ingest, inizia da qui per le istruzioni su come trovare l'URL del tuo workspace, creare la tabella Delta target e configurare un service principal. Questi passaggi sono condivisi su tutte le interfacce.
- Quote di acquisizione di Zerobus: Verificare le quote predefinite di Zerobus prima di distribuire nell'ambiente di produzione.
- Utilizza Arrow Flight con Zerobus Ingest: per l'ingestione colonnare o in batch tramite gRPC.