Usa Zerobus Ingest

Questa pagina descrive come ingerire dati utilizzando Zerobus Ingest in Lakeflow Connect.

Inizia con Zerobus Ingest

Prima di iniziare, verifica che Zerobus Ingest sia disponibile nella regione del tuo workspace. Consultare Disponibilità dell'inserimento.

  1. Ottenere un URL di inserimento Zerobus.
  2. Creare o identificare la tabella in cui inserire i dati.
  3. Creare un principal del servizio e concedere alla tabella privilegi.
  4. Connetti un client o un esportatore per avviare l'invio di dati.

Scegliere la guida per il caso d'uso:

  • Inserire dati personalizzati: usare gli SDK di inserimento Zerobus o l'API REST con uno schema definito. Seguire le istruzioni riportate in questa pagina.

  • Inserire dati OpenTelemetry: usare gli SDK OpenTelemetry standard o gli agenti di raccolta per inviare tracce, log e metriche in schemi di tabella predefiniti. Per istruzioni complete, vedere Ingestione dati OpenTelemetry con Zerobus Ingest.

Scegliere un'interfaccia

Zerobus Ingest supporta diverse interfacce, che scrivono tutte direttamente nelle tabelle Delta di Unity Catalog. In breve:

  • SDK su gRPC: massima velocità effettiva sostenuta, ideale per i produttori di flussi di dati ad alto volume.
  • REST: senza stato, ideale per grandi flotte di dispositivi edge leggeri o ad alta frequenza di comunicazione.
  • OpenTelemetry (OTLP): per i sistemi che già emettono tracce, log e metriche OpenTelemetry. Vedi Ingestire dati OpenTelemetry con Zerobus Ingest.

Per un confronto completo e come scegliere, vedi protocolli API. Oltre agli SDK, puoi anche scegliere un formato di record (JSON, Protocol Buffers (protobuf) o Apache Arrow). Vedi Tipi di messaggio. Il resto di questa pagina utilizza gli SDK e l'API REST.

Ottieni l'URL dell'area di lavoro e l'endpoint Zerobus di inserimento

L'URL dell'area di lavoro viene visualizzato nel browser quando si accede. Mentre l'URL completo segue il formato https://<databricks-instance>.net/o=XXXXX, l'URL dell'area di lavoro è costituito da tutti gli elementi prima di /o=XXXXX. Ad esempio, dato l'URL completo seguente, è possibile determinare l'URL dell'area di lavoro e l'ID dell'area di lavoro.

  • URL completo: https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864#
  • URL area di lavoro: https://abcd-teste2-test-spcse2.azuredatabricks.net
  • ID area di lavoro: 2281745829657864

L'endpoint del server dipende dallo spazio di lavoro e dalla regione.

  • Endpoint server: <workspace-id>.zerobus.<region>.azuredatabricks.net

Per trovare la tua area di lavoro, apri il commutatore di spazio di lavoro nella barra di navigazione superiore dell'interfaccia di Databricks. La regione è visualizzata sotto ogni nome di spazio di lavoro (ad esempio, eastus). È anche possibile trovarlo nella console dell'account in Aree di lavoro.

Per la disponibilità delle regioni, vedi le quote di Zerobus Ingest.

Creare o identificare la tabella di destinazione

Identificare la tabella di destinazione in cui inserire i dati. Per creare una nuova tabella di destinazione, eseguire il CREATE TABLE comando SQL. Ad esempio, creare una nuova tabella denominata unity.default.air_quality.

    CREATE TABLE unity.default.air_quality (
    device_name STRING, temp INT, humidity LONG);

Zerobus Ingest può scrivere sia su tabelle Delta gestite che su tabelle di streaming, che funzionano allo stesso modo, con gli stessi limiti e quote.

Annotazioni

Per l'inserimento di OpenTelemetry, le tabelle devono usare schemi predefiniti per ogni tipo di segnale (tracce, log, metriche). Vedere Creare tabelle di destinazione in Unity Catalog.

