Pular para o conteúdo principal

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

nota

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.

  1. Obtenha uma URL do Zerobus Ingest.
  2. Crie ou identifique a tabela na qual você deseja ingerir dados.
  3. Crie um Service Principal e conceda privilégios à tabela.
  4. 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:

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.

SQL
    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.

nota

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.

  1. Para criar um Service Principal, vá para Settings > Identity and Access .

  2. Em entidade de serviço , selecione gerenciar .

  3. Clique em Adicionar entidade de serviço .

  4. Na janela Adicionar entidade de serviço , crie uma nova entidade de serviço clicando em Adicionar nova .

  5. Gere e salve o ID do cliente e o segredo do cliente para a entidade de serviço.

  6. Conceda as permissões necessárias para o catálogo, o esquema e a tabela à entidade de serviço.

    1. Na página Service principal , vá para a guia tab .
    2. Copie o ID do aplicativo (UUID).
    3. 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.
    SQL
    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>`;

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.

É 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.

Bash
pip install databricks-zerobus-ingest-sdk

Exemplo de JSON:

Python
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.

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:

Python
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.

Próximos passos