Observação
O acesso a essa página exige autorização. Você pode tentar entrar ou alterar diretórios.
O acesso a essa página exige autorização. Você pode tentar alterar os diretórios.
Esta página descreve como ingerir dados usando o Zerobus Ingest no Lakeflow Connect.
Comece com o Zerobus Ingest
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 de Ingestão do Zerobus.
- Crie ou identifique a tabela na qual você deseja ingerir dados.
- Crie um principal de serviço e conceda privilégios à tabela.
- Conecte um cliente ou exportador para começar a enviar dados.
Escolha o guia para seu caso de uso:
Ingerir seus próprios dados: use os SDKs de Ingestão do Zerobus ou a API REST com um esquema que você define. Siga as instruções nesta página.
Ingerir dados OpenTelemetry: use SDKs ou coletores OpenTelemetry padrão para enviar rastreamentos, logs e métricas em esquemas de tabela predefinidos. Para obter instruções completas, consulte Ingestão de dados OpenTelemetry com Zerobus Ingest.
Escolher uma interface
O Zerobus Ingest suporta várias interfaces, todas as quais gravam diretamente em tabelas Delta do Unity Catalog. Resumindo:
- SDKs via gRPC: maior taxa de transferência sustentada, melhor para produtores de streaming com alto volume.
- REST: sem estado, ideal para grandes frotas de dispositivos de borda leves ou que geram muitas comunicações.
- OpenTelemetry (OTLP): para sistemas que já emitem rastros, logs e métricas do OpenTelemetry. Veja Ingerir dados do OpenTelemetry com o Zerobus Ingest.
- APIs compatíveis com Kafka (Beta): para produtores que já falam o protocolo Kafka. Veja Usar APIs compatíveis com Kafka com o Zerobus Ingest.
Para uma comparação completa e como escolher, veja protocolos API. Entre os SDKs, você também pode escolher um formato de registro (JSON, Protocol Buffers (protobuf) ou Apache Arrow). Veja Tipos de mensagem. O restante desta página usa os SDKs e a API REST.
Obter o ponto de extremidade do Zerobus Ingest e o URL do workspace
A URL do workspace é exibida no navegador quando você faz logon. Embora a URL completa siga o formato https://<databricks-instance>.net/o=XXXXX, a URL do workspace consiste em tudo antes do /o=XXXXX. Por exemplo, considerando-se o URL completo a seguir, você pode determinar o ID e o URL do workspace.
- URL completa:
https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864# - URL do workspace:
https://abcd-teste2-test-spcse2.azuredatabricks.net - ID do workspace:
2281745829657864
O endpoint do servidor depende do workspace e da região.
- Ponto de extremidade do servidor:
<workspace-id>.zerobus.<region>.azuredatabricks.net
Para encontrar a região do seu workspace, abra o seletor de workspace na barra de navegação superior da interface do usuário do Databricks. A região é exibida abaixo de cada nome de espaço de trabalho (por exemplo, eastus). Você também pode encontrá-lo no console da conta em Workspaces.
Para ver a disponibilidade por região, consulte as cotas de ingestão do Zerobus.
Criar ou identificar a tabela de destino
Identifique a tabela de destino na qual você deseja ingerir dados. Para criar uma nova tabela de destino, execute o CREATE TABLE comando SQL. 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 escrever tanto em tabelas Delta gerenciadas quanto em tabelas de streaming, que funcionam da mesma forma, com os mesmos limites e cotas.
Observação
Para a ingestão de dados do OpenTelemetry, as tabelas precisam utilizar esquemas predefinidos para cada tipo de sinal (rastros, logs, métricas). Consulte Criar tabelas de destino no Unity Catalog.
O esquema da sua tabela é o contrato que define o que o Zerobus Ingest aceita, e o Zerobus Ingest nunca o evolui automaticamente. Planeje as alterações no esquema de forma proativa: evolua primeiro a tabela e, depois, atualize os produtores. O Zerobus Ingest grava em um local de fallback persistente os registros que deixam de ser compatíveis após uma alteração interruptiva na tabela, em vez de descartá-los. Consulte Gerenciamento de esquema e Recuperando dados do local de fallback persistente.
Por padrão, 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 resgate. Veja a coluna de resgate do Zerobus.
Criar um service principal e conceder permissões
Um serviço principal é uma identidade especializada que fornece mais segurança do que contas personalizadas. Para mais informações sobre entidades de serviço e como usá-las na autenticação, consulte Autorizar o acesso de entidade de serviço ao Azure Databricks com OAuth.
Você pode criar e gerenciar princípios de serviço programaticamente com a API ou SDKs REST do Azure Databricks, ou através da interface do workspace, conforme descrito abaixo. As permissões concedidas ao final desta seção são SQL que você pode rodar de qualquer cliente.
Para criar um principal de serviço, acesse Configurações>de Identidade e Acesso.
Em Entidades 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 o principal do serviço.
Conceda as permissões necessárias para o catálogo, o esquema e a tabela ao principal de serviço.
- Na página do principal do serviço , vá até a aba Configurações .
- Copie a ID do Aplicativo (UUID).
- Use o SQL a seguir para conceder permissões, substituindo o exemplo de UUID e catálogo, nome do esquema e nomes de tabela, se necessário.
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>`;
Gravar um cliente
Use um SDK do Zerobus em sua linguagem de programação preferida ou na API REST para ingerir dados em sua tabela de destino. Os SDKs são código aberto. Para a biblioteca completa, documentação específica de linguagem e exemplos adicionais, veja o repositório Zerobus SDK.
Os exemplos abaixo usam ingest_record_offset, que preserva a ordem em que você envia os registros.
SDK do Python
Python 3.9 ou superior é necessário. O SDK oferece alta taxa de transferência e E/S de rede eficiente por meio de um tempo de execução assíncrono. Ele dá suporte a JSON (mais simples) e Buffers de Protocolo (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 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.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()
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 de um offset e como acompanhar o progresso com um retorno de chamada de confirmação, consulte Bloqueio de mensagens e confirmação.
Protocol Buffers: Para ingestão com segurança de tipos, passe um descritor do protobuf para TableProperties (o formato é selecionado automaticamente). Gere um esquema a partir da sua tabela usando a generate_proto ferramenta, compile com protoc, e depois passe o descritor compilado para criar o fluxo.
Arrow Flight: Para ingestão colunar ou em lote de dados do Apache Arrow RecordBatch pela 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 em lote e exemplos de Protocol Buffer, consulte o repositório do SDK Python.
SDK Rust
O Rust 1.70 ou superior é necessário. O SDK usa E/S assíncrona e gRPC para ingestão de alto desempenho. Ele dá suporte a JSON (mais simples) e Buffers de Protocolo (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://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(())
}
Buffers de Protocolo: para ingestão com segurança de tipos, 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 generate_proto ferramenta e importe para seu projeto.
Arrow Flight: Para ingestão colunar ou em lote de dados Apache Arrow RecordBatch pela mesma conexão gRPC, veja Como usar o Arrow Flight com o Zerobus Ingest. Ative usando o recurso Cargo: cargo add databricks-zerobus-ingest-sdk --features arrow-flight.
Para obter documentação completa, opções de configuração, ingestão em lote, generate_proto ferramenta e exemplos de Buffer de Protocolo, consulte o repositório do Rust SDK.
SDK do Java
Java 8 ou superior é necessário. O SDK oferece baixa latência e operações de E/S de rede eficientes para ingestão com alta taxa de transferência. Ele dá suporte a JSON (mais simples) e Buffers de Protocolo (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.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();
}
}
}
Protocol Buffers: para realizar a ingestão com segurança de tipos, crie ZerobusProtoStream com streamBuilder() e .compiledProto(...). Gere um esquema da tabela usando a ferramenta JAR agrupada e compile-o com protoc.
Arrow Flight: Para ingestão de dados do Apache Arrow RecordBatch de forma colunar ou em lote pela mesma conexão gRPC, consulte Usar o Arrow Flight com o Zerobus Ingest.
Para obter a documentação completa, opções de configuração, ingestão em lote e exemplos de Buffer de Protocolo, consulte o repositório do SDK Java.
SDK do Go
O Go 1.21 ou posterior é necessário. O SDK oferece alta taxa de transferência e desempenho para ingestão de streaming. Ele dá suporte a JSON (mais simples) e Buffers de Protocolo (recomendado para produção).
go get github.com/databricks/zerobus-sdk/go@latest
Exemplo de JSON:
Para simplificar, os erros são ignorados aqui. No 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.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")
}
Buffers de Protocolo: para ingestão de tipo seguro, use Buffers de Protocolo com RecordTypeProto (padrão) e forneça um descriptorProto na tabela de propriedades. Crie um arquivo .proto que corresponda ao esquema de tabela e execute generate_proto o script para ajudá-lo a importar os arquivos para seu projeto.
Arrow Flight: Para ingestão colunar ou orientada a lotes de dados do Apache Arrow RecordBatch pela mesma conexão gRPC, consulte Como usar o Arrow Flight com a ingestão do Zerobus.
Para obter a documentação completa, opções de configuração, ingestão em lote, a ferramenta generate_proto e exemplos de Protocolo de Buffer, consulte o repositório do SDK Go.
C++ SDK
Importante
O SDK do C++ está em Beta.
C++17 ou superior é necessário. O SDK oferece streaming nativo gRPC, OAuth e recuperação automática por meio de uma interface RAII C++. Ele dá 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 distribuição pré-compilado específico para cada plataforma, portanto você não precisa de uma toolchain Rust para usá-lo. Baixe o pacote para sua plataforma (macOS, Linux incluindo musl ou Windows) da página de versões, extraia-o e, em seguida, aponte o CMake para o arquivo FFI empacotado. O arquivo é nomeado libzerobus_ffi.a no macOS e no 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 uma cópia do código-fonte no seu próprio projeto CMake, adicione-o como um subdiretório e vincule-o ao alvo. Você também pode usar FetchContent para buscá-lo durante a configuração. Isso compila a FFI a partir do código-fonte em Rust, portanto requer uma toolchain do Rust:
add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)
Para consumir um bundle pré-compilado por meio de add_subdirectory, defina primeiro os caminhos de FFI para que o CMake vincule o arquivo compactado, em vez de tentar criá-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 realizada em estágios. Os métodos ingest_* colocam um registro na fila e retornam imediatamente. Enfileire o lote e chame flush() uma vez, em vez de aguardar depois de cada registro.
#include "zerobus/zerobus.hpp"
#include <string>
#include <vector>
int main() {
// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const std::string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
const std::string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
const std::string TABLE_NAME = "main.default.air_quality";
const std::string CLIENT_ID = "your-client-id";
const std::string CLIENT_SECRET = "your-client-secret";
zerobus::Sdk sdk = zerobus::Sdk::builder()
.endpoint(SERVER_ENDPOINT)
.unity_catalog_url(DATABRICKS_WORKSPACE_URL)
.application_name("my-app")
.build();
zerobus::TableProperties table;
table.table_name = TABLE_NAME; // empty descriptor => JSON stream
zerobus::StreamOptions options;
options.record_type = zerobus::RecordType::Json;
zerobus::Stream stream =
sdk.create_stream(table, CLIENT_ID, CLIENT_SECRET, options);
std::vector<std::string> batch = {
R"({"device_name": "sensor-001", "temp": 20, "humidity": 60})",
R"({"device_name": "sensor-002", "temp": 22, "humidity": 55})",
};
stream.ingest_json_records(batch); // queue the batch — no per-record wait
stream.flush(); // wait once for all acks
stream.close();
return 0;
}
Toda falha lança zerobus::ZerobusException, que contém uma mensagem e um sinalizador is_retryable(). Para monitorar a durabilidade em um fluxo contínuo sem bloqueios, registre um AckCallback via StreamOptions::ack_callback. Os retornos de chamada são executados em série em uma thread em segundo plano e devem ser noexcept. Consulte a documentação do SDK de C++ para ver o modelo completo de threads, a política de drenagem e o contrato de ciclo de vida.
Para uma ingestão de tipo seguro, você pode usar Buffers de Protocolo de duas maneiras:
Gerar 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, por isso não precisa de arquivo.protonem deprotoc:- Busque o JSON de metadados da tabela na API Obter uma tabela (
GET /api/2.1/unity-catalog/tables/{full_name}). A entidade de serviço precisaSELECTna tabela. - Passe os metadados para
ProtoSchema::from_uc_json()criar o descritor e o codificador. - Configure
TableProperties::descriptor_protoe, em seguida, faça a ingestão comingest_proto_records().
- Busque o JSON de metadados da tabela na API Obter uma tabela (
Compile um
.protosubmetido comprotocpara tipagem em tempo de compilação.
Para um tutorial executável, consulte os exemplos de Protocol Buffers.
Para a ingestão orientada a lotes ou colunas de lotes de registro do Apache Arrow através da mesma conexão gRPC, consulte Usar Arrow Flight com Zerobus Ingest.
Para obter documentação completa, opções de configuração, ingestão em lote e exemplos de Buffer de Protocolo, consulte o repositório do SDK do C++.
C# SDK
Importante
O SDK C# / .NET está em Beta. O pacote Databricks.Zerobus está em versão prévia.
.NET 8.0 ou superior é necessário. O SDK oferece streaming nativo de gRPC, OAuth e recuperação automática. Ele dá suporte a JSON para configurações simples e buffers de protocolo para cargas de trabalho de produção. Arrow Flight não está disponível no SDK C#.
Adicione o pacote de 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.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 retorna o offset do registro e WaitForOffset bloqueia até que esse registro seja persistente. Para realizar a ingestão de um lote, use IngestRecords, que recebe uma matriz de registros e retorna o último offset. O bloqueio no offset é opcional. Veja Bloqueio e reconhecimento de mensagens.
Protocol Buffers: para realizar a ingestão com segurança de tipos, crie um fluxo com sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), em que descriptorProto são os bytes serializados de DescriptorProto para sua mensagem compilada. Em seguida, faça a ingestão com stream.IngestRecord(protoBytes).
Para documentação completa, opções de configuração e exemplos de Protocol Buffer, veja o repositório do SDK C#.
TypeScript SDK
Node.js 16 ou superior é necessário. O SDK oferece alto desempenho com suporte assíncrono por meio do JavaScript Promises. Ele dá 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.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();
}
Buffers de Protocolo: para ingestão de tipo seguro, use Buffers de Protocolo com RecordType.Proto (padrão) e forneça um descriptorProto na tabela de propriedades.
Arrow Flight: Para ingestão colunar ou orientada a lotes de dados em Apache Arrow RecordBatch na mesma conexão gRPC, consulte Use o Arrow Flight com o Zerobus Ingest.
Para obter documentação completa, opções de configuração, ingestão em lote e exemplos de Buffer de Protocolo, consulte o repositório do SDK do TypeScript.
API REST
A API REST permite que você ingera um único registro enviando uma solicitação HTTP POST para o /zerobus/v1/tables/<table-name>/insert ponto de extremidade. O registro em si está incluído no corpo da solicitação e deve estar no formato JSON.
Este exemplo explica como usar CURL para enviar dados para a 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.
-
Tipo de conteúdo: application/json
- Campo obrigatório para especificar o tipo de conteúdo. Atualmente, JSON é o único formato de mensagem com suporte.
-
Autorização: <token> de portador do ARM
- Substitua o <token> pelo token OAuth que você obteve usando o comando curl fornecido a seguir.
Buscar token OAuth: esses tokens expiram a cada hora e devem ser atualizados. Você pode atualizá-los buscando novamente o token 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')
Ingestão de Registros:
Preencha os seguintes parâmetros:
$ZEROBUS_ENDPOINT- Conforme definido na seção Obter o URL do seu workspace e o ponto de extremidade do Zerobus Ingest.
-
$CATALOG,$SCHEMA,$TABLE, ,$WORKSPACE_ID$WORKSPACE_URL $OAUTH_TOKEN- Isso foi criado na etapa anterior.
O corpo da solicitaçã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 de 200.
Lidar com erros
Os exemplos acima mostram o caminho feliz. Em produção, envolva a ingestão no tratamento de erros. O SDK faz novas tentativas automaticamente em caso de erros transitórios, como problemas de rede, por meio de seu mecanismo integrado de recuperação. 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 automaticamente de falhas transitórias e permitem que você resgate registros não confirmados quando um fluxo falha permanentemente. Para padrões para clientes resilientes e a referência completa de erros, consulte Padrões de recuperação e nova tentativa e Tratamento de erros do Zerobus Ingest.