Il tuo schema di tabella è il contratto per ciò che Zerobus Ingest accetta, e Zerobus Ingest non lo evolve mai automaticamente. Pianifica i cambiamenti dello schema in modo proattivo: evolvere prima la tabella, poi aggiornare i produttori. Zerobus Ingest scrive i record che non sono più compatibili dopo una modifica incompatibile della tabella in una posizione di fallback persistente invece di eliminarli. Vedi la gestione dello schema e il recupero dei dati dalla posizione di fallback permanente.

Per impostazione predefinita, Zerobus Ingest rifiuta i record con campi che non corrispondono allo schema della tabella di destinazione. Per acquisire tali campi invece di perderli, configurare una colonna di salvataggio. Vedi colonna di soccorso Zerobus.

Creare un principale del servizio e assegnare autorizzazioni

Il principale del servizio è un'identità specializzata che offre maggiore sicurezza rispetto agli account personalizzati. Per maggiori informazioni sui principi di servizio e su come usarli per l'autenticazione, vedi Autorizzare l'accesso del principale servizio ad Azure Databricks con OAuth.

Puoi creare e gestire i principali di servizio in modo programmativo con l'API o gli SDK REST di Azure Databricks, oppure tramite l'interfaccia dello spazio di lavoro come descritto di seguito. I permessi concessi alla fine di questa sezione sono SQL che puoi eseguire da qualsiasi client.

  1. Per creare un principale di servizio, vai su Impostazioni>Identità e Accesso.

  2. In Entità Principali di Servizio, selezionare Gestisci.

  3. Fare clic su Aggiungi entità servizio.

  4. Nella finestra Aggiungi entità servizio creare una nuova entità servizio facendo clic su Aggiungi nuovo.

  5. Generare e salvare l'ID client e il client secret per il service principal.

  6. Concedere le autorizzazioni necessarie per il catalogo, lo schema e la tabella al principale del servizio.

    1. Nella pagina del principale del servizio , vai alla scheda Configurazioni .
    2. Copiare l'ID applicazione (UUID).
    3. Usare il codice SQL seguente per concedere le autorizzazioni, sostituendo l'UUID e il catalogo di esempio, il nome dello schema e i nomi di tabella, se necessario.
    GRANT USE CATALOG ON CATALOG <catalog> TO `<UUID>`;
    GRANT USE SCHEMA ON SCHEMA <catalog.schema> TO `<UUID>`;
    GRANT MODIFY, SELECT ON TABLE <catalog.schema.table_name> TO `<UUID>`;
    

Scrivere un client

Usare Zerobus SDK nel linguaggio di programmazione preferito o nell'API REST per inserire i dati nella tabella di destinazione. Gli SDK sono open source. Per la libreria completa, la documentazione specifica per linguaggio e altri esempi, vedi il repository Zerobus SDK.

Gli esempi qui sotto usano ingest_record_offset, che preserva l'ordine in cui invii i documenti.

PYTHON SDK

Python 3.9 o versione successiva è richiesto. L’SDK garantisce un throughput elevato e un I/O di rete efficiente tramite un runtime asincrono. Supporta JSON (più semplice) e buffer di protocollo (consigliato per la produzione). L'SDK supporta inoltre sia implementazioni di sync che asincrone, oltre ai metodi di ingestione basati su offset e basati sul futuro.

pip install databricks-zerobus-ingest-sdk

Esempio di JSON:

import logging
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties

# See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
SERVER_ENDPOINT="https://1234567890123456.zerobus.eastus.azuredatabricks.net"
DATABRICKS_WORKSPACE_URL="https://adb-1234567890123456.12.azuredatabricks.net"
TABLE_NAME="main.default.air_quality"
CLIENT_ID="your-client-id"
CLIENT_SECRET="your-client-secret"

sdk = ZerobusSdk(
    SERVER_ENDPOINT,
    DATABRICKS_WORKSPACE_URL
)

table_properties = TableProperties(TABLE_NAME)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

try:
    for i in range(1000):
        record_dict = {
            "device_name": f"sensor-{i}",
            "temp": 20 + i % 15,
            "humidity": 50 + i % 40
        }
        stream.ingest_record_offset(record_dict)
finally:
    stream.close()

