Usa Zerobus Ingest

Esta página describe cómo ingerir datos usando Zerobus Ingest en Lakeflow Connect.

Empieza con Zerobus Ingest

Antes de empezar, confirma que Zerobus Ingest está disponible en la región de tu espacio de trabajo. Consulte la disponibilidad de ingestión.

  1. Obtenga una URL de ingestión de Zerobus.
  2. Cree o identifique la tabla en la que desea ingerir datos.
  3. Cree una entidad de servicio y conceda privilegios a la tabla.
  4. Conecte un cliente o exportador para empezar a enviar datos.

Elija la guía para su caso de uso:

  • Ingesta de datos propios: utilice los SDK de ingesta de Zerobus o la API REST con un esquema que defina. Siga las instrucciones de esta página.

  • Ingesta de datos de OpenTelemetry: Use SDK o recopiladores estándar de OpenTelemetry para enviar seguimientos, registros y métricas a esquemas de tabla predefinidos. Para obtener instrucciones completas, consulte Ingesta de datos de OpenTelemetry con ingesta de Zerobus.

Elegir una interfaz

Zerobus Ingest admite varias interfaces, y todas escriben directamente en tablas Delta de Unity Catalog. En resumen:

  • SDKs con gRPC: mayor rendimiento sostenido, ideal para productores de streaming de gran volumen.
  • REST: sin estado, ideal para grandes flotas de dispositivos ligeros o "habladores" en el borde.
  • OpenTelemetry (OTLP): para sistemas que ya emiten trazas, registros y métricas de OpenTelemetry. Consulte Ingerir datos de OpenTelemetry con Zerobus Ingest.

Para una comparación completa y cómo elegir, consulta protocolos API. Sobre los SDK, también puedes elegir un formato de registro (JSON, Protocol Buffers (protobuf) o Apache Arrow). Ver Tipos de mensajes. El resto de esta página utiliza los SDKs y la API REST.

Obtenga la dirección URL del área de trabajo y el punto de conexión de Ingesta de Zerobus

La dirección URL del área de trabajo aparece en el explorador al iniciar sesión. Aunque la dirección URL completa sigue el formato https://<databricks-instance>.net/o=XXXXX, la dirección URL del área de trabajo consta de todo antes de /o=XXXXX. Por ejemplo, dada la siguiente dirección URL completa, puede determinar la dirección URL del área de trabajo y el identificador del área de trabajo.

  • Dirección URL completa: https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864#
  • Dirección URL del área de trabajo: https://abcd-teste2-test-spcse2.azuredatabricks.net
  • Id. del área de trabajo: 2281745829657864

El punto de conexión del servidor depende del área de trabajo y la región:

  • Punto de conexión del servidor: <workspace-id>.zerobus.<region>.azuredatabricks.net

Para encontrar la región de tu espacio de trabajo, abre el conmutador de espacios de trabajo en la barra de navegación superior de la interfaz de Databricks. La región se muestra debajo de cada nombre de espacio de trabajo (por ejemplo, eastus). También puede encontrarlo en la consola de la cuenta en Áreas de trabajo.

Para consultar la disponibilidad por región, véase cuotas de ingesta de Zerobus.

Creación o identificación de la tabla de destino

Identifique la tabla de destino en la que desea ingerir datos. Para crear una nueva tabla de destino, ejecute el CREATE TABLE comando SQL. Por ejemplo, cree una nueva tabla denominada unity.default.air_quality.

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

Zerobus Ingest puede escribir tanto en tablas Delta gestionadas como en tablas de streaming, que funcionan de la misma manera, con los mismos límites y cuotas.

Nota:

Para la ingesta de OpenTelemetry, las tablas deben usar esquemas predefinidos para cada tipo de señal (seguimientos, registros, métricas). Consulte Creación de tablas de destino en el catálogo de Unity.

Tu esquema de tabla es el contrato para lo que Zerobus Ingest acepta, y Zerobus Ingest nunca lo evoluciona automáticamente. El esquema de planificación cambia de forma proactiva: evoluciona primero la tabla y luego actualiza a los productores. Zerobus Ingest escribe en una ubicación de respaldo persistente los registros que dejan de ser compatibles tras un cambio incompatible en la tabla, en lugar de descartarlos. Consulta Gestión de esquemas y Recuperación de datos de la ubicación de respaldo persistente.

Por defecto, Zerobus Ingest rechaza registros con campos que no coinciden con el esquema de la tabla objetivo. Para capturar esos campos en lugar de perderlos, configure una columna de rescate. Véase la columna de rescate de Zerobus.

Crear un principal de servicio y conceder permisos

