Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Diese Seite beschreibt, wie man Daten mit Zerobus Ingest in Lakeflow Connect einführt.
Starten Sie mit Zerobus Ingest
Bevor Sie beginnen, bestätigen Sie, dass Zerobus Ingest in der Region Ihres Arbeitsbereichs verfügbar ist. Siehe Verfügbarkeit der Datenerfassung.
- Rufen Sie eine Zerobus Ingest-URL ab.
- Erstellen oder identifizieren Sie die Tabelle, in die Sie Daten aufnehmen möchten.
- Erstellen Sie einen Dienstprinzipal, und gewähren Sie der Tabelle Berechtigungen.
- Verbinden Sie einen Client oder Exporter, um mit dem Senden von Daten zu beginnen.
Wählen Sie den Leitfaden für Ihren Anwendungsfall aus:
Eigene Daten aufnehmen: Verwenden Sie die Zerobus Ingest-SDKs oder die REST-API mit einem Schema, das Sie definieren. Befolgen Sie die Anweisungen auf dieser Seite.
OpenTelemetry-Daten erfassen: Verwenden Sie standardmäßige OpenTelemetry-SDKs oder ‑Sammler, um Ablaufverfolgungen, Protokolle und Metriken in vordefinierte Tabellenschemas zu senden. Vollständige Anweisungen finden Sie unter Ingest OpenTelemetry-Daten mit Zerobus Ingest.
Auswählen einer Schnittstelle
Zerobus Ingest unterstützt mehrere Schnittstellen, die alle direkt in Unity Catalog Delta-Tabellen schreiben. Kurz gesagt:
- SDKs gegenüber gRPC: höchster anhaltender Durchsatz, am besten für Streaming-Produzenten mit hohem Volumen.
- REST: zustandslos, am besten geeignet für große Bestände leichter oder „gesprächiger“ Edge-Geräte.
- OpenTelemetry (OTLP): für Systeme, die bereits OpenTelemetrie-Traces, Logs und Metriken aussenden. Siehe OpenTelemetry-Daten mit Zerobus Ingest erfassen.
- Kafka-kompatible APIs (Beta): für Produzenten, die bereits das Kafka-Protokoll sprechen. Siehe Use Kafka-compatible APIs with Zerobus Ingest.
Für einen vollständigen Vergleich und die Auswahl siehe API-Protokolle. Im Gegensatz zu den SDKs können Sie auch ein Datensatzformat wählen (JSON, Protocol Buffers (protobuf) oder Apache Arrow). Siehe Nachrichtentypen. Der Rest dieser Seite verwendet die SDKs und die REST-API.
Rufen Sie die Arbeitsbereichs-URL und den Zerobus-Ingestion-Endpunkt ab.
Ihre Arbeitsbereichs-URL wird im Browser angezeigt, wenn Sie sich anmelden. Während die vollständige URL dem Format https://<databricks-instance>.net/o=XXXXXfolgt, besteht die Arbeitsbereichs-URL aus allem vor dem /o=XXXXX. Mit der folgenden vollständigen URL können Sie beispielsweise die Arbeitsbereichs-URL und die Arbeitsbereichs-ID ermitteln.
- Vollständige URL:
https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864# - Arbeitsbereichs-URL:
https://abcd-teste2-test-spcse2.azuredatabricks.net - Arbeitsbereichs-ID:
2281745829657864
Der Serverendpunkt hängt vom Arbeitsbereich und der Region ab:
- Serverendpunkt:
<workspace-id>.zerobus.<region>.azuredatabricks.net
Um Ihre Arbeitsbereichsregion zu finden, öffnen Sie den Arbeitsbereichswechsel in der oberen Navigationsleiste der Databricks-Benutzeroberfläche. Der Bereich wird unter jedem Arbeitsbereichsnamen angezeigt (zum Beispiel ). eastus Sie finden sie auch in der Kontokonsole unter "Arbeitsbereiche".
Informationen zur regionalen Verfügbarkeit finden Sie unter Zerobus Ingest-Quoten.
Erstellen oder Identifizieren der Zieltabelle
Identifizieren Sie die Zieltabelle, in die Sie Daten aufnehmen möchten. Führen Sie den CREATE TABLE SQL-Befehl aus, um eine neue Zieltabelle zu erstellen. Erstellen Sie beispielsweise eine neue Tabelle mit dem Namen unity.default.air_quality.
CREATE TABLE unity.default.air_quality (
device_name STRING, temp INT, humidity LONG);
Zerobus Ingest kann sowohl in verwaltete Delta-Tabellen als auch in Streaming-Tabellen schreiben, die auf dieselbe Weise arbeiten, mit denselben Limits und Quoten.
Hinweis
Für die Aufnahme von OpenTelemetry müssen Tabellen vordefinierte Schemata für jeden Signaltyp (Ablaufverfolgungen, Protokolle, Metriken) verwenden. Siehe Erstellen von Zieltabellen im Unity-Katalog.
Dein Tabellenschema ist der Vertrag für das, was Zerobus Ingest akzeptiert, und Zerobus Ingest entwickelt es nie automatisch weiter. Planen Sie Schemaänderungen proaktiv: Entwickeln Sie zuerst die Tabelle weiter und aktualisieren Sie dann die Datenproduzenten. Zerobus Ingest schreibt Datensätze, die nach einem Breaking Change an der Tabelle nicht mehr in die Tabelle passen, an einen dauerhaften Fallback-Speicherort, anstatt sie zu verwerfen. Siehe Schema-Management und Wiederherstellung von Daten aus dem dauerhaften Rückfall-Standort.
Standardmäßig lehnt Zerobus Ingest Datensätze mit Feldern ab, die nicht mit dem Schema der Zieltabelle übereinstimmen. Um diese Felder zu erfassen, anstatt sie zu verlieren, konfigurieren Sie eine Rettungsspalte. Siehe Zerobus-Rettungssäule.
Erstellen Sie einen Dienstprinzipal und erteilen Sie Berechtigungen
Ein Dienstprinzipal ist eine spezielle Identität, die mehr Sicherheit bietet als personalisierte Konten. Weitere Informationen zu Service Principals und deren Verwendung für die Authentifizierung finden Sie unter Service Principal Access zu Azure Databricks mit OAuth autorisieren.
Sie können Service Principals programmatisch mit der Azure Databricks REST API oder SDKs erstellen und verwalten oder über die unten beschriebene Workspace-UI. Die Berechtigungen am Ende dieses Abschnitts sind SQL, das du von jedem Client ausführen kannst.
Um einen Service Principal zu erstellen, gehe zu Einstellungen>Identität und Zugriff.
Wählen Sie unter "Dienstprinzipale" die Option "Verwalten" aus.
Klicken Sie auf Dienstprinzipal hinzufügen.
Erstellen Sie im Fenster "Dienstprinzipal hinzufügen " einen neuen Dienstprinzipal, indem Sie auf "Neu hinzufügen" klicken.
Generieren und speichern Sie die Client-ID und den geheimen Clientschlüssel für den Dienstprinzipal.
Erteilen Sie die erforderlichen Berechtigungen für den Katalog, das Schema und die Tabelle an den Dienstprinzipal.
- Wechseln Sie auf der Seite Service principal zur Registerkarte Konfigurationen.
- Kopieren Sie die Anwendungs-ID (UUID).
- Verwenden Sie die folgende SQL-Datei, um Berechtigungen zu erteilen, und ersetzen Sie bei Bedarf die Beispiel-UUID und den Katalognamen, den Schemanamen und die Tabellennamen.
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>`;
Entwicklung eines Clients
Verwenden Sie ein Zerobus SDK in Ihrer bevorzugten Programmiersprache oder der REST-API, um Daten in Ihre Zieltabelle aufzunehmen. Die SDKs sind Open Source. Für die vollständige Bibliothek, sprachspezifische Dokumentation und weitere Beispiele siehe das Zerobus SDK-Repository.
Die untenstehenden Beispiele verwenden ingest_record_offset, was die Reihenfolge beibehält, in der Sie Datensätze senden.
Python SDK
Python 3,9 oder höher ist erforderlich. Das SDK bietet einen hohen Durchsatz und effiziente Netzwerk-I/O über eine asynchrone Laufzeit. Es unterstützt JSON (einfachste) und Protokollpuffer (empfohlen für die Produktion). Das SDK unterstützt außerdem sowohl synchrone als auch asynchrone Implementierungen sowie offsetbasierte und Future-basierte Erfassungsmethoden.
pip install databricks-zerobus-ingest-sdk
JSON-Beispiel:
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()
Die obigen Beispiele verwenden die offsetbasierte ingest_record_offset Methode, ohne auf den zurückgegebenen Offset zu warten. Um mehr über die verfügbaren Erfassungsmethoden zu erfahren, wann auf die Bestätigung der Dauerhaftigkeit für einen Offset gewartet werden sollte und wie sich der Fortschritt mit einem Bestätigungs-Callback verfolgen lässt, finden Sie weitere Informationen unter Blockierung und Bestätigung von Nachrichten.
Protokollpuffer: Für eine typsichere Datenerfassung übergeben Sie einen Protobuf-Deskriptor an TableProperties (das Format wird automatisch ausgewählt). Generiere ein Schema aus deiner Tabelle mit dem generate_proto Tool, kompiliere es mit protoc, und gib dann den kompilierten Deskriptor weiter, um den Stream zu erstellen.
Arrow Flight: Für die spaltenbasierte oder batchorientierte Ingestion von Apache Arrow-RecordBatchDaten über dieselbe gRPC-Verbindung siehe Arrow Flight mit Zerobus Ingest verwenden. Erfordert das [arrow]- Extra: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.
Vollständige Dokumentation, Konfigurationsoptionen, Batchaufnahme und Protokollpufferbeispiele finden Sie im Python SDK-Repository.
Rust SDK
Rust 1.70 oder höher ist erforderlich. Das SDK verwendet asynchrone I/O und gRPC für Hochdurchsatzaufnahme. Es unterstützt JSON (einfachste) und Protokollpuffer (empfohlen für die Produktion).
Importieren Sie zuerst das Paket.
cargo add databricks-zerobus-ingest-sdk
Oder fügen Sie es zu Ihrem Cargo.toml.
[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication
JSON-Beispiel:
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: Für eine typensichere Erfassung verwenden Sie Protocol Buffers über .compiled_proto(descriptor) im Stream-Builder anstelle von .json(), wobei descriptor ein prost_types::DescriptorProto ist. Generieren Sie die erforderlichen Dateien mithilfe des generate_proto Tools, und importieren Sie sie in Ihr Projekt.
Arrow Flight: Informationen zur spaltenorientierten oder batchorientierten Erfassung von Apache Arrow-RecordBatchDaten über dieselbe gRPC-Verbindung finden Sie unter Verwenden von Arrow Flight mit Zerobus Ingest. Mit dem Cargo-Feature aktivieren: cargo add databricks-zerobus-ingest-sdk --features arrow-flight.
Vollständige Dokumentation, Konfigurationsoptionen, Batchaufnahme, generate_proto Tool- und Protokollpufferbeispiele finden Sie im Rust SDK-Repository.
Java SDK
Java 8 oder höher ist erforderlich. Das SDK bietet eine niedrige Latenz und effiziente Netzwerk-I/O für eine hohe Durchsatzaufnahme. Es unterstützt JSON (einfachste) und Protokollpuffer (empfohlen für die Produktion).
Maven:
<dependency>
<groupId>com.databricks</groupId>
<artifactId>zerobus-ingest-sdk</artifactId>
<version>0.2.0</version>
</dependency>
JSON-Beispiel:
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();
}
}
}
Protokollpuffer: Für die typsichere Aufnahme erstellen Sie ein ZerobusProtoStream mit streamBuilder() und .compiledProto(...). Generieren Sie ein Schema aus Ihrer Tabelle mithilfe des gebündelten JAR-Tools, und kompilieren Sie es mit protoc.
Arrow Flight: Informationen zur spaltenbasierten oder batchorientierten Ingestion von Apache Arrow-RecordBatchDaten über dieselbe gRPC-Verbindung finden Sie unter Arrow Flight mit Zerobus Ingest verwenden.
Vollständige Dokumentation, Konfigurationsoptionen, Batchaufnahme- und Protokollpufferbeispiele finden Sie im Java SDK-Repository.
Go Software Development Kit (SDK)
Go 1.21 oder höher ist erforderlich. Das SDK bietet einen hohen Durchsatz und eine hohe Leistung für die Streaming-Aufnahme. Es unterstützt JSON (einfachste) und Protokollpuffer (empfohlen für die Produktion).
go get github.com/databricks/zerobus-sdk/go@latest
JSON-Beispiel:
Aus Gründen der Einfachheit werden hier Fehler ignoriert. Überprüfen Sie im Produktionscode immer Fehler.
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")
}
Protokollpuffer: Verwenden Sie für typsicheres Einlesen Protokollpuffer mit RecordTypeProto (Standard) und stellen Sie descriptorProto in den Tabelleneigenschaften bereit. Erstellen Sie eine .proto-Datei, die Ihrem Tabellenschema entspricht, und führen Sie das Skript generate_proto aus, um die Dateien in Ihr Projekt zu importieren.
Arrow Flight: Für die spaltenorientierte oder batchorientierte Ingestion von Apache Arrow-RecordBatchDaten über dieselbe gRPC-Verbindung siehe Arrow Flight mit Zerobus Ingest verwenden.
Vollständige Dokumentation, Konfigurationsoptionen, Batchaufnahme, generate_proto Tool- und Protokollpufferbeispiele finden Sie im Go SDK-Repository.
C++-SDK
Von Bedeutung
Das C++-SDK befindet sich in der Betaversion.
C++17 oder höher ist erforderlich. Das SDK bietet natives gRPC-Streaming, OAuth und automatische Wiederherstellung über eine RAII C++-Schnittstelle. Es unterstützt JSON für einfache Setups und Protokollpuffer für Produktionsworkloads.
Das SDK wird als vorgefertigtes, plattformbasiertes Release-Bundle ausgeliefert, sodass Sie keine Rust-Toolkette benötigen, um es zu verwenden. Laden Sie das Bundle für Ihre Plattform (macOS, Linux einschließlich Musl oder Windows) von der Veröffentlichungsseite herunter, extrahieren Sie es, und zeigen Sie dann CMake im gebündelten FFI-Archiv. Das Archiv heißt unter macOS und Linux libzerobus_ffi.a und unter Windows zerobus_ffi.lib:
# 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
Zeigen Sie stattdessen auf Windows (PowerShell) auf das .lib Archiv:
cmake -S cpp -B build `
-DZEROBUS_FFI_LIBRARY="$PWD/lib/zerobus_ffi.lib" `
-DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j
Um das SDK in Ihrem eigenen CMake-Projekt aus ausgechecktem Quellcode zu erstellen, fügen Sie es als Unterverzeichnis hinzu und verlinken Sie das Ziel. Sie können auch FetchContent verwenden, um es während der Konfiguration abzurufen. Dies erstellt die FFI aus dem Rust-Quellcode, daher wird eine Rust-Toolchain benötigt:
add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)
Um stattdessen über add_subdirectory ein vorgefertigtes Paket zu verwenden, legen Sie zuerst die FFI-Pfade fest, damit CMake mit dem gebündelten Archiv verlinkt, anstatt zu versuchen, es aus nicht vorhandenem Rust-Quellcode zu erstellen. Verwendung von zerobus_ffi.lib unter 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)
JSON-Beispiel:
Die Erfassung erfolgt asynchron und über eine Pipeline. Die Methoden ingest_* reihen einen Datensatz in die Warteschlange ein und kehren sofort zurück. Stellen Sie den Batch in die Warteschlange und rufen Sie flush() einmal auf, anstatt nach jedem Datensatz zu warten.
#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;
}
Bei jedem Fehler wird zerobus::ZerobusException ausgelöst, das eine Nachricht und ein is_retryable()-Flag enthält. Um die Zuverlässigkeit in einem kontinuierlichen Datenstrom ohne Blockierung zu verfolgen, registrieren Sie ein AckCallback über StreamOptions::ack_callback. Callbacks werden serialisiert in einem Hintergrund-Thread ausgeführt und müssen noexcept sein. Siehe die C++-SDK-Dokumentation für den vollständigen Vertrag zu Threading, Drain-Richtlinie und Lebensdauer.
Für die typsichere Aufnahme können Sie Protokollpuffer auf eine von zwei Arten verwenden:
Generieren Sie das Schema aus dem Unity-Katalog mit
ProtoSchema::from_uc_json(). Dies erstellt direkt aus den Metadaten der Tabelle einen Deskriptor und einen JSON-zu-Proto-Encoder, sodass keine.proto-Datei und keinprotocerforderlich sind:- Rufen Sie das Metadaten-JSON der Tabelle aus der API Tabelle abrufen ab (
GET /api/2.1/unity-catalog/tables/{full_name}). Der Dienstprinzipal benötigtSELECTfür die Tabelle. - Übergeben Sie die Metadaten an
ProtoSchema::from_uc_json(), um den Deskriptor und den Encoder zu erstellen. - Legen Sie
TableProperties::descriptor_protofest und führen Sie anschließend die Erfassung mitingest_proto_records()durch.
- Rufen Sie das Metadaten-JSON der Tabelle aus der API Tabelle abrufen ab (
Kompilieren Sie ein eingechecktes
.protomitprotocfür die Typisierung bei der Kompilierung.
Eine ausführbare Schritt-für-Schritt-Anleitung finden Sie in den Protocol Buffers-Beispielen.
Informationen zur spaltenorientierten oder batchorientierten Erfassung von Apache-Arrow-Datensatz-Batches über dieselbe gRPC-Verbindung finden Sie unter Verwenden von Arrow Flight mit Zerobus Ingest.
Vollständige Dokumentation, Konfigurationsoptionen, Batchaufnahme- und Protokollpufferbeispiele finden Sie im C++-SDK-Repository.
C# SDK
Von Bedeutung
Das C# / .NET SDK befindet sich in der Beta-Phase. Das Paket Databricks.Zerobus ist eine Vorabversion.
.NET 8.0 oder höher ist erforderlich. Das SDK bietet natives gRPC-Streaming, OAuth und automatische Wiederherstellung. Es unterstützt JSON für einfache Setups und Protokollpuffer für Produktionsworkloads. Arrow Flight ist im C# SDK nicht verfügbar.
Fügen Sie dem Projekt das Databricks.Zerobus-Paket hinzu:
dotnet add package Databricks.Zerobus
JSON-Beispiel:
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 gibt den Offset des Datensatzes zurück und WaitForOffset blockiert, bis dieser Datensatz dauerhaft ist. Um einen Batch einzulesen, verwenden Sie IngestRecords, das ein Array von Datensätzen entgegennimmt und den letzten Offset zurückgibt. Das Blocken im Offset ist optional. Siehe Nachrichtensperre und Bestätigung.
Protocol Buffers: Für eine typsichere Ingestion erstellen Sie einen Stream mit sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), wobei descriptorProto die serialisierten DescriptorProto-Bytes für Ihre kompilierte Nachricht sind, und lesen Sie sie dann mit stream.IngestRecord(protoBytes) ein.
Für vollständige Dokumentation, Konfigurationsoptionen und Protokollpuffer-Beispiele siehe das C# SDK-Repository.
TypeScript SDK
Node.js 16 oder höher ist erforderlich. Das SDK bietet hohe Leistung mit asynchroner Unterstützung durch JavaScript Promises. Es unterstützt JSON (einfachste) und Protokollpuffer (empfohlen für die Produktion).
npm install @databricks/zerobus-ingest-sdk
JSON-Beispiel:
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();
}
Protokollpuffer: Verwenden Sie für typsicheres Einlesen Protokollpuffer mit RecordType.Proto (Standard) und stellen Sie descriptorProto in den Tabelleneigenschaften bereit.
Arrow Flight: Informationen zur spaltenorientierten oder batchorientierten Ingestion von Apache Arrow-RecordBatchDaten über dieselbe gRPC-Verbindung finden Sie unter Arrow Flight mit Zerobus Ingest verwenden.
Vollständige Dokumentation, Konfigurationsoptionen, Batchaufnahme- und Protokollpufferbeispiele finden Sie im TypeScript SDK-Repository.
REST API
Mit der REST-API können Sie einen einzelnen Datensatz aufnehmen, indem Sie eine HTTP POST-Anforderung an den /zerobus/v1/tables/<table-name>/insert Endpunkt senden. Der Datensatz selbst ist im Anforderungstext enthalten und muss im JSON-Format vorliegen.
Dieses Beispiel führt Sie durch die Verwendung von CURL, um Daten mithilfe der REST-API an Zerobus Ingest zu übertragen.
Headers
Für die Anforderung sind zwei bestimmte HTTP-Header erforderlich, um die Anforderung korrekt zu authentifizieren und zu formatieren.
-
Content-Type: application/json
- Obligatorisches Feld zum Angeben des Inhaltstyps. Derzeit ist JSON das einzige unterstützte Nachrichtenformat.
-
Autorisierung: Bearertoken <>
- Ersetzen Sie <das Token durch das OAuth-Token> , das Sie mithilfe des später bereitgestellten Curl-Befehls abgerufen haben.
OAuth-Token abrufen: Diese Token laufen jede Stunde ab und müssen aktualisiert werden. Sie können sie aktualisieren, indem Sie das OAuth-Token erneut abrufen.
Geben Sie die folgenden Parameter ein:
-
$CATALOG, ,$SCHEMA$TABLE, ,$WORKSPACE_ID$WORKSPACE_URL -
$DATABRICKS_CLIENT_IDund$DATABRICKS_CLIENT_SECRET- Diese beiden Parameter entsprechen dem von Ihnen erstellten Dienstprinzip.
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')
Datensatzerfassung:
Geben Sie die folgenden Parameter ein:
$ZEROBUS_ENDPOINT- Wie im Abschnitt „Get your workspace URL und Zerobus Ingest Endpoint“ definiert.
-
$CATALOG, ,$SCHEMA$TABLE, ,$WORKSPACE_ID$WORKSPACE_URL $OAUTH_TOKEN- Dies wurde im vorherigen Schritt erstellt.
Der Anforderungstext muss eine Liste von JSON-Objekten sein.
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 }]'
Wenn alle Informationen richtig ausgefüllt sind, sollten Sie eine leere JSON-Antwort mit einem HTTP-Statuscode von 200 erhalten.
Fehler behandeln
Die obigen Beispiele zeigen den glücklichen Weg. Im Produktivbetrieb sollte die Datenerfassung mit einer Fehlerbehandlung versehen werden. Das SDK überprüft vorübergehende Fehler wie Netzwerkprobleme automatisch über seine integrierte Wiederherstellung. Fehler, von denen es sich nicht erholen kann, wie ungültige Zugangsdaten oder eine fehlende Tabelle, werden als ZerobusException angezeigt:
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.
...
Die SDKs erholen sich außerdem automatisch von vorübergehenden Ausfällen und ermöglichen es Ihnen, nicht bestätigte Datensätze zu retten, wenn ein Stream dauerhaft ausfällt. Zu Mustern für resiliente Clients und der vollständigen Fehlerreferenz siehe Muster für Wiederherstellung und Wiederholungsversuche sowie Fehlerbehandlung für Zerobus Ingest.