Usar o Zerobus Ingest
Esta página descreve como ingerir dados usando o Zerobus Ingest no Lakeflow Connect.
Começar a usar o Zerobus Ingest
Se você tiver um firewall no lado do cliente, adicione o endereço IP usado pelo Zerobus Ingest à sua lista de permissões. Para visualizar endereços IP por região, consulte Endereços IP e domínios para serviços e ativos do Databricks.
Antes de começar, confirme se o Zerobus Ingest está disponível na região do seu workspace. Consulte Disponibilidade de ingestão.
- Obtenha uma URL do Zerobus Ingest.
- Crie ou identifique a tabela na qual você deseja ingerir dados.
- Crie um Service Principal e conceda privilégios à tabela.
- Conecte um cliente ou exportador para começar a enviar dados.
Escolha o guia para o seu caso de uso:
-
Ingira seus próprios dados : use os SDKs do Zerobus Ingest ou a API REST com um esquema definido por você. Siga as instruções nesta página.
-
Ingerir dados do OpenTelemetry : use SDKs ou coletores padrão do OpenTelemetry para enviar rastreamentos, logs e métricas para esquemas de tabela predefinidos. Para obter instruções completas, consulte Ingerir dados do OpenTelemetry com o Zerobus Ingest.
Escolha uma interface
O Zerobus Ingest oferece suporte a várias interfaces, todas gravando diretamente em tabelas Delta do Unity Catalog. Em resumo:
-
SDKs sobre gRPC : throughput sustentado mais alto, ideal para produtores de transmissão de alto volume.
-
REST : stateless, ideal para grandes frotas de dispositivos de ponta leves ou "tagarelas".
-
OpenTelemetry (OTLP) : para sistemas que já emitem rastreamentos, logs e métricas do OpenTelemetry. Veja Ingest OpenTelemetry data with Zerobus Ingest.
-
APIs compatíveis com Kafka (Beta): para produtores que já utilizam o protocolo Kafka. Consulte Usar APIs compatíveis com Kafka com o Zerobus Ingest.
Para uma comparação completa e como escolher, consulte protocolos de API. Por meio dos SDKs, você também pode escolher um formato de registro (JSON, Protocol Buffers (protobuf) ou Apache Arrow). Consulte Tipos de mensagem. O restante desta página usa os SDKs e a API REST.
Obtenha o URL do seu workspace e endpointde ingestão do Zerobus.
O URL do seu workspace aparece no navegador quando você log in. Embora o URL completo siga o formato https://<databricks-instance>.com/o=XXXXX, o URL workspace consiste em tudo o que vem antes do /o=XXXXX. Por exemplo, dado o seguinte URL completo, você pode determinar o URL workspace e o ID workspace .
- URL completa:
https://abcd-teste2-test-spcse2.cloud.databricks.com/?o=2281745829657864# - URL do espaço de trabalho:
https://abcd-teste2-test-spcse2.cloud.databricks.com - ID do espaço de trabalho:
2281745829657864
O endpoint do servidor depende do workspace e da região:
- endpoint do servidor:
<workspace-id>.zerobus.<region>.cloud.databricks.com
Para encontrar a região do seu workspace, abra o alternador de workspace na barra de navegação superior da interface do Databricks. A região é exibida abaixo de cada nome de workspace (por exemplo, us-west-2). Você também pode encontrá-lo no console da conta em Workspaces .
Para disponibilidade de região, consulte cotas do Zerobus Ingest.
Criar ou identificar a tabela de destino
Identifique a tabela de destino na qual você deseja inserir os dados. Para criar uma nova tabela de destino, execute o comando SQL CREATE TABLE . Por exemplo, crie uma nova tabela chamada unity.default.air_quality.
CREATE TABLE unity.default.air_quality (
device_name STRING, temp INT, humidity LONG);
O Zerobus Ingest pode gravar tanto em tabelas Delta gerenciadas quanto em tabelas de transmissão, que funcionam da mesma maneira, com os mesmos limites e cotas.
Para a ingestão de dados do OpenTelemetry, as tabelas devem usar esquemas predefinidos para cada tipo de sinal (rastreamentos, logs, métricas). Consulte Criar tabelas de destino no Unity Catalog.
O esquema da sua tabela é o contrato para o que o Zerobus Ingest aceita, e o Zerobus Ingest nunca o evolui automaticamente. Planeje as alterações de esquema proativamente: evolua a tabela primeiro e, em seguida, atualize os produtores. O Zerobus Ingest grava registros que não se encaixam mais após uma alteração de tabela disruptiva em um local de fallback durável, em vez de descartá-los. Consulte Gerenciamento de esquema e Recuperação de dados do local de fallback durável.
Por default, o Zerobus Ingest rejeita registros com campos que não correspondem ao esquema da tabela de destino. Para capturar esses campos em vez de perdê-los, configure uma coluna de dados resgatados. Consulte coluna de dados resgatados do Zerobus.
Crie uma entidade de serviço e conceda permissões.
Um Service Principal é uma identidade especializada que oferece mais segurança do que contas personalizadas. Para obter mais informações sobre entidades de serviço e como usá-las para autenticação, consulte Autorizar o acesso de entidades de serviço ao Databricks com OAuth.
Você pode criar e gerenciar Service Principal programaticamente com a API REST ou SDKs do Databricks, ou por meio da interface do usuário do Workspace, conforme descrito abaixo. As concessões de permissão no final desta seção são comandos SQL que você pode executar a partir de qualquer cliente.
-
Para criar um Service Principal, vá para Settings > Identity and Access .
-
Em entidade de serviço , selecione gerenciar .
-
Clique em Adicionar entidade de serviço .
-
Na janela Adicionar entidade de serviço , crie uma nova entidade de serviço clicando em Adicionar nova .
-
Gere e salve o ID do cliente e o segredo do cliente para a entidade de serviço.
-
Conceda as permissões necessárias para o catálogo, o esquema e a tabela à entidade de serviço.
- Na página Service principal , vá para a guia tab .
- Copie o ID do aplicativo (UUID).
- Utilize o seguinte SQL para conceder permissões, substituindo o UUID de exemplo e os nomes do catálogo, do esquema e das tabelas, se necessário.
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>`;
Escreva para um cliente
Use um Zerobus SDK na sua linguagem de programação preferida ou a API REST para ingerir dados na sua tabela de destino. Os SDKs são de código aberto. Para a biblioteca completa, documentação específica de linguagem e exemplos adicionais, consulte o repository do Zerobus SDK.
Os exemplos abaixo usam ingest_record_offset, que preserva a ordem na qual você envia os registros.
- Python SDK
- Rust SDK
- Java SDK
- Go SDK
- C++ SDK
- C# SDK
- TypeScript SDK
- REST API
É necessário Python 3.9 ou superior. O SDK oferece throughput elevado e E/S de rede eficiente por meio de um runtime assíncrono. Ele oferece suporte a JSON (mais simples) e Protocol Buffers (recomendado para produção). O SDK também oferece suporte a implementações síncronas e assíncronas, bem como aos métodos de ingestão baseados em offset e baseados em future.
pip install databricks-zerobus-ingest-sdk
Exemplo de 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()
Os exemplos acima usam o método ingest_record_offset baseado em offset sem aguardar o offset retornado. Para saber mais sobre os métodos de ingestão disponíveis, quando aguardar a confirmação de durabilidade em um offset e como acompanhar o progresso com um callback de confirmação, consulte Bloqueio e confirmação de mensagens.
Buffers de protocolo: para ingestão com segurança de tipo, passe um descritor protobuf para TableProperties (o formato é selecionado automaticamente). Gere um esquema a partir de sua tabela usando a ferramenta generate_proto, compile-o com protoc e, em seguida, passe o descritor compilado para criar a transmissão.
Arrow Flight (Beta): Para ingestão colunar ou orientada a lotes de dados Apache Arrow RecordBatch na mesma conexão gRPC, consulte Usar Arrow Flight com Zerobus Ingest. Requer o [arrow] extra: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.
Para obter documentação completa, opções de configuração, ingestão de lotes e exemplos do Protocol Buffer, consulte o repositório SDK Python.
É necessária a versão 1.70 ou superior do Rust. O SDK utiliza I/O assíncrono e gRPC para ingestão de alto throughput. Ele oferece suporte a JSON (mais simples) e Protocol Buffers (recomendado para produção).
Primeiro, importe o pacote.
cargo add databricks-zerobus-ingest-sdk
Ou adicione-o ao seu Cargo.toml.
[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication
Exemplo de 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(())
}
Buffers de protocolo: Para ingestão com segurança de tipo, use Buffers de protocolo por meio de .compiled_proto(descriptor) no construtor de transmissão em vez de .json(), em que descriptor é um prost_types::DescriptorProto. Gere os arquivos necessários usando a ferramenta generate_proto e importe para seu projeto.
Arrow Flight (Beta): Para ingestão colunar ou orientada a lotes de dados RecordBatch do Apache Arrow pela mesma conexão gRPC, consulte Usar o Arrow Flight com o Zerobus Ingest. Ative com o recurso Cargo: cargo add databricks-zerobus-ingest-sdk --features arrow-flight.
Para obter documentação completa, opções de configuração, ingestão de lotes, ferramenta generate_proto e exemplos de Protocol Buffer, consulte o repositório SDK Rust.
É necessário Java 8 ou superior. O SDK oferece baixa latência e E/S de rede eficiente para ingestão de alto throughput. Ele oferece suporte a JSON (mais simples) e Protocol Buffers (recomendado para produção).
Maven:
<dependency>
<groupId>com.databricks</groupId>
<artifactId>zerobus-ingest-sdk</artifactId>
<version>0.2.0</version>
</dependency>
Exemplo de 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();
}
}
}
Buffers de protocolo: Para ingestão com segurança de tipo, crie um ZerobusProtoStream com streamBuilder() e .compiledProto(...). Gere um esquema a partir da sua tabela usando a ferramenta JAR incluída e, em seguida, compile-o com protoc.
Arrow Flight (Beta): Para ingestão colunar ou orientada a lotes de dados Apache Arrow RecordBatch na mesma conexão gRPC, consulte Usar Arrow Flight com Zerobus Ingest.
Para obter documentação completa, opções de configuração, ingestão de lotes e exemplos do Protocol Buffer, consulte o repositório SDK Java.
É necessário Go 1.21 ou superior. O SDK oferece alto throughput e desempenho para ingestão de transmissão. Ele oferece suporte a JSON (mais simples) e Protocol Buffers (recomendado para produção).
go get github.com/databricks/zerobus-sdk/go@latest
Exemplo de JSON:
Para simplificar, os erros serão ignorados aqui. Em código de produção, sempre verifique os erros.
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: Para ingestão segura de tipos, use Protocol Buffers com RecordTypeProto (default) e forneça um descriptorProto nas propriedades da tabela. Crie um arquivo .proto arquivo correspondente ao esquema da sua tabela e script de execução generate_proto para ajudar você a importar os arquivos para o seu projeto.
Arrow Flight (Beta): Para ingestão colunar ou orientada a lotes de dados Apache Arrow RecordBatch na mesma conexão gRPC, consulte Usar Arrow Flight com Zerobus Ingest.
Para obter documentação completa, opções de configuração, ingestão de lotes, ferramenta generate_proto e exemplos do Protocol Buffer, consulte o repositório SDK Go.
Beta
O SDK C++ está em Beta.
É necessário C++17 ou superior. O SDK oferece transmissão gRPC nativa, OAuth e recuperação automática por meio de uma interface C++ RAII. Ele oferece suporte a JSON para configurações simples e Buffers de protocolo para cargas de trabalho de produção.
O SDK é fornecido como um pacote de lançamento pré-compilado por plataforma, portanto, você não precisa de uma cadeia de ferramentas Rust para usá-lo. Faça o download do pacote para sua plataforma (macOS, Linux incluindo musl ou Windows) na página de lançamentos, extraia-o e, em seguida, aponte o CMake para o arquivo compactado FFI do pacote. O arquivo compactado é nomeado libzerobus_ffi.a no macOS e Linux e zerobus_ffi.lib no 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
No Windows (PowerShell), aponte para o arquivo .lib em vez disso:
cmake -S cpp -B build `
-DZEROBUS_FFI_LIBRARY="$PWD/lib/zerobus_ffi.lib" `
-DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j
Para criar o SDK a partir de um checkout de origem em seu próprio projeto CMake, adicione-o como um subdiretório e vincule o destino. Você também pode usar FetchContent para buscá-lo no momento da configuração. Isso cria a FFI a partir da origem Rust, portanto, requer uma cadeia de ferramentas Rust:
add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)
Para consumir um pacote pré-compilado através de add_subdirectory, defina primeiro os caminhos FFI para que o CMake vincule o arquivo compactado do pacote em vez de tentar compilá-lo a partir do código-fonte Rust que não está presente. Use zerobus_ffi.lib no 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)
Exemplo de JSON:
A ingestão é assíncrona e em pipeline. Os métodos ingest_* enfileiram um registro e retornam imediatamente. Enfileire o lote e chame flush() uma única vez, em vez de aguardar após 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.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;
}
Cada falha gera zerobus::ZerobusException, que contém uma mensagem e um sinalizador is_retryable(). Para rastrear a durabilidade em uma transmissão contínua sem bloqueio, registre um AckCallback via StreamOptions::ack_callback. Os callbacks são executados serializados em um thread em segundo plano e devem ser noexcept. Consulte a documentação do SDK C++ para obter o contrato completo de threading, política de drenagem e tempo de vida.
Para ingestão com segurança de tipo, você pode usar Protocol Buffers de uma das duas maneiras:
-
Gere o esquema a partir do Unity Catalog com
ProtoSchema::from_uc_json(). Isso cria um descritor e um codificador de JSON para proto diretamente dos metadados da tabela, portanto, não precisa de nenhum arquivo.protoouprotoc:- Busque o JSON de metadados da tabela na API Get a table (
GET /api/2.1/unity-catalog/tables/{full_name}). O service principal precisa deSELECTna tabela. - Passe os metadados para
ProtoSchema::from_uc_json()para criar o descritor e o codificador. - Defina
TableProperties::descriptor_proto, depois ingira comingest_proto_records().
- Busque o JSON de metadados da tabela na API Get a table (
-
Compile um
.protoverificado comprotocpara tipagem em tempo de compilação.
Para um passo a passo executável, veja os exemplos de Protocol Buffers.
Para ingestão colunar ou orientada a lotes de lotes de registros do Apache Arrow pela mesma conexão gRPC, consulte Use Arrow Flight with Zerobus Ingest.
Para documentação completa, opções de configuração, ingestão em lotes e exemplos de Protocol Buffer, consulte o repository do SDK C++.
Beta
O SDK para C# / .NET está em Beta. O pacote Databricks.Zerobus é uma versão de pré-lançamento.
.NET 8.0 ou superior é necessário. O SDK oferece transmissão gRPC nativa, OAuth e recuperação automática. Ele oferece suporte a JSON para configurações simples e Buffers de protocolo para cargas de trabalho de produção. O Arrow Flight não está disponível no SDK C#.
Adicione o pacote Databricks.Zerobus ao seu projeto:
dotnet add package Databricks.Zerobus
Exemplo de 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 retorna o deslocamento do registro, e WaitForOffset bloqueia até que esse registro seja durável. Para ingerir um lote, use IngestRecords, que recebe uma matriz de registros e retorna o último deslocamento. O bloqueio no deslocamento é opcional. Consulte Bloqueio e reconhecimento de mensagens.
Protocol Buffers: Para ingestão com segurança de tipo, crie uma transmissão com sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), em que descriptorProto são os DescriptorProto bytes serializados para sua mensagem compilada, e então ingira com stream.IngestRecord(protoBytes).
Para documentação completa, opções de configuração e exemplos de Protocol Buffer, consulte o repository do SDK C#.
É necessário Node.js 16 ou superior. O SDK oferece alto desempenho com suporte assíncrono por meio de Promises do JavaScript. Ele oferece suporte a JSON (mais simples) e Buffers de protocolo (recomendado para produção).
npm install @databricks/zerobus-ingest-sdk
Exemplo de 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: Para ingestão segura de tipos, use Protocol Buffers com RecordType.Proto (default) e forneça um descriptorProto nas propriedades da tabela.
Arrow Flight (Beta): Para ingestão colunar ou orientada a lotes de dados Apache Arrow RecordBatch na mesma conexão gRPC, consulte Usar Arrow Flight com Zerobus Ingest.
Para obter documentação completa, opções de configuração, ingestão de lotes e exemplos do Protocol Buffer, consulte o repositório SDK do TypeScript.
A API REST permite que você insira um único registro enviando uma solicitação HTTP POST para o endpoint /zerobus/v1/tables/<table-name>/insert . O próprio registro está incluído no corpo da solicitação e deve estar no formato JSON.
Este exemplo mostra como usar o CURL para enviar dados para o Zerobus Ingest usando a API REST.
Cabeçalhos
A solicitação requer dois cabeçalhos HTTP específicos para autenticar e formatar a solicitação corretamente.
-
Content-Type: application/json- Campo obrigatório para especificar o tipo de conteúdo. Atualmente, JSON é o único formato de mensagem suportado.
-
Authorization: Bearer <token>- Substitua
<token>pelos tokens OAuth que você obteve usando o comando curl fornecido posteriormente.
- Substitua
Obter tokens OAuth : Esses tokens expiram a cada hora e precisam ser renovados. Você pode refresh -los buscando novamente os tokens OAuth .
Preencha os seguintes parâmetros:
$CATALOG,$SCHEMA,$TABLE,$WORKSPACE_ID,$WORKSPACE_URL$DATABRICKS_CLIENT_IDe$DATABRICKS_CLIENT_SECRET- Esses dois parâmetros correspondem ao princípio de serviço que você criou.
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')
Registro de ingestão:
Preencha os seguintes parâmetros:
-
$ZEROBUS_ENDPOINT- Conforme definido na seção Obtenha o URL do seu workspace e o endpoint de ingestão do Zerobus .
-
$CATALOG,$SCHEMA,$TABLE,$WORKSPACE_ID,$WORKSPACE_URL -
$OAUTH_TOKEN- Isso foi criado no passo anterior.
O corpo da requisição deve ser uma 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 }]'
Se todas as informações forem preenchidas corretamente, você deverá receber uma resposta JSON vazia com um código de status HTTP 200.
Tratamento de erros
Os exemplos acima mostram o caminho feliz. Em produção, envolva a ingestão em tratamento de erros. O SDK tenta novamente erros transitórios, como problemas de rede, automaticamente por meio de sua recuperação integrada. Falhas das quais não é possível se recuperar, como credenciais inválidas ou uma tabela ausente, aparecem 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.
...
Os SDKs também se recuperam de falhas transitórias automaticamente e permitem que você resgate registros não confirmados quando uma transmissão falha permanentemente. Para padrões de cliente resilientes e a referência completa de erros, consulte Padrões de recuperação e repetição e Tratamento de erros do Zerobus Ingest.