Gli esempi sopra utilizzano il metodo ingest_record_offset basato sull'offset senza attendere l'offset restituito. Per conoscere i metodi di ingestione disponibili, quando attendere la conferma di durabilità su un offset e come monitorare l'avanzamento tramite un callback di conferma, consulta Blocco dei messaggi e conferma.

Buffer di protocollo: Per un'ingestione sicura del tipo, passa un descrittore protobuf a TableProperties (il formato viene selezionato automaticamente). Genera uno schema dalla tua tabella usando lo generate_proto strumento, compilalo con protoc, poi passa il descrittore compilato per creare il flusso.

Arrow Flight: Per l'ingestione colonnare o in batch di dati Apache Arrow RecordBatch tramite la stessa connessione gRPC, vedi Usare Arrow Flight con Zerobus Ingest. Richiede il componente aggiuntivo [arrow]: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.

Per la documentazione completa, le opzioni di configurazione, l'inserimento in batch e gli esempi di buffer del protocollo, vedere il repository Python SDK.

Rust SDK

È necessario Rust 1.70 o versione successiva. L'SDK utilizza I/O asincroni e gRPC per l'ingestione ad alta produttività. Supporta JSON (più semplice) e buffer di protocollo (consigliato per la produzione).

Prima di tutto, importare il pacchetto.

cargo add databricks-zerobus-ingest-sdk

In alternativa, aggiungerlo a Cargo.toml.

