Use APIs compatíveis com Kafka com o Zerobus Ingest
Beta
Este recurso está em Beta e está disponível apenas na AWS e no Azure.
Zerobus Ingest oferece APIs de produtor compatíveis com Kafka que permitem a ingestão usando qualquer cliente produtor do Apache Kafka sem um SDK do Databricks. Você aponta um produtor Kafka existente para o endpoint do Zerobus e produz para um tópico nomeado de acordo com sua tabela de destino, e os registros são armazenados em uma tabela Delta do Unity Catalog. As APIs compatíveis com Kafka são adequadas quando você já possui um produtor Kafka, um coletor que utiliza o protocolo Kafka ou ferramentas que emitem para o Kafka, e deseja rotear esses dados para o Delta com o mínimo de alterações no código.
As APIs compatíveis com Kafka estão em execução sobre TCP com SASL_SSL e o mecanismo OAUTHBEARER e implementam o subconjunto do lado do produtor do protocolo Kafka, que cobre Produce, Metadata, ApiVersions e as APIs de handshake SASL.
[[ ## completed ##]] As APIs de consumidor, administrador e transacionais não estão disponíveis. As APIs são somente para escrita.
Quando usar as APIs compatíveis com Kafka
As APIs compatíveis com Kafka são a melhor opção nos seguintes cenários:
- Você deseja enviar dados para o Delta sem adotar um SDK do Zerobus, e você já tem em execução um produtor Kafka ou um aplicativo, agente ou coletor que emite para o Kafka. [[ ## completed ##]]
- Você deseja reutilizar sua configuração de produtor Kafka, processamento em lote e ferramentas operacionais existentes.
- Você envia registros JSON e não precisa de Protocol Buffers ou Apache Arrow.
Se você estiver criando um novo cliente do zero e quiser o maior throughput, confirmações por registro e recuperação automática, use um SDK do Zerobus via gRPC em vez das APIs compatíveis com Kafka. Consulte Escolher uma interface. Para cargas de trabalho colunares ou em lotes, consulte Usar o Arrow Flight com o Zerobus Ingest.
Como funciona o modelo de ingestão
As APIs compatíveis com Kafka mapeiam os conceitos do Kafka para o Zerobus Ingest da seguinte forma:
- Tópicos. Um nome de tópico Kafka é o nome completo da tabela do Unity Catalog de três níveis (
catalog.schema.table). A tabela de destino já deve existir, pois o Zerobus nunca cria tópicos. - Registros. O Zerobus ingere apenas o valor do registro, que deve ser um objeto JSON codificado em UTF-8 que corresponda ao esquema da tabela Delta de destino. Ele ignora chaves de registro, cabeçalhos e a partição e o timestamp fornecidos pelo cliente, e não os persiste.
- Autenticação. Cada conexão é autenticada com um token OAuth do Databricks com escopo para a tabela de destino, apresentado via SASL/
OAUTHBEARER. Consulte Autenticação. - Reconhecimentos. O Zerobus retorna uma resposta
Producesomente após persistir os registros de forma durável. Configure seu produtor comacks=all.
O Zerobus não possui partições por design. O endpoint do Zerobus é um único broker lógico e uma única partição, portanto, cada registro de um tópico é confirmado na partição 0. Os producers não precisam dar account disso. O Zerobus Ingest escala horizontalmente para lidar com a carga de entrada.
Além disso, o Zerobus Ingest oferece entrega pelo menos uma vez. Uma única conexão sustenta aproximadamente 50.000 mensagens por segundo, o que é inferior ao caminho gRPC do SDK. Para obter a maior taxa de transferência (throughput), use um SDK do Zerobus em vez das APIs compatíveis com Kafka. Os limites de latência, cota, tamanho de registro e tabela particionada são compartilhados com o restante do Zerobus Ingest. Consulte cotas do conector do Zerobus Ingest.
Autenticação
As APIs compatíveis com Kafka usam SASL/OAUTHBEARER. O token de portador é um token OAuth do Databricks que você obtém com as credenciais de cliente de uma Service Principal com acesso à tabela de destino.
[[ ## completed ##]] O token tem escopo para essa tabela por meio do OAuth authorization_details e usa o recurso zerobusDirectWriteApi, o mesmo fluxo da API REST do Zerobus.
Os tokens OAuth expiram após uma hora, portanto, forneça o token por meio do callback do provedor de token do seu cliente Kafka em vez de como uma string estática. O cliente então busca um novo token sempre que se reconecta. As conexões também têm um tempo de vida limitado no lado do servidor. Quando uma conexão atinge esse limite, o Zerobus a fecha, e o produtor se reconecta e se reautentica por conta própria. A reautenticação em uma conexão ativa não é suportada.
Conceda à Service Principal os privilégios necessários do Unity Catalog na tabela de destino antes de se conectar. Consulte Criar um Service Principal e conceder permissões.
Gravar um cliente
O exemplo abaixo produz para a mesma tabela air_quality usada nos exemplos Usar o conector do Zerobus Ingest. Ele usa kafka-python, mas qualquer cliente produtor Kafka que suporte SASL_SSL com o mecanismo OAUTHBEARER funciona. Adapte o padrão de provedor de token à sua biblioteca cliente.
O produtor conecta-se ao servidor bootstrap do Zerobus na porta 9092. Encontre o ID e a região do seu workspace conforme descrito em Obter a URL do workspace e o endpoint do Zerobus Ingest.
- Servidor de bootstrap:
<workspace-id>.zerobus.<region>.cloud.databricks.com:9092
pip install kafka-python requests
o passo 1: Criar um provedor de tokens
[[ ## completed ##]]
O Zerobus autentica cada conexão com um token OAuth do Databricks de curta duração e escopo de tabela. Como o token expira, passe um callback que gera um novo token sob demanda em vez de um token estático.
A função fetch_zerobus_token() troca suas credenciais de Service Principal por um token com escopo para a tabela de destino, e ZerobusTokenProvider o encapsula na interface de retorno de chamada que kafka-python espera.
[[ ## completed ##]]
import json
import requests
from kafka.sasl.oauth import AbstractTokenProvider
# See "Get your workspace URL and Zerobus Ingest endpoint" in zerobus-ingest.md.
WORKSPACE_ID = "1234567890123456"
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"
def fetch_zerobus_token():
catalog, schema, table = TABLE_NAME.split(".")
authorization_details = [
{
"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": f"{catalog}.{schema}",
},
{
"type": "unity_catalog_privileges",
"privileges": ["SELECT", "MODIFY"],
"object_type": "TABLE",
"object_full_path": TABLE_NAME,
},
]
response = requests.post(
f"{WORKSPACE_URL}/oidc/v1/token",
auth=(CLIENT_ID, CLIENT_SECRET),
data={
"grant_type": "client_credentials",
"scope": "all-apis",
"resource": f"api://databricks/workspaces/{WORKSPACE_ID}/zerobusDirectWriteApi",
"authorization_details": json.dumps(authorization_details),
},
timeout=30,
)
response.raise_for_status()
return response.json()["access_token"]
# kafka-python calls token() whenever it needs a fresh OAuth token.
class ZerobusTokenProvider(AbstractTokenProvider):
def token(self):
return fetch_zerobus_token()
o passo 2: configurar o producer e enviar registros
Aponte o produtor para o servidor bootstrap, configure SASL_SSL com o mecanismo OAUTHBEARER e passe o provedor de tokens do passo 1. Use acks="all" para que cada lote seja confirmado somente após ser persistido de forma durável, e envie os registros sem compressão.
from kafka import KafkaProducer
BOOTSTRAP_SERVERS = "1234567890123456.zerobus.us-west-2.cloud.databricks.com:9092"
producer = KafkaProducer(
bootstrap_servers=BOOTSTRAP_SERVERS,
security_protocol="SASL_SSL",
sasl_mechanism="OAUTHBEARER",
sasl_oauth_token_provider=ZerobusTokenProvider(),
# Wait for durable acknowledgement before treating a record as ingested.
acks="all",
# Compression is not supported by the endpoint; send records uncompressed.
compression_type=None,
)
# Each send() returns a future immediately. The topic is the full table name.
futures = [
producer.send(
topic=TABLE_NAME,
value=json.dumps(
{"device_name": f"sensor-{i}", "temp": 20 + i % 15, "humidity": 50 + i % 40}
).encode("utf-8"),
)
for i in range(1000)
]
producer.flush()
# Block on each future to confirm every record was durably acknowledged.
for future in futures:
future.get(timeout=30)
producer.close()
print("All records ingested successfully")
Opções de configuração
As APIs compatíveis com Kafka implementam o subconjunto de produtor do protocolo Kafka. Configure seu produtor de acordo com as seguintes opções.
Opção | Detalhes |
|---|---|
Formato de registro | Somente JSON. Cada valor de registro deve ser um objeto JSON codificado em UTF-8 que corresponda ao esquema da tabela de destino. Para enviar Protocol Buffers ou Avro, use um SDK do Zerobus. |
Compressão | Não suportado. Envie lotes não compactados, por exemplo |
Campos de registro | Apenas valor. O Zerobus ingere o valor do registro e não persiste key, cabeçalhos, atribuições de partição ou Timestamp. |
Suporte a API | Somente escrita. O Zerobus aceita solicitações |
Link privado front-end | Não suportado. Conecte-se pelo endpoint público. |
Esquema | Aplicado. O Zerobus rejeita registros com campos que não correspondem ao esquema da tabela de destino e trata colunas Delta anuláveis extras como uma alteração não disruptiva. Para capturar campos não correspondentes em vez de rejeitá-los, configure uma coluna de resgate. |
Para obter mais informações sobre o Link privado front-end, consulte Conceitos de Private Link.
Para limites de latência, cota, tamanho de registro e tabela particionada, consulte cotas do conector Zerobus Ingest.
Práticas recomendadas
Siga estas diretrizes para obter o melhor desempenho e confiabilidade das APIs compatíveis com Kafka. Uma única conexão pode sustentar aproximadamente 50.000 mensagens por segundo.
- Reutilize um producer de longa duração em muitos registros em vez de criar um por lotes, pois a criação do producer e o handshake SASL acarretam custos de configuração.
- Permita que o produtor acumule registros em lotes, por exemplo, ajustando
linger.msebatch.size, em vez de liberar após cada registro. O processamento em lotes é a maior alavanca para o throughput. - Use
acks=allpara obter uma confirmação durável para cada lote, correspondendo à semântica de "pelo menos uma vez" do Zerobus. - Obtenha tokens OAuth por meio do callback do provedor de token do seu cliente para que eles recebam refresh automaticamente na reconexão, em vez de passar um token estático que expira. [[ ## completed ##]]
- Faça a execução do produtor na mesma região de cloud que o Endpoint do Zerobus para obter o throughput máximo.
Tratamento de erros
O Zerobus relata falhas usando códigos de erro padrão do Kafka no tópico e na partição afetados. Os códigos comuns incluem:
Erro do Kafka | Significado |
|---|---|
| O token OAuth está ausente, inválido ou não possui os privilégios necessários do Unity Catalog na tabela. |
| A tabela de destino não existe, foi descartada ou o token não está autorizado a gravar nela. |
| Um registro falhou na validação do esquema ou não pôde ser decodificado como JSON UTF-8. |
| Um único registro excede o limite de tamanho de registro de 10 MB. Consulte Tamanho do registro. |
| O lote foi comprimido. Envie registros não comprimidos. |
Após uma solicitação Produce com falha, o Zerobus retorna os códigos de erro por partição e fecha a conexão. Os produtores Kafka reconectam-se automaticamente, mas projete seu cliente para exibir falhas de envio, por exemplo, inspecionando o resultado de cada envio, para que os registros não sejam descartados silenciosamente.
Recursos adicionais
- Use o conector do Zerobus Ingest: se você ainda não configurou o Zerobus Ingest, comece aqui para obter instruções sobre como encontrar a URL do seu workspace, criar a tabela Delta de destino e configurar um service principal. Estes os passos são compartilhados entre todas as interfaces.
- Quotas do conector do Zerobus Ingest: revise as quotas e os limites do Zerobus antes de implantar em produção.
- Use o Arrow Flight com o Zerobus Ingest: Para ingestão colunar ou em lote via gRPC.