Nota
O acesso a esta página requer autorização. Pode tentar iniciar sessão ou alterar os diretórios.
O acesso a esta página requer autorização. 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 espaço de trabalho. Ver Disponibilidade de Ingestão.
- Obtenha um URL de Zerobus para ingestão.
- Crie ou identifique a tabela na qual você deseja ingerir dados.
- Crie um principal de serviço e conceda privilégios à tabela.
- Ligue um cliente ou exportador para começar a enviar dados.
Escolha o guia para o seu caso de uso:
Ingere os seus próprios dados: Use os SDKs de Ingesta Zerobus ou a API REST com um esquema que defina. Siga as instruções nesta página.
Ingerir dados OpenTelemetry: Use SDKs ou coletores OpenTelemetry padrão para enviar traços, registos e métricas para esquemas de tabela pré-definidos. Para instruções completas, consulte Ingest OpenTelemetry data with Zerobus Ingest.
Escolha uma interface
O Zerobus Ingest é compatível com várias interfaces, todas escrevendo diretamente em tabelas Delta do Unity Catalog. Em resumo:
- SDKs com gRPC: maior taxa de transferência sustentada, ideal para produtores de streaming com elevado volume.
- REST: sem estado, ideal para grandes frotas de dispositivos periféricos ligeiros ou com muitas comunicações.
- OpenTelemetry (OTLP): para sistemas que já emitem rastreamentos, registos e métricas OpenTelemetry. Veja Ingerir dados OpenTelemetry com Zerobus Ingest.
- APIs compatíveis com Kafka (Beta): para produtores que já falam o protocolo Kafka. Veja Usar APIs compatíveis com Kafka com Zerobus Ingest.
Para uma comparação completa e como escolher, consulte protocolos API. Em vez dos SDKs, também pode escolher um formato de registo (JSON, Protocol Buffers (protobuf) ou Apache Arrow). Ver Tipos de mensagem. O resto desta página utiliza os SDKs e a API REST.
Obtenha o URL do seu espaço de trabalho e o endpoint Zerobus Ingest
O URL do seu espaço de trabalho aparece no navegador quando inicia sessão. Enquanto a URL completa segue o formato https://<databricks-instance>.net/o=XXXXX, a URL do espaço de trabalho consiste em tudo o que está antes do /o=XXXXX. Por exemplo, dado o seguinte URL completo, pode determinar o URL do espaço de trabalho e o ID do espaço de trabalho.
- URL completo:
https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864# - URL do espaço de trabalho:
https://abcd-teste2-test-spcse2.azuredatabricks.net - ID do Espaço de Trabalho:
2281745829657864
O endpoint do servidor depende do espaço de trabalho e da região:
- Endpoint do servidor:
<workspace-id>.zerobus.<region>.azuredatabricks.net
Para encontrar a sua região de espaço de trabalho, abra o comutador de espaço de trabalho na barra de navegação superior da interface do Databricks. A região é exibida abaixo de cada nome de espaço de trabalho (por exemplo, eastus). Também podes encontrá-lo na consola da conta, em Espaços de Trabalho.
Para ver as regiões disponíveis, consulte quotas do Zerobus Ingest.
Criar ou identificar a tabela de destino
Identifique a tabela alvo onde quer 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 geridas como em tabelas de streaming, que funcionam da mesma forma, com os mesmos limites e quotas.
Observação
Para a ingestão do OpenTelemetry, as tabelas devem usar esquemas pré-definidos para cada tipo de sinal (traços, logs, métricas). Veja Criar tabelas de alvo no Catálogo Unity.
O teu esquema de tabela é o contrato para o que o Zerobus Ingest aceita, e o Zerobus Ingest nunca o evolui automaticamente. Planeie antecipadamente as alterações ao esquema: evolua primeiro a tabela e, em seguida, atualize os produtores. O Zerobus Ingest escreve os registos que deixam de ser compatíveis após uma alteração incompatível da tabela num local de contingência persistente, em vez de os descartar. Consulte Gestão de esquemas e Recuperação de dados a partir do local de contingência persistente.
Por defeito, o Zerobus Ingest rejeita registos com campos que não correspondem ao esquema da tabela de destino. Para capturar esses campos em vez de os perder, configure uma coluna de resgate. Ver coluna de resgate do Zerobus.
Criar um principal de serviço e conceder permissões
Um principal de serviço é uma identidade especializada que oferece mais segurança do que contas personalizadas. Para mais informações sobre princípios de serviço e como os usar para autenticação, consulte Autorizar o acesso do principal de serviço ao Azure Databricks com OAuth.
Pode criar e gerir 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 no final desta secção são SQL que podes executar a partir de qualquer cliente.
Para criar um principal de serviço, vá a Definições>Identidade e Acesso.
Na secção de princípios de serviço, selecione Gerir.
Clique em Adicionar principal do serviço.
Na janela Adicionar principal de serviço , crie um novo principal de serviço clicando em Adicionar novo.
Gerar e guardar o ID do cliente e o segredo do cliente para o principal de serviço.
Conceda as permissões necessárias para o catálogo, o esquema e a tabela ao principal do serviço.
- Na página de Princípio do Serviço , vai ao separador Configurações .
- Copie a ID do aplicativo (UUID).
- Use o SQL seguinte para conceder permissões, substituindo o UUID de exemplo e o catálogo, nome do esquema e nomes das tabelas, 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>`;
Desenvolver um cliente
Use um SDK Zerobus na sua linguagem de programação preferida ou a API REST para ingerir dados na sua tabela de destino. Os SDKs são open source. Para a biblioteca completa, documentação específica da linguagem e exemplos adicionais, consulte o repositório Zerobus SDK.
Os exemplos abaixo usam ingest_record_offset, que preserva a ordem em que envia os registos.
Python SDK
Python 3.9 ou superior é obrigatório. O SDK proporciona alta taxa de transferência e E/S de rede eficiente através de um tempo de execução assíncrono. Suporta JSON (o mais simples) e Protocol Buffers (recomendados para produção). O SDK também suporta implementações síncronas e assíncronas, bem como métodos de ingestão baseados em offset e em futures.
pip install databricks-zerobus-ingest-sdk
Exemplo 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 utilizam o método ingest_record_offset baseado em offset sem aguardar o offset devolvido. Para conhecer os métodos de ingestão disponíveis, quando aguardar a confirmação de durabilidade de um offset e como acompanhar o progresso com uma função de callback de confirmação, consulte Bloqueio e confirmação de mensagens.
Protocol Buffers: Para uma ingestão com segurança de tipos, passe um descritor de protobuf para TableProperties (o formato é selecionado automaticamente). Gera um esquema a partir da tua tabela usando a generate_proto ferramenta, compila-o com protoc, e depois passa o descritor compilado para criar o fluxo.
Arrow Flight: Para a ingestão em colunas ou orientada para lotes de dados Apache Arrow RecordBatch através da mesma ligação gRPC, consulte Utilizar Arrow Flight com o Zerobus Ingest. Necessita do [arrow] extra: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.
Para documentação completa, opções de configuração, ingestão em lote e exemplos de Protocol Buffer, consulte o repositório SDK Python.
Rust SDK
É necessário Rust 1.70 ou superior. O SDK utiliza I/O assíncrona e gRPC para ingestão de alto rendimento. Suporta JSON (o mais simples) e Protocol Buffers (recomendados para produção).
Primeiro, importa o pacote.
cargo add databricks-zerobus-ingest-sdk
Ou adicioná-lo ao Cargo.toml.
[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication
Exemplo 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(())
}
Protocol Buffers: Para ingestão com segurança de tipos, use Protocol Buffers através de .compiled_proto(descriptor) no construtor de fluxo em vez de .json(), em que descriptor é um prost_types::DescriptorProto. Gera os ficheiros necessários usando a generate_proto ferramenta e importa-os para o teu projeto.
Arrow Flight: Para a ingestão em formato colunar ou orientada por lotes de dados Apache Arrow RecordBatch através da mesma ligação gRPC, consulte Utilizar o Arrow Flight com o Zerobus Ingest. Ative com a função Cargo: cargo add databricks-zerobus-ingest-sdk --features arrow-flight.
Para documentação completa, opções de configuração, ingestão em lote, generate_proto exemplos de ferramentas e Protocol Buffer, consulte o repositório Rust SDK.
SDK de Java
Java 8 ou superior é obrigatório. O SDK proporciona baixa latência e operações de entrada/saída de rede eficientes para ingestão de elevado débito. Suporta JSON (o mais simples) e Protocol Buffers (recomendados para produção).
Maven:
<dependency>
<groupId>com.databricks</groupId>
<artifactId>zerobus-ingest-sdk</artifactId>
<version>0.2.0</version>
</dependency>
Exemplo 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 ingestão com segurança de tipos, crie um ZerobusProtoStream com streamBuilder() e .compiledProto(...). Gera um esquema a partir da tua tabela usando a ferramenta JAR incluída, depois compila-o com protoc.
Arrow Flight: Para a ingestão colunar ou orientada para lotes de dados Apache Arrow RecordBatch através da mesma ligação gRPC, consulte Usar o Arrow Flight com o Zerobus Ingest.
Para documentação completa, opções de configuração, ingestão em lote e exemplos de Protocol Buffer, consulte o repositório SDK Java.
Kit de Desenvolvimento de Software Go
É necessário o Go 1.21 ou superior. O SDK proporciona alto débito e desempenho para ingestão de streaming. Suporta JSON (o mais simples) e Protocol Buffers (recomendados para produção).
go get github.com/databricks/zerobus-sdk/go@latest
Exemplo JSON:
Para simplificar, os erros são ignorados aqui. No código de produção, verifica sempre 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 segura por tipos, use Protocol Buffers com RecordTypeProto (por defeito) e forneça um descriptorProto nas propriedades da tabela. Cria um ficheiro .proto que corresponda ao esquema da tua tabela e executa generate_proto um script para te ajudar a importar os ficheiros para o teu projeto.
Arrow Flight: Para a ingestão em colunas ou em lote de dados Apache Arrow RecordBatch através da mesma ligação gRPC, consulte Utilizar o Arrow Flight com o Zerobus Ingest.
Para a documentação completa, opções de configuração, ingestão em lote, a ferramenta generate_proto e exemplos de Protocol Buffers, consulte o repositório Go SDK.
C++ SDK
Importante
O SDK em C++ está em Beta.
É necessário C++17 ou superior. O SDK fornece streaming nativo gRPC, OAuth e recuperação automática através de uma interface RAII C++. Suporta JSON para configurações simples e Protocol Buffers para cargas de trabalho de produção.
O SDK vem como um pacote pré-construído, lançado por plataforma, por isso não precisas de uma cadeia de ferramentas Rust para o usar. Descarregue o pacote para a sua plataforma (macOS, Linux incluindo musl, ou Windows) da página de lançamentos, extraia-o e depois aponte o CMake para o arquivo FFI incluído. O arquivo está 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), indique o ficheiro .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 compilar o SDK a partir de um checkout do código-fonte no seu próprio projeto CMake, adicione-o como subdiretório e associe o alvo. Também podes usar FetchContent para o obter durante a configuração. Isto constrói o FFI a partir da fonte Rust, por isso requer uma cadeia de ferramentas Rust:
add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)
Para consumir, em vez disso, um pacote pré-compilado através de add_subdirectory, defina primeiro os caminhos de FFI para que o CMake ligue ao arquivo incluído no pacote, em vez de tentar compilá-lo a partir do código-fonte em Rust, que não está presente. Utilização 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 JSON:
A ingestão é assíncrona e em pipeline. Os ingest_* métodos colocam um registo em fila e retornam imediatamente. Coloque o lote na fila e chame flush() uma vez, em vez de esperar após cada registo.
#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;
}
Cada falha lança zerobus::ZerobusException, que inclui uma mensagem e um sinalizador is_retryable(). Para monitorizar a durabilidade num fluxo contínuo, sem bloqueio, regista um AckCallback através de StreamOptions::ack_callback. Os callbacks executam-se serializados num thread em segundo plano e devem ser noexcept. Consulte a documentação do SDK de C++ para obter informações completas sobre o modelo de processamento por threads, a política de drenagem e o contrato de tempo de vida.
Para uma ingestão segura por tipos, pode usar Protocol Buffers de duas formas:
Gerar o esquema a partir do Unity Catalog com
ProtoSchema::from_uc_json(). Isto constrói um descritor e um codificador JSON-para-proto diretamente a partir dos metadados da tabela, pelo que não precisa de.protoficheiro nemprotoc:- Obtenha os metadados JSON da tabela a partir da API Get a table (
GET /api/2.1/unity-catalog/tables/{full_name}). O principal de serviço precisa deSELECTna mesa. - Passa os metadados para
ProtoSchema::from_uc_json()construir o descritor e o codificador. - Defina
TableProperties::descriptor_proto, e depois ingera comingest_proto_records().
- Obtenha os metadados JSON da tabela a partir da API Get a table (
Compilar um checked-in
.protocomprotocpara a tipagem em tempo de compilação.
Para um guia executável, consulte os exemplos de Protocol Buffers.
Para a ingestão colunar ou orientada por lotes de lotes de registos do Apache Arrow através da mesma ligação gRPC, veja Use Arrow Flight with Zerobus Ingest.
Para documentação completa, opções de configuração, ingestão em lote e exemplos de Protocol Buffer, consulte o repositório SDK C++.
C# SDK
Importante
O SDK C# / .NET está em Beta. O pacote Databricks.Zerobus está em pré-lançamento.
.NET 8.0 ou superior é obrigatório. O SDK fornece streaming nativo gRPC, OAuth e recuperação automática. Suporta JSON para configurações simples e Protocol Buffers 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 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 devolve o deslocamento do disco e WaitForOffset bloqueia até que esse disco seja durável. Para ingerir um lote, use IngestRecords, que recebe um array de registos e devolve o último offset. O bloqueio no offset é opcional.
Veja Bloqueio e reconhecimento de mensagens.
Protocol Buffers: Para ingestão com segurança de tipos, crie um stream com sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET), em que descriptorProto são os bytes DescriptorProto serializados da sua mensagem compilada, e depois faça a ingestão com stream.IngestRecord(protoBytes).
Para documentação completa, opções de configuração e exemplos de Protocol Buffer, consulte o repositório SDK C#.
TypeScript SDK
Node.js 16 ou mais é obrigatório. O SDK oferece alto desempenho com suporte assíncrono através do JavaScript Promises. Suporta JSON (o mais simples) e Protocol Buffers (recomendados para produção).
npm install @databricks/zerobus-ingest-sdk
Exemplo 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 segura por tipos, use Protocol Buffers com RecordType.Proto (por defeito) e forneça um descriptorProto nas propriedades da tabela.
Arrow Flight: Para a ingestão em colunas ou em lote de dados do Apache Arrow RecordBatch através da mesma ligação gRPC, consulte Utilizar o Arrow Flight com o Zerobus Ingest.
Para documentação completa, opções de configuração, ingestão em lote e exemplos de Protocol Buffer, consulte o repositório TypeScript SDK.
API REST
A API REST permite ingerir um único registo enviando um pedido HTTP POST ao /zerobus/v1/tables/<table-name>/insert endpoint. O próprio registo está incluído no corpo do pedido e deve estar em formato JSON.
Este exemplo explica-te como usar o CURL para enviar dados para o Zerobus Ingest usando a API REST.
Cabeçalhos
O pedido requer dois cabeçalhos HTTP específicos para autenticar e formatar corretamente o pedido.
-
Tipo-Conteúdo: application/json
- Campo obrigatório para especificar o tipo de conteúdo. Atualmente, o JSON é o único formato de mensagem suportado.
-
Autorização: Bearer <token>
- Substitui <o token> pelo token OAuth que obtiveste usando o comando curl fornecido mais tarde.
Obter o Token OAuth: Estes tokens expiram a cada hora e devem ser renovados. Podes atualizá-los voltando a buscar o token OAuth.
Preencha os seguintes parâmetros:
-
$CATALOG,$SCHEMA,$TABLE,$WORKSPACE_ID,$WORKSPACE_URL -
$DATABRICKS_CLIENT_IDe$DATABRICKS_CLIENT_SECRET- Estes dois parâmetros correspondem ao princípio de serviço que 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')
Importação de Registos:
Preencha os seguintes parâmetros:
$ZEROBUS_ENDPOINT- Conforme definido na secção Obter o URL do seu espaço de trabalho e o ponto de ingestão Zerobus.
-
$CATALOG,$SCHEMA,$TABLE,$WORKSPACE_ID,$WORKSPACE_URL $OAUTH_TOKEN- Isto foi criado na etapa anterior.
O corpo do pedido 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 toda a informação estiver preenchida corretamente, deverá receber uma resposta JSON vazia com um código de estado HTTP de 200.
Lidar com erros
Os exemplos acima mostram o caminho feliz. Em produção, envolve a ingestão em gestão de erros. O SDK repete automaticamente as tentativas em caso de erros transitórios, como problemas de rede, através do seu mecanismo de recuperação incorporado. Falhas das quais não consegue recuperar, como credenciais inválidas ou uma tabela em falta, surgem 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 recuperam automaticamente falhas transitórias e permitem-lhe resgatar registos não confirmados quando um fluxo falha permanentemente. Para padrões para clientes resilientes e a referência completa de erros, consulte Recuperação e padrões de repetição e Tratamento de erros do Zerobus Ingest.