Utiliser Zerobus Ingest
Cette page décrit comment ingérer des données à l'aide de Zerobus Ingest dans Lakeflow Connect.
Premiers pas avec Zerobus Ingest
Si vous disposez d'un pare-feu côté client, ajoutez l'adresse IP utilisée par Zerobus Ingest à votre liste d'autorisation. Pour afficher les adresses IP par région, voir Adresses IP et domaines pour les services et assets Databricks.
Avant de start, confirmez que Zerobus Ingest est disponible dans la région de votre Workspace. Voir Disponibilité de l'ingestion.
- Obtenez une URL Zerobus Ingest.
- Créez ou identifiez la table dans laquelle vous souhaitez ingérer des données.
- Créez un Service Principal et accordez des privilèges sur la table.
- Connectez un client ou un exportateur pour start à envoyer des données.
Choisissez le guide correspondant à votre cas d'utilisation :
-
Ingérez vos propres données : utilisez les SDK ou l'API REST de Zerobus Ingest avec un schéma que vous définissez. Veuillez suivre les instructions figurant sur cette page.
-
Ingérer des données OpenTelemetry : utilisez des SDK ou des collecteurs OpenTelemetry standard pour envoyer des traces, des logs et des métriques dans des schémas de table prédéfinis. Pour obtenir des instructions complètes, consultez Ingest OpenTelemetry data with Zerobus Ingest.
Choisir une interface
Zerobus Ingest prend en charge plusieurs interfaces, écrivant toutes directement dans des tables Delta de Unity Catalog. En bref :
-
SDK via gRPC : throughput soutenu le plus élevé, idéal pour les producteurs de streaming à haut volume.
-
REST : sans état, idéal pour les grands parcs de dispositifs de périphérie légers ou « bavards ».
-
OpenTelemetry (OTLP) : pour les systèmes émettant déjà des traces, des logs et des métriques OpenTelemetry. Consultez Ingérer des données OpenTelemetry avec Zerobus Ingest.
-
APIs compatibles Kafka (Bêta) : pour les producteurs qui utilisent déjà le protocole Kafka. Consultez Utiliser les APIs compatibles avec Kafka avec Zerobus Ingest.
Pour une comparaison complète et savoir comment choisir, consultez les protocoles API. Via les SDK, vous pouvez également choisir un format d'enregistrement (JSON, Protocol Buffers (protobuf) ou Apache Arrow). Voir Types de messages. Le reste de cette page utilise les SDK et l'API REST.
Obtenez l'URL de votre Workspace et l'Endpoint Zerobus Ingest
Votre URL de workspace apparaît dans le navigateur lorsque vous vous connectez. Bien que l'URL complète suive le format https://<databricks-instance>.com/o=XXXXX, l'URL du workspace comprend tout ce qui précède le /o=XXXXX. Par exemple, étant donné l'URL complète suivante, vous pouvez déterminer l'URL du workspace et l'ID du workspace.
- URL complète :
https://abcd-teste2-test-spcse2.cloud.databricks.com/?o=2281745829657864# - URL du workspace :
https://abcd-teste2-test-spcse2.cloud.databricks.com - ID du Workspace :
2281745829657864
L'endpoint de serveur dépend du workspace et de la région :
- Endpoint du serveur :
<workspace-id>.zerobus.<region>.cloud.databricks.com
Pour trouver la région de votre workspace, ouvrez le sélecteur de workspace dans la barre de navigation supérieure de l'interface utilisateur Databricks. La région est affichée sous chaque nom de workspace (par exemple, us-west-2). Vous pouvez également la trouver dans la console du compte sous Workspaces .
Pour connaître la disponibilité par région, consultez les quotas de Zerobus Ingest.
Créez ou identifiez la table cible
Identifiez la table cible dans laquelle vous souhaitez ingérer des données. Pour créer une nouvelle table cible, exécutez la commande CREATE TABLE SQL. Par exemple, créez une nouvelle table nommée unity.default.air_quality.
CREATE TABLE unity.default.air_quality (
device_name STRING, temp INT, humidity LONG);
Zerobus Ingest peut écrire à la fois dans des tables Delta gérées et dans des tables de streaming, qui fonctionnent de la même manière, avec les mêmes limites et quotas.
Pour l'ingestion OpenTelemetry, les tables doivent utiliser des schémas prédéfinis pour chaque type de signal (traces, logs, métriques). Consultez Créer des tables cibles dans Unity Catalog.
Le schéma de votre table constitue le contrat de ce que Zerobus Ingest accepte, et Zerobus Ingest ne le fait jamais évoluer automatiquement. Planifiez les changements de schéma de manière proactive : faites d'abord évoluer la table, puis mettez à jour les producteurs. Zerobus Ingest écrit les enregistrements qui ne correspondent plus après un changement de table critique vers un emplacement de fallback durable au lieu de les supprimer. Voir Gestion de schémas et Récupération des données depuis l'emplacement de fallback durable.
Par default, Zerobus Ingest rejette les enregistrements dont les champs ne correspondent pas au schéma de la table cible. Pour capturer ces champs au lieu de les perdre, configurez une colonne de données sauvées. Voir colonne de données sauvées Zerobus.
Créer un Service Principal et accorder des autorisations
Un service principal est une identité spécialisée qui offre une meilleure sécurité que les comptes personnalisés. Pour plus d’information sur les Service Principal et sur la manière de les utiliser pour l’authentification, consultez Autoriser l’accès d’un Service Principal à Databricks avec OAuth.
Vous pouvez créer et gérer des Service Principals par programmation avec l'API REST ou les SDK Databricks, ou via l'interface utilisateur du workspace comme décrit ci-dessous. Les attributions de permissions à la fin de cette section sont du SQL que vous pouvez exécuter depuis n'importe quel client.
-
Pour créer un Service Principal, accédez à Settings > Identity and Access .
-
Sous Service Principal , sélectionnez Gérer .
-
Cliquez sur Ajouter un Service Principal .
-
Dans la fenêtre Ajouter un Service Principal , créez un nouveau Service Principal en cliquant sur Ajouter .
-
Générez et enregistrez l'ID client et le secret client pour le Service Principal.
-
Accordez au service principal les autorisations requises pour le catalogue, le schéma et la table.
- Sur la page Service Principal , accédez à l'onglet tab .
- Copiez l' ID d'application (UUID).
- Utilisez le SQL suivant pour accorder des autorisations, en remplaçant l'UUID d'exemple et le catalogue, le nom du schéma et les noms de table si nécessaire.
SQLGRANT 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>`;
Écrire un client
Utilisez un SDK Zerobus dans votre langage de programmation préféré ou l'API REST pour ingérer des données dans votre table cible. Les SDK sont open source. Pour la bibliothèque complète, la documentation spécifique au langage et des exemples supplémentaires, consultez le repository du SDK Zerobus.
Les exemples ci-dessous utilisent ingest_record_offset, qui préserve l’ordre dans lequel vous envoyez les enregistrements.
- Python SDK
- Rust SDK
- Java SDK
- Go SDK
- C++ SDK
- C# SDK
- TypeScript SDK
- REST API
Python 3.9 ou une version ultérieure est requis. Le SDK offre un throughput élevé et des E/S réseau efficaces grâce à un runtime asynchrone. Il prend en charge JSON (le plus simple) et Protocol Buffers (recommandé pour la production). Le SDK prend également en charge les implémentations synchrones et asynchrones, ainsi que les méthodes d'ingestion basées sur les offsets et sur les futures.
pip install databricks-zerobus-ingest-sdk
Exemple 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.us-west-2.cloud.databricks.com"
DATABRICKS_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"
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()
Les exemples ci-dessus utilisent la méthode ingest_record_offset basée sur le décalage sans attendre le décalage renvoyé. Pour en savoir plus sur les méthodes d'ingestion disponibles, sur le moment où attendre une confirmation de durabilité sur un décalage et sur la façon de suivre la progression avec un rappel d'accusé de réception, consultez Message blocking and acknowledgment.
Protocol Buffers : pour une ingestion typée en toute sécurité, transmettez un descripteur protobuf à TableProperties (le format est sélectionné automatiquement). Générez un schéma à partir de votre table en utilisant l’outil generate_proto, compilez-le avec protoc, puis transmettez le descripteur compilé pour créer le Stream.
Arrow Flight : pour une ingestion en colonnes ou par batch de données Apache Arrow RecordBatch via la même connexion gRPC, consultez Utiliser Arrow Flight avec Zerobus Ingest. Nécessite le supplément [arrow] : pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.
Pour une documentation complète, les options de configuration, l'ingestion par lots et les exemples de tampon de protocole, consultez le repository Python du SDK.
Rust 1.70 ou une version ultérieure est requis. Le SDK utilise des E/S asynchrones et gRPC pour une ingestion à haut throughput. Il prend en charge JSON (le plus simple) et Protocol Buffers (recommandé pour la production).
Tout d'abord, importez le package.
cargo add databricks-zerobus-ingest-sdk
Ou ajoutez-le à votre Cargo.toml.
[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication
Exemple 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://dbc-a1b2c3d4-e5f6.cloud.databricks.com";
const SERVER_ENDPOINT: &str = "https://1234567890123456.zerobus.us-west-2.cloud.databricks.com";
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 : pour une ingestion typée, utilisez Protocol Buffers via .compiled_proto(descriptor) sur le générateur de Stream au lieu de .json(), où descriptor est un prost_types::DescriptorProto. Générez les fichiers nécessaires à l'aide de l'outil generate_proto et importez-les dans votre projet.
Arrow Flight : pour une ingestion en colonnes ou orientée batch de données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Utiliser Arrow Flight avec Zerobus Ingest. Activez avec la fonctionnalité Cargo : cargo add databricks-zerobus-ingest-sdk --features arrow-flight.
Pour une documentation complète, les options de configuration, l'ingestion par batch, l'outil generate_proto et les exemples de Protocol Buffer, consultez le repository du SDK Rust.
Java 8 ou une version ultérieure est requis. Le SDK offre une faible latence et des E/S réseau efficaces pour une ingestion à haut throughput. Il prend en charge JSON (le plus simple) et Protocol Buffers (recommandé pour la production).
Maven:
<dependency>
<groupId>com.databricks</groupId>
<artifactId>zerobus-ingest-sdk</artifactId>
<version>0.2.0</version>
</dependency>
Exemple 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.us-west-2.cloud.databricks.com";
private static final String DATABRICKS_WORKSPACE_URL =
"https://dbc-a1b2c3d4-e5f6.cloud.databricks.com";
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();
}
}
}
Protocol Buffers : pour une ingestion typée, créez un ZerobusProtoStream avec streamBuilder() et .compiledProto(...). Générez un schéma à partir de votre table à l’aide de l’outil JAR fourni, puis compilez-le avec protoc.
Arrow Flight : pour une ingestion orientée colonnes ou batch de données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Use Arrow Flight with Zerobus Ingest.
Pour une documentation complète, les options de configuration, l'ingestion batch et des exemples de Protocol Buffer, consultez le repository Java SDK.
Go 1.21 ou une version ultérieure est requis. Le SDK fournit un throughput et des performances élevés pour l'ingestion en streaming. Il prend en charge JSON (le plus simple) et Protocol Buffers (recommandé pour la production).
go get github.com/databricks/zerobus-sdk/go@latest
Exemple JSON :
Par souci de simplicité, les erreurs sont ignorées ici. Dans le code de production, vérifiez toujours les erreurs.
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.us-west-2.cloud.databricks.com"
DatabricksWorkspaceURL = "https://dbc-a1b2c3d4-e5f6.cloud.databricks.com"
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 : Pour une ingestion sécurisée des types, utilisez les Protocol Buffers avec RecordTypeProto (default) et fournissez un descriptorProto dans les propriétés de la table. Créer un .proto fichier correspondant à votre schéma de table et exécutez le script generate_proto pour vous aider à importer les fichiers dans votre projet.
Arrow Flight : pour une ingestion orientée colonnes ou batch de données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Use Arrow Flight with Zerobus Ingest.
Pour une documentation complète, les options de configuration, l'ingestion par batch, l'outil generate_proto et des exemples de Protocol Buffer, consultez le repository du SDK Go.
Bêta
Le SDK C++ est en bêta.
C++17 ou une version ultérieure est requis. Le SDK fournit le streaming gRPC natif, OAuth et une récupération automatique via une interface C++ RAII. Il prend en charge JSON pour les configurations simples et Protocol Buffers pour les workloads de production.
Le SDK est fourni sous forme de bundle de version préconstruit et par plateforme ; vous n’avez donc pas besoin d’une chaîne d’outils Rust pour l’utiliser. download le bundle pour votre plateforme (macOS, Linux incluant musl, ou Windows) depuis la page des versions, extrayez-le, puis pointez CMake vers l’archive FFI bundle. L’archive est nommée libzerobus_ffi.a sur macOS et Linux et zerobus_ffi.lib sur 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
Sur Windows (PowerShell), pointez plutôt vers l’archive .lib :
cmake -S cpp -B build `
-DZEROBUS_FFI_LIBRARY="$PWD/lib/zerobus_ffi.lib" `
-DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j
Pour générer le SDK à partir d’une extraction de source dans votre propre projet CMake, ajoutez-le en tant que sous-répertoire et liez la cible. Vous pouvez également utiliser FetchContent pour le récupérer au moment de la configuration. Ceci génère le FFI à partir de la source Rust, il nécessite donc une chaîne d’outils Rust :
add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)
Pour consommer un bundle préconstruit via add_subdirectory à la place, définissez d’abord les chemins FFI afin que CMake lie l’archive bundle plutôt que d’essayer de la construire à partir d’une source Rust qui n’est pas présente. Utilisez zerobus_ffi.lib sur 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)
Exemple JSON :
L’ingestion est asynchrone et pipelinée. Les méthodes ingest_* mettent un enregistrement en file d’attente et renvoient immédiatement. Mettez le batch en file d’attente et appelez flush() une fois, plutôt que d’attendre après chaque enregistrement.
#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.us-west-2.cloud.databricks.com";
const std::string DATABRICKS_WORKSPACE_URL = "https://dbc-a1b2c3d4-e5f6.cloud.databricks.com";
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;
}
Chaque échec génère zerobus::ZerobusException, qui contient un message et un indicateur is_retryable(). Pour suivre la durabilité sur un flux continu sans blocage, enregistrez un AckCallback via StreamOptions::ack_callback. Les rappels (callbacks) s’exécutent de manière sérialisée sur un thread en arrière-plan et doivent être noexcept. Consultez la documentation du SDK C++ pour obtenir les informations complètes sur le threading, la politique de vidage (drain-policy) et le contrat de durée de vie.
Pour une ingestion à typage sécurisé, vous pouvez utiliser Protocol Buffers de deux manières :
-
Générez le schéma à partir d’Unity Catalog avec
ProtoSchema::from_uc_json(). Ceci génère un descripteur et un encodeur JSON-to-proto directement à partir des métadonnées de la table, il ne nécessite donc aucun fichier.protoouprotoc:- Récupérez le JSON des métadonnées de la table à partir de l’API Get a table (
GET /api/2.1/unity-catalog/tables/{full_name}). Le Service Principal a besoin deSELECTsur la table. - Transmettez les métadonnées à
ProtoSchema::from_uc_json()pour construire le descripteur et l’encodeur. - Définissez
TableProperties::descriptor_proto, puis ingérez avecingest_proto_records().
- Récupérez le JSON des métadonnées de la table à partir de l’API Get a table (
-
Compilez un
.protoarchivé avecprotocpour le typage au moment de la compilation.
Pour un guide pas à pas exécutable, consultez les exemples Protocol Buffers.
Pour une ingestion orientée colonne ou par batch de batchs d’enregistrements Apache Arrow sur la même connexion gRPC, consultez Utiliser Arrow Flight avec Zerobus Ingest.
Pour obtenir une documentation complète, des options de configuration, des informations sur l’ingestion par batch et des exemples de Protocol Buffer, consultez le repository du SDK C++.
Bêta
Le SDK C# / .NET est en bêta. Le package Databricks.Zerobus est en version préliminaire.
.NET 8.0 ou une version ultérieure est requis. Le SDK fournit le streaming gRPC natif, OAuth et une récupération automatique. Il prend en charge JSON pour les configurations simples et Protocol Buffers pour les charges de travail de production. Arrow Flight n’est pas disponible dans le SDK C#.
Ajoutez le package Databricks.Zerobus à votre projet :
dotnet add package Databricks.Zerobus
Exemple 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.us-west-2.cloud.databricks.com";
const string DATABRICKS_WORKSPACE_URL = "https://dbc-a1b2c3d4-e5f6.cloud.databricks.com";
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 renvoie le décalage de l'enregistrement, et WaitForOffset se bloque jusqu'à ce que cet enregistrement soit durable. Pour ingérer un batch, utilisez IngestRecords, qui prend un tableau d'enregistrements et renvoie le dernier décalage. Le blocage sur le décalage est facultatif. Voir Blocage et confirmation des messages.
Protocol Buffers : pour une ingestion à typage sécurisé, créez un Stream avec sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), où descriptorProto correspond aux DescriptorProto octets sérialisés de votre message compilé, puis procédez à l'ingestion avec stream.IngestRecord(protoBytes).
Pour une documentation complète, des options de configuration et des exemples de Protocol Buffer, consultez le repository du SDK C#.
Node.js 16 ou une version supérieure est requis. Le SDK offre des performances élevées avec une prise en charge asynchrone via les promesses JavaScript. Il prend en charge JSON (le plus simple) et Protocol Buffers (recommandé pour la production).
npm install @databricks/zerobus-ingest-sdk
Exemple 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.us-west-2.cloud.databricks.com';
const DATABRICKS_WORKSPACE_URL = 'https://dbc-a1b2c3d4-e5f6.cloud.databricks.com';
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 : Pour une ingestion sécurisée, utilisez Protocol Buffers avec RecordType.Proto (default) et fournissez un descriptorProto dans les propriétés de la table.
Arrow Flight : pour une ingestion orientée colonnes ou batch de données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Use Arrow Flight with Zerobus Ingest.
Pour une documentation complète, des options de configuration, l'ingestion par batch et des exemples de Protocol Buffer, consultez le repository du SDK TypeScript.
L'API REST vous permet d'ingérer un seul enregistrement en envoyant une requête HTTP POST à l'Endpoint /zerobus/v1/tables/<table-name>/insert. L'enregistrement lui-même est inclus dans le corps de la requête et doit être au format JSON.
Cet exemple vous explique comment utiliser CURL pour pousser des données vers Zerobus Ingest à l'aide de l'API REST.
En-têtes
La requête nécessite deux en-têtes HTTP spécifiques pour authentifier et formater correctement la requête.
-
Content-Type: application/json- Champ obligatoire pour spécifier le type de contenu. Actuellement, JSON est le seul format de message pris en charge.
-
Authorization: Bearer <token>- Remplacez
<token>par le jeton OAuth que vous avez récupéré à l’aide de la commande curl fournie plus loin.
- Remplacez
Récupérer le jeton OAuth : Ces jetons expirent toutes les heures et doivent être actualisés. Vous pouvez les refresh en récupérant de nouveau le jeton OAuth.
Remplissez les paramètres suivants :
$CATALOG,$SCHEMA,$TABLE,$WORKSPACE_ID,$WORKSPACE_URL$DATABRICKS_CLIENT_IDet$DATABRICKS_CLIENT_SECRET- Ces deux parameters correspondent au principe de service que vous avez créé.
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')
Ingestion d'enregistrements :
Remplissez les paramètres suivants :
-
$ZEROBUS_ENDPOINT- Tel que défini dans la section Obtenez l'URL de votre Workspace et l'Endpoint Zerobus Ingest.
-
$CATALOG,$SCHEMA,$TABLE,$WORKSPACE_ID,$WORKSPACE_URL -
$OAUTH_TOKEN- Cela a été créé à l’étape précédente.
Le corps de la requête doit être une liste d'objets 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 toutes les informations sont correctement remplies, vous devriez recevoir une réponse JSON vide avec un code d’état HTTP de 200.
Gérer les erreurs
Les exemples ci-dessus illustrent le chemin nominal. En production, encapsulez l'ingestion dans une gestion des erreurs. Le SDK réessaie automatiquement les erreurs transitoires, telles que les problèmes réseau, grâce à sa récupération intégrée. Les échecs dont il ne peut pas se remettre, tels que des identifiants non valides ou une table manquante, apparaissent sous la forme 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.
...
Les SDK se remettent également automatiquement des défaillances transitoires et vous permettent de récupérer les enregistrements non acquittés lorsqu'un Stream échoue de manière permanente. Pour les modèles de client résilient et la référence complète des erreurs, consultez Modèles de récupération et de nouvelle tentative et Gestion des erreurs Zerobus Ingest.