Utilisez le connecteur Zerobus Ingest
Cette page décrit comment ingérer des données à l'aide du connecteur Zerobus Ingest dans Lakeflow Connect.
Choisir une interface
Zerobus Ingest prend en charge les interfaces gRPC, REST et OpenTelemetry (OTLP). Les SDK fournissent des clients personnalisés basés sur gRPC avec une interface conviviale pour les développeurs afin de créer des applications à haut throughput. L'interface REST gère les contraintes architecturales lorsqu'il s'agit de flottes massives de dispositifs "volumineux". L'interface OTLP accepte les données OpenTelemetry standard sans nécessiter de bibliothèques personnalisées.
-
Les SDK avec la « taxe de connexion » gRPC : gRPC se spécialise dans les performances à haut throughput grâce à des connexions persistantes. Chaque Stream ouvert est comptabilisé dans vos quotas de simultanéité. Les SDK prennent en charge trois formats d'enregistrement sur gRPC :
- JSON : le plus simple, aucune définition de schéma requise.
- Protocol Buffers : Recommandé pour les charges de travail de production orientées lignes.
- Apache Arrow Flight : recommandé pour le format de colonne ou les charges de travail par batch. Voir Utiliser Arrow Flight avec Zerobus Ingest.
-
La « taxe de throughput » REST : REST nécessite une négociation complète pour chaque mise à jour, ce qui le rend sans état. REST convient bien aux cas d'utilisation des appareils en périphérie où l'état est rarement rapporté.
-
OpenTelemetry (OTLP) : Si vous utilisez déjà les SDK ou collecteurs OpenTelemetry, l'endpoint OTLP ingère les traces, les logs et les métriques dans les tables Delta de Unity Catalog sans aucune intégration personnalisée requise. Pour plus d'informations, consultez Ingérer des données OpenTelemetry avec Zerobus Ingest.
Utilisez les SDK avec proto pour les flux à volume élevé orientés lignes, ou les SDK avec Arrow Flight pour le format de colonne ou les charges de travail en batch. Utilisez REST pour les flottes d'appareils massives à basse fréquence, et OTLP pour les environnements déjà instrumentés avec OpenTelemetry.
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 de Databricks — la région est affichée sous chaque nom de Workspace (par exemple, us-west-2). Vous pouvez également le trouver dans la console de compte sous **Workspaces**.
Pour la disponibilité régionale, consultez les limites du connecteur 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);
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.
By default, les enregistrements dont les champs ne correspondent pas au schéma de la table cible sont rejetés. Pour capturer ces champs au lieu de les perdre, configurez une colonne de récupération. Consultez Capturez les champs non conformes avec la colonne de récupération Zerobus.
Ingérer dans une table de streaming
Bêta
L'ingestion dans les tables en streaming à l'aide du connecteur Zerobus Ingest est en bêta. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité depuis la page Aperçus . Consultez Gérer les aperçus Databricks.
Pour créer une nouvelle table de streaming, exécutez la commande SQL CREATE STREAMING TABLE. Par exemple :
CREATE STREAMING TABLE unity.default.air_quality (
device_name STRING, temp INT, humidity LONG);
Une fois la table de streaming créée, ingérez-y des données en utilisant l’une des interfaces de Écrire un client, exactement comme vous le feriez pour une table Delta standard.
Créer un Service Principal et accorder des autorisations
Un Service Principal est une identité spécialisée qui offre plus de sécurité que les comptes personnalisés. Pour plus d'information concernant les service principals et comment les utiliser pour l'authentification, consultez Autoriser l'accès des Service Principal à Databricks avec OAuth.
-
Pour créer un Service Principal, accédez à **Paramètres** > **Identité et accès**.
-
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.
- Dans la page **Service principal**, accédez à l'**tab** Configurations.
- 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.
- Python SDK
- Rust SDK
- Java SDK
- Go SDK
- TypeScript SDK
- REST API
Python 3.9 ou une version ultérieure est requis. Le SDK utilise les liaisons PyO3 vers le SDK Rust haute performance, offrant un throughput jusqu'à 40 fois supérieur à celui de Python pur et des E/S réseau efficaces grâce au runtime asynchrone de Rust. Il prend en charge le JSON (le plus simple) et les Protocol Buffers (recommandé pour la production). Le SDK prend également en charge les implémentations synchrones et asynchrones, ainsi que 3 méthodes d'ingestion différentes (basées sur le futur, basées sur le décalage et « fire-and-forget »).
pip install databricks-zerobus-ingest-sdk
Exemple JSON :
import logging
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import RecordType, StreamConfigurationOptions, 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)
options = StreamConfigurationOptions(record_type=RecordType.JSON)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties, options)
try:
for i in range(1000):
record_dict = {
"device_name": f"sensor-{i}",
"temp": 20 + i % 15,
"humidity": 50 + i % 40
}
offset = stream.ingest_record_offset(record_dict)
# Optional: Wait for durability confirmation
stream.wait_for_offset(offset)
finally:
stream.close()
**Rappel d'accusé de réception :** Pour suivre la progression de l'ingestion de manière asynchrone, utilisez ack_callback l'option. Passez une sous-classe de AckCallback avec les méthodes on_ack(offset: int) et on_error(offset: int, error_message: str), qui sont appelées lorsque les enregistrements sont soit accusés de réception, soit en échec.
from zerobus.sdk.shared import AckCallback, StreamConfigurationOptions, RecordType
class MyAckCallback(AckCallback):
def on_ack(self, offset: int) -> None:
print(f"Record acknowledged at offset: {offset}")
def on_error(self, offset: int, error_message: str) -> None:
print(f"Error at offset {offset}: {error_message}")
options = StreamConfigurationOptions(
record_type = RecordType.JSON,
ack_callback = MyAckCallback()
)
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 (Beta) : Pour une ingestion orientée colonne ou par batch des données Apache Arrow RecordBatch sur la même connexion gRPC, voir Utiliser Arrow Flight avec Zerobus Ingest. Nécessite le [arrow] supplémentaire : 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 supérieure est requis. Le SDK tire parti des E/S asynchrones et de gRPC pour une ingestion à haut throughput et sert de base à tous les autres SDK. Il prend en charge le JSON (le plus simple) et les Protocol Buffers (recommandés 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?;
let offset = stream.ingest_record_offset(
JsonString("{
\"device_name\": \"sensor\",
\"temp\": 22,
\"humidity\": 55}".to_string())).await?;
stream.wait_for_offset(offset).await?;
println!("Record ingested successfully");
stream.close().await?;
println!("Stream closed successfully");
Ok(())
}
Protocol Buffers : Pour une ingestion de type sécurisé, utilisez les 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 (Bêta) : pour l'ingestion orientée colonne ou par batch de données Apache Arrow RecordBatch sur la même connexion gRPC, voir Utiliser Arrow Flight avec Zerobus Ingest. Activer 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 utilise des liaisons JNI (Java Native Interface) vers le SDK Rust haute performance, offrant une latence inférieure à celle du gRPC Java pur et un I/O réseau efficace via le runtime asynchrone Rust. Il prend en charge le JSON (le plus simple) et les Protocol Buffers (recommandés 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.createJsonStream(
TABLE_NAME, CLIENT_ID, CLIENT_SECRET
).join();
try {
long lastOffset = 0;
for (int i = 0; i < 100; i++) {
String record = String.format(
"{\"device_name\": \"sensor-%d\", \"temp\": 22, \"humidity\": 55}", i
);
lastOffset = stream.ingestRecordOffset(record);
}
stream.waitForOffset(lastOffset);
} finally {
stream.close();
}
}
}
Protocol Buffers : pour une ingestion sécurisée des types, utilisez ZerobusProtoStream avec createProtoStream(). Générez un schéma à partir de votre table à l'aide de l'outil JAR fourni, puis compilez-le avec protoc.
Arrow Flight (Bêta) : Pour l'ingestion orientée colonne ou par batch des données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Utiliser Arrow Flight avec 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 version supérieure est requis. Le SDK enveloppe le SDK Rust haute performance en utilisant CGO et FFI, offrant le même throughput et les mêmes performances. Il prend en charge le JSON (le plus simple) et les Protocol Buffers (recommandés 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()
offset, _ := stream.IngestRecordOffset(`{
"device_name": "sensor-001",
"temp": 20,
"humidity": 60
}`)
_ = stream.WaitForOffset(offset)
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 (Bêta) : Pour l'ingestion orientée colonne ou par batch des données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Utiliser Arrow Flight avec 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.
Node.js 16 ou version supérieure est requis. Le SDK enveloppe le SDK Rust haute performance à l'aide de liaisons natives NAPI-RS, offrant des performances natives avec des futures Rust mappées sur des Promises JavaScript. Il prend en charge le JSON (le plus simple) et les Protocol Buffers (recommandés 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 {
let lastOffset = BigInt(0);
for (let i = 0; i < 100; i++) {
const record = { device_name: `sensor-${i}`, temp: 22, humidity: 55 };
lastOffset = await stream.ingestRecordOffset(record);
}
await stream.waitForOffset(lastOffset);
} 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 (Bêta) : Pour l'ingestion orientée colonne ou par batch des données Apache Arrow RecordBatch sur la même connexion gRPC, consultez Utiliser Arrow Flight avec 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.