[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication

Esempio di JSON:

  use databricks_zerobus_ingest_sdk::{JsonString, ZerobusSdk};
  use std::error::Error;

  // See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
  const DATABRICKS_WORKSPACE_URL: &str = "https://adb-1234567890123456.12.azuredatabricks.net";
  const SERVER_ENDPOINT: &str = "1234567890123456.zerobus.eastus.azuredatabricks.net";
  const TABLE_NAME: &str = "main.default.air_quality";
  const CLIENT_ID: &str = "your-client-id";
  const CLIENT_SECRET: &str = "your-client-secret";


  #[tokio::main]
  async fn main() -> Result<(), Box<dyn Error>> {
      let sdk_handle = ZerobusSdk::builder()
          .endpoint(SERVER_ENDPOINT)
          .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
          .build()?;

      let mut stream = sdk_handle
          .stream_builder()
          .table(TABLE_NAME)
          .oauth(CLIENT_ID, CLIENT_SECRET)
          .json()
          .max_inflight_requests(100)
          .build()
          .await?;

      stream.ingest_record_offset(
        JsonString("{
          \"device_name\": \"sensor\",
          \"temp\": 22,
          \"humidity\": 55}".to_string())).await?;

      println!("Record ingested successfully");
      stream.close().await?;
      println!("Stream closed successfully");

      Ok(())
  }

Protocol Buffers: Per l'ingestione con controllo dei tipi, usare Protocol Buffers tramite .compiled_proto(descriptor) nel generatore del flusso invece di .json(), dove descriptor è un prost_types::DescriptorProto. Generare i file necessari usando lo generate_proto strumento e importarli nel progetto. Arrow Flight: Per l'ingestione colonnare o in batch di dati Apache Arrow RecordBatch tramite la stessa connessione gRPC, consulta Usa Arrow Flight con Zerobus Ingest. Abilitare con la funzionalità Cargo: cargo add databricks-zerobus-ingest-sdk --features arrow-flight.

Per la documentazione completa, le opzioni di configurazione, l'inserimento batch, generate_proto gli strumenti e gli esempi di Protocol Buffer, vedere il repository Rust SDK.

JAVA SDK

è necessario Java 8 o versione successiva. L’SDK garantisce una bassa latenza e operazioni di I/O di rete efficienti per l’ingestione ad alta velocità. Supporta JSON (più semplice) e buffer di protocollo (consigliato per la produzione).

Intenditore:

<dependency>
    <groupId>com.databricks</groupId>
    <artifactId>zerobus-ingest-sdk</artifactId>
    <version>0.2.0</version>
</dependency>

Esempio di JSON:

import com.databricks.zerobus.*;

public class ZerobusClient {

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
    private static final String SERVER_ENDPOINT =
        "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
    private static final String DATABRICKS_WORKSPACE_URL =
        "https://adb-1234567890123456.12.azuredatabricks.net";
    private static final String TABLE_NAME = "main.default.air_quality";
    private static final String CLIENT_ID = "your-client-id";
    private static final String CLIENT_SECRET = "your-client-secret";

    public static void main(String[] args) throws Exception {
        ZerobusSdk sdk = new ZerobusSdk(
            SERVER_ENDPOINT,
            DATABRICKS_WORKSPACE_URL
        );

        ZerobusJsonStream stream = sdk.streamBuilder()
            .table(TABLE_NAME)
            .oauth(CLIENT_ID, CLIENT_SECRET)
            .json()
            .build()
            .join();

        try {
            for (int i = 0; i < 100; i++) {
                String record = String.format(
                    "{\"device_name\": \"sensor-%d\", \"temp\": 22, \"humidity\": 55}", i
                );
                stream.ingestRecordOffset(record);
            }
        } finally {
            stream.close();
        }
    }
}

Buffer di protocollo: Per un'ingestione sicura per tipi, crea un ZerobusProtoStream con streamBuilder() e .compiledProto(...). Generare uno schema dalla tabella usando lo strumento JAR in bundle, quindi compilarlo con protoc.

Arrow Flight: Per l'ingestione colonnare o orientata ai batch dei dati Apache Arrow RecordBatch tramite la stessa connessione gRPC, consulta Usare Arrow Flight con Zerobus Ingest.

Per la documentazione completa, le opzioni di configurazione, l'inserimento batch e gli esempi di buffer del protocollo, vedere il repository Java SDK.

SDK di sviluppo Go

È necessario passare alla versione 1.21 o successiva. L'SDK offre un alto throughput e prestazioni per l'ingestione in streaming. Supporta JSON (più semplice) e buffer di protocollo (consigliato per la produzione).

go get github.com/databricks/zerobus-sdk/go@latest

Esempio di JSON:

Per semplicità, gli errori vengono ignorati qui. Nel codice di produzione controllare sempre gli errori.

package main

import (
	"fmt"

	zerobus "github.com/databricks/zerobus-sdk/go"
)

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const (
	ServerEndpoint         = "https://1234567890123456.zerobus.eastus.azuredatabricks.net"
	DatabricksWorkspaceURL = "https://adb-1234567890123456.12.azuredatabricks.net"
	TableName              = "main.default.air_quality"
	ClientID               = "your-client-id"
	ClientSecret           = "your-client-secret"
)

func main() {
	sdk, _ := zerobus.NewZerobusSdk(
		ServerEndpoint,
		DatabricksWorkspaceURL,
	)
	defer sdk.Free()

	options := zerobus.DefaultStreamConfigurationOptions()
	options.RecordType = zerobus.RecordTypeJson

	stream, _ := sdk.CreateStream(
		zerobus.TableProperties{
			TableName: TableName,
		},
		ClientID,
		ClientSecret,
		options,
	)
	defer stream.Close()

	_, _ = stream.IngestRecordOffset(`{
		"device_name": "sensor-001",
		"temp": 20,
		"humidity": 60
	}`)

  fmt.Println("Record ingested successfully")

  _ = stream.Close()
  fmt.Println("Stream closed successfully")
}

Protocol Buffers: Per l'inserimento sicuro rispetto ai tipi, usare Protocol Buffers con RecordTypeProto (impostazione predefinita) e specificare un oggetto descriptorProto nelle proprietà della tabella. Creare un file .proto che corrisponda al tuo schema di tabella ed eseguire lo script generate_proto per aiutarti a importare i file nel progetto.

Arrow Flight: Per l'ingestione in formato colonnare o orientata ai batch dei dati Apache Arrow RecordBatch sulla stessa connessione gRPC, consulta Usa Arrow Flight con Zerobus Ingest.

Per la documentazione completa, le opzioni di configurazione, l'inserimento in batch, gli strumenti di generate_proto e gli esempi di buffer del protocollo, vedere il repository Go SDK.

C++ SDK

Important

L'SDK C++ è in fase Beta.

È richiesta C++17 o superiore. L'SDK fornisce streaming nativo gRPC, OAuth e recupero automatico tramite un'interfaccia RAII C++. Supporta JSON per configurazioni semplici e Protocol Buffer per carichi di lavoro di produzione.

L’SDK viene distribuito come pacchetto precompilato specifico per ciascuna piattaforma, quindi non è necessaria una toolchain di Rust per utilizzarlo. Scarica il bundle per la tua piattaforma (macOS, Linux, musl compreso, o Windows) dalla pagina delle release, estrailo, quindi fa' puntare CMake all'archivio FFI incluso. L'archivio è denominato libzerobus_ffi.a su macOS e Linux e zerobus_ffi.lib su Windows:

# macOS and Linux
cmake -S cpp -B build \
  -DZEROBUS_FFI_LIBRARY="$PWD/lib/libzerobus_ffi.a" \
  -DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j

Su Windows (PowerShell), punta invece all'archivio.lib:

cmake -S cpp -B build `
  -DZEROBUS_FFI_LIBRARY="$PWD/lib/zerobus_ffi.lib" `
  -DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j

Per compilare l’SDK a partire da un checkout del codice sorgente all’interno del tuo progetto CMake, aggiungilo come sottodirectory e collega il target. Puoi anche usarlo FetchContent per recuperarlo al momento della configurazione. Questo costruisce l'FFI dalla fonte Rust, quindi richiede una toolchain Rust:

add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)