Un principal de servicio es una identidad especializada que ofrece más seguridad que las cuentas personalizadas. Para más información sobre las entidades de servicio y cómo usarlas para la autenticación, consulte Autorizar el acceso de la entidad de servicio a Azure Databricks con OAuth.

Puedes crear y administrar entidades de servicio mediante programación con la API REST o los SDK de Azure Databricks, o a través de la interfaz de usuario del espacio de trabajo, como se describe a continuación. Los permisos concedidos al final de esta sección son SQL que puedes ejecutar desde cualquier cliente.

  1. Para crear un principal de servicio, ve a Configuración>Identidad y Acceso.

  2. En Entidades de servicio, seleccione Administrar.

  3. Haga clic en Agregar entidad de servicio.

  4. En la ventana Agregar principal de servicio, cree un nuevo principal de servicio haciendo clic en Agregar nuevo.

  5. Genere y guarde el identificador de cliente y el secreto de cliente para la entidad de servicio.

  6. Conceda los permisos necesarios para el catálogo, el esquema y la tabla al principal de servicio.

    1. En la página de Principal de Servicio , ve a la pestaña Configuraciones .
    2. Copie el id. de aplicación (UUID).
    3. Use el siguiente código SQL para conceder permisos, reemplazando el UUID de ejemplo y el catálogo, el nombre del esquema y los nombres de tabla si es necesario.
    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>`;
    

Escribir un cliente de software

Use un SDK de Zerobus en el lenguaje de programación preferido o la API REST para ingerir datos en la tabla de destino. Los SDKs son de código abierto. Para la biblioteca completa, documentación específica del lenguaje y ejemplos adicionales, véase el repositorio Zerobus SDK.

Los ejemplos siguientes utilizan ingest_record_offset, que conserva el orden en que envías los registros.

SDK de Python

se requiere Python 3.9 o superior. El SDK proporciona un alto rendimiento y una E/S de red eficiente a través de un tiempo de ejecución asincrónico. Admite JSON (más sencillo) y búferes de protocolo (recomendados para producción). El SDK también soporta tanto implementaciones de sincronización como asíncronas, así como métodos de ingestión basados en desplazamientos y futuros.

pip install databricks-zerobus-ingest-sdk

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

Los ejemplos anteriores utilizan el método basado en desplazamiento ingest_record_offset sin esperar al desplazamiento devuelto. Para conocer los métodos de ingesta disponibles, cuándo esperar la confirmación de durabilidad de un desplazamiento y cómo hacer un seguimiento del progreso con una función de devolución de llamada de confirmación, consulte Bloqueo y confirmación de mensajes.

Búferes de protocolo: Para una ingesta segura de tipos, pasa un descriptor protobuf a TableProperties (el formato se selecciona automáticamente). Genera un esquema a partir de tu tabla usando la generate_proto herramienta, compiláalo con protoc, y luego pasa el descriptor compilado para crear el flujo.

Arrow Flight: Para la ingestión en columnas o por lotes de datos Apache Arrow RecordBatch a través de la misma conexión gRPC, consulta Usar Arrow Flight con Zerobus Ingest. Requiere el complemento [arrow]: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.

Para obtener documentación completa, opciones de configuración, ingesta por lotes y ejemplos de búfer de protocolo, consulte el repositorio del SDK de Python.

Rust SDK

Se requiere Rust 1.70 o superior. El SDK utiliza E/S asíncrona y gRPC para la ingestión de alto rendimiento. Admite JSON (más sencillo) y búferes de protocolo (recomendados para producción).

En primer lugar, importe el paquete.

cargo add databricks-zerobus-ingest-sdk

O agréguelo a su Cargo.toml.

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

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

Búferes de protocolo: Para la ingesta con seguridad de tipos, use Búferes de protocolo mediante .compiled_proto(descriptor) en el constructor de flujos en lugar de .json(), donde descriptor es un prost_types::DescriptorProto. Genere los archivos necesarios mediante la herramienta generate_proto e impórtelos en el proyecto. Arrow Flight: Para la ingestión orientada a columnas o por lotes de datos de Apache Arrow RecordBatch a través de la misma conexión gRPC, consulte Usar Arrow Flight con Zerobus Ingest. Habilítelo con la función Cargo: cargo add databricks-zerobus-ingest-sdk --features arrow-flight.

Para obtener documentación completa, opciones de configuración, ingesta por lotes, generate_proto herramientas y ejemplos de búfer de protocolo, consulte el repositorio del SDK de Rust.

SDK de Java

se requiere Java 8 o superior. El SDK proporciona una latencia baja y una E/S de red eficiente para la ingesta de alto rendimiento. Admite JSON (más sencillo) y búferes de protocolo (recomendados para producción).

Maven:

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

Ejemplo 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();
        }
    }
}

Búferes de protocolo: Para una ingestión segura por tipo, se crea un ZerobusProtoStream con streamBuilder() y .compiledProto(...). Genere un esquema a partir de la tabla mediante la herramienta JAR agrupada y, a continuación, compílelo con protoc.

Arrow Flight: Para la ingestión columnar o por lotes de datos de Apache Arrow RecordBatch a través de la misma conexión gRPC, consulta Usar Arrow Flight con Zerobus Ingest.

Para obtener documentación completa, opciones de configuración, ingesta por lotes y ejemplos de búfer de protocolo, consulte el repositorio Java SDK.

Kit de Desarrollo de Software Go

Se requiere Go 1.21 o una versión superior. El SDK proporciona un alto rendimiento y una gran capacidad de procesamiento para la ingesta en streaming. Admite JSON (más sencillo) y búferes de protocolo (recomendados para producción).

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

Ejemplo JSON:

Para simplificar, los errores se omiten aquí. En el código de producción, compruebe siempre los errores.

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")
}

Búferes de Protocolo: Para la ingesta segura de tipos, usa Búferes de protocolo con RecordTypeProto (valor predeterminado) y proporciona una descriptorProto en las propiedades de la tabla. Cree un archivo .proto que coincida con el esquema de la tabla y ejecute el generate_proto script para ayudarle a importar los archivos en el proyecto.

Arrow Flight: Para la ingestión en columnas o por lotes de datos de Apache Arrow RecordBatch a través de la misma conexión gRPC, consulte Usar Arrow Flight con Zerobus Ingest.

Para obtener documentación completa, opciones de configuración, ingesta por lotes, herramientas de generate_proto y ejemplos de búfer de protocolo, consulte el repositorio del SDK de Go.

C++ SDK

Importante

El SDK de C++ está en beta.

Se requiere C++17 o posterior. El SDK proporciona streaming nativo gRPC, OAuth y recuperación automática a través de una interfaz RAII C++. Admite JSON para configuraciones sencillas y búferes de protocolo para cargas de trabajo de producción.

El SDK se distribuye como una agrupación de versiones precompilada por plataforma, por lo que no necesita una cadena de herramientas de Rust para usarla. Descargue el paquete para su plataforma (macOS, Linux, incluida la variante musl, o Windows) desde la página de lanzamientos, extráigalo y, a continuación, indique a CMake la ubicación del archivo FFI incluido en el paquete. El archivo se denomina libzerobus_ffi.a en macOS y Linux y zerobus_ffi.lib en 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

En Windows (PowerShell), apunte al archivo .lib en su lugar:

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

Para compilar el SDK a partir de una copia del código fuente en tu propio proyecto CMake, añádelo como un subdirectorio y enlaza el objetivo. También puede usar FetchContent para recuperarlo durante la configuración. Esto compila el FFI desde el origen de Rust, por lo que requiere una cadena de herramientas de Rust:

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

Para consumir un paquete precompilado a través de add_subdirectory en su lugar, establezca primero las rutas de FFI para que CMake enlace el archivo incluido en lugar de intentar compilarlo a partir del código fuente de Rust, que no está presente. Use zerobus_ffi.lib en 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)

Ejemplo JSON:

La ingesta es asíncrona y está canalizada. Los ingest_* métodos ponen en cola un registro y devuelven inmediatamente. Pon en cola el lote y llama a flush() una sola vez, en lugar de esperar después de cada registro.

#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;
}

Cada fallo lanza zerobus::ZerobusException, que incluye un mensaje y un indicador is_retryable(). Para realizar un seguimiento de la durabilidad en un flujo continuo sin bloquear, registre un AckCallback mediante StreamOptions::ack_callback. Las llamadas de retorno se ejecutan en serie en un hilo en segundo plano y deben ser noexcept. Consulta la documentación del SDK de C++ para conocer todos los detalles sobre subprocesos, directivas de vaciado y contratos de vida útil.

Para una ingesta segura en cuanto a tipos, puedes utilizar Protocol Buffers de una de estas dos formas:

  • Genere el esquema desde el catálogo de Unity con ProtoSchema::from_uc_json(). Esto compila un descriptor y un codificador JSON a proto directamente desde los metadatos de la tabla, por lo que no necesita ningún .proto archivo o protoc:

    1. Obtén el JSON de metadatos de la tabla de la API Get a table (GET /api/2.1/unity-catalog/tables/{full_name}). El principal del servicio necesita SELECT en la tabla.
    2. Pase los metadatos a ProtoSchema::from_uc_json() para crear el descriptor y el codificador.
    3. Establece TableProperties::descriptor_proto y, a continuación, realiza la ingesta con ingest_proto_records().
  • Compila un .proto registrado con protoc para la tipificación en tiempo de compilación.

Para ver un tutorial ejecutable, consulte los ejemplos de búferes de protocolo.

Para la ingesta columnar u orientada a lotes de registros de Apache Arrow a través de la misma conexión gRPC, consulta Usar Arrow Flight con Zerobus Ingest.

Para obtener documentación completa, opciones de configuración, ingesta por lotes y ejemplos de búfer de protocolo, consulte el repositorio del SDK de C++.

C# SDK

Importante

El SDK de C# / .NET está en Beta. El paquete Databricks.Zerobus es una versión preliminar.

.NET 8.0 o superior es obligatorio. El SDK proporciona streaming nativo gRPC, OAuth y recuperación automática. Admite JSON para configuraciones sencillas y búferes de protocolo para cargas de trabajo de producción. Arrow Flight no está disponible en el SDK de C#.

Agregue el paquete Databricks.Zerobus al proyecto:

dotnet add package Databricks.Zerobus

Ejemplo 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 devuelve el desplazamiento del disco y WaitForOffset bloquea hasta que ese registro sea duradero. Para ingerir un lote, usa IngestRecords, que toma un conjunto de registros y devuelve el último desplazamiento. El bloqueo por desplazamiento es opcional. Consulta Bloqueo de mensajes y acuse de recibo.

Búferes de protocolo: Para la ingestión segura por tipos, crea un flujo con sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), donde descriptorProto son los bytes serializados DescriptorProto de tu mensaje compilado, y luego ingira con stream.IngestRecord(protoBytes).

Para documentación completa, opciones de configuración y ejemplos de búfer de protocolo, consulte el repositorio SDK de C#.

TypeScript SDK

se requiere Node.js 16 o superior. El SDK ofrece un alto rendimiento con soporte asíncrono a través de JavaScript Promises. Admite JSON (más sencillo) y búferes de protocolo (recomendados para producción).

npm install @databricks/zerobus-ingest-sdk

Ejemplo 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();
}

Búferes de Protocolo: Para la ingesta segura de tipos, usa Búferes de protocolo con RecordType.Proto (valor predeterminado) y proporciona una descriptorProto en las propiedades de la tabla.

Arrow Flight: Para la ingesta de datos de Apache Arrow RecordBatch en formato columnar o por lotes a través de la misma conexión gRPC, consulta Usar Arrow Flight con Zerobus Ingest.

Para obtener documentación completa, opciones de configuración, ingesta por lotes y ejemplos de búfer de protocolo, consulte el repositorio del SDK de TypeScript.

API de REST

La API REST permite ingerir un único registro mediante el envío de una solicitud HTTP POST al /zerobus/v1/tables/<table-name>/insert punto de conexión. El propio registro se incluye en el cuerpo de la solicitud y debe estar en formato JSON.

En este ejemplo se explica cómo usar CURL para enviar datos a Zerobus Ingest mediante la API REST.

Headers

La solicitud requiere dos encabezados HTTP específicos para autenticar y dar formato a la solicitud correctamente.

  • Content-Type: application/json
    • Campo obligatorio para especificar el tipo de contenido. Actualmente, JSON es el único formato de mensaje admitido.
  • Autorización: token de portador <>
    • Reemplaza <token> por el token OAuth que has obtenido mediante el comando curl proporcionado más adelante.

Capturar token de OAuth: Estos tokens expiran cada hora y se deben actualizar. Para actualizarlos, vuelva a capturar el token de OAuth.

Rellene los parámetros siguientes:

  • $CATALOG, $SCHEMA, $TABLE, , $WORKSPACE_ID, $WORKSPACE_URL
  • $DATABRICKS_CLIENT_ID y $DATABRICKS_CLIENT_SECRET
    • Estos dos parámetros corresponden al principio de servicio que creó.
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')

Ingestión de registros:

Rellene los parámetros siguientes:

El cuerpo de la solicitud debe ser una lista de objetos 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 }]'

Si toda la información se rellena correctamente, debe recibir una respuesta JSON vacía con un código de estado HTTP de 200.

Gestión de errores

Los ejemplos anteriores muestran el camino feliz. En producción, envuelve la ingestión en gestión de errores. El SDK reintenta automáticamente las operaciones cuando se producen errores transitorios, como los problemas de red, gracias a su mecanismo de recuperación integrado. Los fallos de los que no puede recuperarse, como credenciales inválidas o una tabla ausente, aparecen como 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.
    ...

Los SDKs también se recuperan automáticamente de fallos transitorios y te permiten recuperar registros no reconocidos cuando un flujo falla de forma permanente. Para patrones de cliente resiliente y la referencia completa de errores, véase Patrones de recuperación y reintentos y gestión de errores con Ingesta Zerobus.

Pasos siguientes