Per utilizzare invece un pacchetto precompilato tramite add_subdirectory, imposta prima i percorsi FFI in modo che CMake colleghi l'archivio incluso nel pacchetto invece di tentare di compilarlo dal codice sorgente Rust, che non è presente. Uso zerobus_ffi.lib su Windows:

set(ZEROBUS_FFI_LIBRARY "/path/to/bundle/lib/libzerobus_ffi.a")
set(ZEROBUS_FFI_HEADER_DIR "/path/to/bundle/lib")
add_subdirectory(path/to/zerobus-sdk/cpp zerobus-cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)

Esempio di JSON:

L'ingestione è asincrona e pipelineata. I ingest_* metodi mettono in coda un record e restituiscono immediatamente. Metti in coda il batch e chiama flush() una volta, invece di aspettare dopo ogni record.

#include "zerobus/zerobus.hpp"
#include <string>
#include <vector>

int main() {
  // See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
  const std::string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
  const std::string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
  const std::string TABLE_NAME = "main.default.air_quality";
  const std::string CLIENT_ID = "your-client-id";
  const std::string CLIENT_SECRET = "your-client-secret";

  zerobus::Sdk sdk = zerobus::Sdk::builder()
                         .endpoint(SERVER_ENDPOINT)
                         .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
                         .application_name("my-app")
                         .build();

  zerobus::TableProperties table;
  table.table_name = TABLE_NAME;   // empty descriptor => JSON stream

  zerobus::StreamOptions options;
  options.record_type = zerobus::RecordType::Json;

  zerobus::Stream stream =
      sdk.create_stream(table, CLIENT_ID, CLIENT_SECRET, options);

  std::vector<std::string> batch = {
      R"({"device_name": "sensor-001", "temp": 20, "humidity": 60})",
      R"({"device_name": "sensor-002", "temp": 22, "humidity": 55})",
  };
  stream.ingest_json_records(batch);   // queue the batch — no per-record wait
  stream.flush();                      // wait once for all acks
  stream.close();

  return 0;
}

Ogni errore genera zerobus::ZerobusException, che include un messaggio e un flag is_retryable(). Per monitorare la persistenza su un flusso continuo senza bloccare, registra un AckCallback tramite StreamOptions::ack_callback. I callback vengono eseguiti in modo seriale in un thread in background e devono essere noexcept. Consulta la documentazione dell'SDK C++ per il contratto completo relativo a threading, politica di drain e ciclo di vita.

Per un'ingestione type-safe, puoi utilizzare i Protocol Buffer in uno dei due modi:

  • Genera lo schema dal Catalogo Unity con ProtoSchema::from_uc_json(). Questo costruisce un descrittore e un encoder JSON-to-proto direttamente dai metadati della tabella, quindi non ha bisogno di .proto file o protoc:

    1. Recupera i metadati JSON della tabella dall'API Get a table (GET /api/2.1/unity-catalog/tables/{full_name}). Il service principal necessita di SELECT sul tavolo.
    2. Passa i metadati a ProtoSchema::from_uc_json() per creare il descrittore e l'encoder.
    3. Imposta TableProperties::descriptor_proto, poi ingeri con ingest_proto_records().
  • Compila un elemento sottoposto a check-in .proto con protoc per ottenere la tipizzazione in fase di compilazione.

Per una guida eseguibile, vedi gli esempi dei Protocol Buffers.

Per l'ingestione colonnare o orientata ai batch dei batch di record di Apache Arrow tramite la stessa connessione gRPC, consulta Usare Arrow Flight con Zerobus Ingest.

Per documentazione completa, opzioni di configurazione, ingestione batch e esempi di Protocol Buffer, consulta il repository SDK C++.

SDK di C#

Important

L'SDK C# / .NET è in beta. Il pacchetto Databricks.Zerobus è in versione prerilascio.

.NET 8.0 o superiore è richiesto. L'SDK fornisce streaming nativo gRPC, OAuth e recupero automatico. Supporta JSON per configurazioni semplici e Protocol Buffer per carichi di lavoro di produzione. Arrow Flight non è disponibile nell'SDK C#.

Aggiungere il pacchetto Databricks.Zerobus al progetto:

dotnet add package Databricks.Zerobus

Esempio di JSON:

using Databricks.Zerobus;

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
const string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
const string TABLE_NAME = "main.default.air_quality";
const string CLIENT_ID = "your-client-id";
const string CLIENT_SECRET = "your-client-secret";

using var sdk = ZerobusSdk.CreateBuilder()
    .Endpoint(SERVER_ENDPOINT)
    .UnityCatalogUrl(DATABRICKS_WORKSPACE_URL)
    .Build();

using var stream = sdk.CreateJsonStream(TABLE_NAME, CLIENT_ID, CLIENT_SECRET);

long offset = stream.IngestRecord(
    """{"device_name": "sensor-1", "temp": 22, "humidity": 55}""");
stream.WaitForOffset(offset);
stream.Close();

IngestRecord restituisce l'offset del record e WaitForOffset blocca finché quel record non sia stato reso persistente. Per acquisire un batch, usa IngestRecords, che accetta un array di record e restituisce l'ultimo offset. Il blocco sull'offset è opzionale. Vedi Blocco e conferma dei messaggi.

Buffer di protocollo: Per l'ingestione type-safe, crea uno stream con sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), dove descriptorProto sono i byte serializzati DescriptorProto del tuo messaggio compilato, poi ingeri con stream.IngestRecord(protoBytes).

Per documentazione completa, opzioni di configurazione e esempi di Protocol Buffer, consulta il repository SDK C#.

TypeScript SDK

è necessario Node.js 16 o versione successiva. L'SDK offre alte prestazioni con supporto asincrono tramite JavaScript Promises. Supporta JSON (più semplice) e buffer di protocollo (consigliato per la produzione).

npm install @databricks/zerobus-ingest-sdk

Esempio di JSON:

import { ZerobusSdk, RecordType } from '@databricks/zerobus-ingest-sdk';

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const SERVER_ENDPOINT = 'https://1234567890123456.zerobus.eastus.azuredatabricks.net';
const DATABRICKS_WORKSPACE_URL = 'https://adb-1234567890123456.12.azuredatabricks.net';
const TABLE_NAME = 'main.default.air_quality';
const CLIENT_ID = 'your-client-id';
const CLIENT_SECRET = 'your-client-secret';

const sdk = new ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL);

const stream = await sdk.createStream({ tableName: TABLE_NAME }, CLIENT_ID, CLIENT_SECRET, {
  recordType: RecordType.Json,
});

try {
  for (let i = 0; i < 100; i++) {
    const record = { device_name: `sensor-${i}`, temp: 22, humidity: 55 };
    await stream.ingestRecordOffset(record);
  }
} finally {
  await stream.close();
}

Protocol Buffers: Per l'inserimento sicuro rispetto ai tipi, usare Protocol Buffers con RecordType.Proto (impostazione predefinita) e specificare un oggetto descriptorProto nelle proprietà della tabella.

Arrow Flight: Per l'ingestione di dati Apache Arrow RecordBatch orientata a colonne o in batch sulla stessa connessione gRPC, consulta Usare Arrow Flight con Zerobus Ingest.

Per la documentazione completa, le opzioni di configurazione, l'inserimento in batch e gli esempi di buffer del protocollo, vedere il repository TypeScript SDK.

REST API

L'API REST consente di inserire un singolo record inviando una richiesta HTTP POST all'endpoint /zerobus/v1/tables/<table-name>/insert . Il record stesso è incluso nel corpo della richiesta e deve essere in formato JSON.

Questo esempio illustra come usare CURL per eseguire il push dei dati in Zerobus Ingest usando l'API REST.

Intestazioni

La richiesta richiede due intestazioni HTTP specifiche per autenticare e formattare correttamente la richiesta.

  • Content-Type: application/json
    • Campo obbligatorio per specificare il tipo di contenuto. Json è attualmente l'unico formato di messaggio supportato.
  • Autorizzazione: Bearer <token>
    • Sostituire <il token> con il token OAuth recuperato usando il comando curl fornito più avanti.

Recuperare il token OAuth: Questi token scadono ogni ora e devono essere aggiornati. È possibile aggiornarli recuperando nuovamente il token OAuth.

Compilare i parametri seguenti:

  • $CATALOG, $SCHEMA, $TABLE, , $WORKSPACE_ID, $WORKSPACE_URL
  • $DATABRICKS_CLIENT_ID e $DATABRICKS_CLIENT_SECRET
    • Questi due parametri corrispondono al principio di servizio che hai creato.
authorization_details=$(cat <<EOF
[{
  "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": "$CATALOG.$SCHEMA"
},
{
  "type": "unity_catalog_privileges",
  "privileges": ["SELECT", "MODIFY"],
  "object_type": "TABLE",
  "object_full_path": "$CATALOG.$SCHEMA.$TABLE"
}]
EOF
)

export OAUTH_TOKEN=$(curl -X POST \
  -u "$DATABRICKS_CLIENT_ID:$DATABRICKS_CLIENT_SECRET" \
  -d "grant_type=client_credentials" \
  -d "scope=all-apis" \
  -d "resource=api://databricks/workspaces/$WORKSPACE_ID/zerobusDirectWriteApi" \
  --data-urlencode "authorization_details=$authorization_details" \
  "$WORKSPACE_URL/oidc/v1/token" | jq -r '.access_token')

Ingestione dei record:

Compilare i parametri seguenti:

Il corpo della richiesta deve essere un elenco di oggetti JSON.

curl -X POST \
  "$ZEROBUS_ENDPOINT/zerobus/v1/tables/$CATALOG.$SCHEMA.$TABLE/insert" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer $OAUTH_TOKEN" \
  -d '[{ "device_name": "device_num_1", "temp": 28, "humidity": 60 },
       { "device_name": "device_num_1", "temp": 28, "humidity": 60 }]'

Se tutte le informazioni vengono compilate correttamente, si dovrebbe ricevere una risposta JSON vuota con un codice di stato HTTP 200.

Gestire gli errori

Gli esempi sopra mostrano il sentiero felice. In produzione, gestisci l'ingestion con una logica di gestione degli errori. L’SDK ritenta automaticamente le operazioni in caso di errori transitori, come problemi di rete, grazie al meccanismo di ripristino integrato. I fallimenti da cui non può recuperare, come credenziali non valide o una tabella mancante, emergono come:ZerobusException

from zerobus.sdk.shared import ZerobusException

try:
    stream.ingest_record_offset(record)
except ZerobusException as e:
    # Handle the failure: log it, fix the cause, recover on a new stream, or stop.
    ...

Gli SDK recuperano automaticamente anche da guasti transitori e ti permettono di recuperare i record non riconosciuti quando un flusso fallisce permanentemente. Per i pattern di client resilienti e il riferimento completo degli errori, vedere Pattern di ripristino e nuovi tentativi e Gestione degli errori di Zerobus Ingest.

Passaggi successivi