Pular para o conteúdo principal

Use o Arrow Flight com o Zerobus Ingest

A ingestão via Arrow Flight permite enviar dados do Apache Arrow RecordBatch diretamente para o Zerobus Ingest em vez de converter cada linha para JSON ou Protocol Buffers (protobuf) primeiro. É uma terceira opção de formato de registro nos SDKs do Zerobus que oferecem suporte ao Arrow Flight, juntamente com JSON e protobuf, e é executada na mesma conexão gRPC. Ele usa o mesmo Endpoint do Zerobus, o mesmo fluxo de OAuth e a mesma convenção de cabeçalho x-databricks-zerobus-table-name. O protocolo de rede é o Arrow Flight DoPut, que transporta mensagens IPC do Arrow via gRPC.

Quando usar o Arrow Flight

Arrow Flight é a melhor opção nos seguintes cenários:

  • Seu aplicativo já produz dados Arrow, como pyarrow.Table ou pyarrow.RecordBatch (Python), arrow_array::RecordBatch dos crates arrow-rs (Rust) ou VectorSchemaRoot (Java). Bibliotecas de DataFrame criadas sobre Arrow, como Polars ou DataFusion, encaixam-se naturalmente neste caminho.
  • Você ingere linhas em lotes, em vez de enviar um registro por vez.
  • Seu esquema é amplo, envolve muitos dados numéricos ou é orientado para análises, onde a serialização linha por linha adiciona uma sobrecarga de CPU considerável.
  • Você está criando coletores ou gateways que agregam dados por um curto período e, em seguida, os enviam como lotes em formato de coluna única.

O Arrow Flight geralmente não é a melhor escolha para tráfego esparso, linha por linha. Nesses casos, JSON ou protobuf via caminho gRPC do SDK são normalmente mais simples. Consulte Escolha uma interface.

Como funciona o modelo de ingestão

Com a ingestão do Arrow Flight, uma transmissão grava em uma tabela de destino. Para ingerir dados, siga esta sequência:

  1. Defina um esquema Arrow que corresponda ao esquema da tabela Delta de destino.
  2. Abra uma transmissão Zerobus Arrow para essa mesa.
  3. Enviar cargas úteis RecordBatch (ou Table).
  4. Aguarde o último deslocamento ou chame flush() para confirmar a durabilidade.
  5. Feche a transmissão.

Se você usar um SDK do Zerobus, o SDK cuidará dos detalhes de baixo nível da conexão do Arrow Flight para você. Ele serializa seus dados Arrow para o formato IPC e divide automaticamente um grande lote em mensagens Flight menores e ordenadas.

O servidor reporta o progresso cumulativo à medida que os registros se tornam duráveis. Um grande lote lógico pode, portanto, ser parcialmente durável se ocorrer uma falha enquanto ele está sendo enviado. ingest_batch() ainda retorna um offset lógico para o lote enviado, e aguardar por esse offset confirma que todos os seus registros são duráveis. Consulte Os lotes do Arrow Flight são a exceção.

O Zerobus corresponde campos Arrow a colunas Delta por nome. O esquema Arrow deve incluir todas as colunas Delta obrigatórias (não anuláveis). Você pode omitir colunas anuláveis, que o Zerobus grava como NULL. Não inclua campos que estejam ausentes na tabela de destino. Os campos incluídos devem seguir a ordem relativa do esquema Delta, corresponder à sua anulabilidade e usar o tipo da coluna de destino. Para obter detalhes, consulte regras de correspondência de esquema.

Um lote Arrow não está sujeito ao limite de tamanho de mensagem gRPC de 10 MB. No entanto, cada linha individual dentro de um RecordBatch deve caber dentro do limite de 10 MB. O SDK divide automaticamente lotes maiores em várias mensagens de rede, mas não pode dividir uma linha muito grande. Veja cotas do Zerobus Ingest.

Escreva para um cliente

Os exemplos abaixo usam os SDKs Python e Rust. Para outras linguagens que suportam Arrow Flight, consulte o repository Zerobus SDK.

O SDK Python aceita um pyarrow.Schema no momento da criação da transmissão e um pyarrow.RecordBatch ou pyarrow.Table para cada chamada de ingestão.

Bash
pip install "databricks-zerobus-ingest-sdk[arrow]"
Python
import pyarrow as pa

from zerobus.sdk.sync import ZerobusSdk

# See "Get your workspace URL and Zerobus Ingest endpoint" in zerobus-ingest.md.
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"

schema = pa.schema(
[
("device_name", pa.large_utf8()),
("temp", pa.int32()),
("humidity", pa.int64()),
]
)

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

stream = sdk.create_arrow_stream(TABLE_NAME, schema, CLIENT_ID, CLIENT_SECRET)

try:
for start in range(0, 10_000, 1_000):
end = start + 1_000
batch = pa.record_batch(
{
"device_name": [f"sensor-{i}" for i in range(start, end)],
"temp": [20 + (i % 5) for i in range(start, end)],
"humidity": [55 + (i % 10) for i in range(start, end)],
},
schema=schema,
)
stream.ingest_batch(batch)
stream.flush()
finally:
stream.close()

stream.ingest_batch() também aceita um pyarrow.Table. O SDK o converte internamente para um único RecordBatch antes do envio. Cada chamada retorna um deslocamento lógico. O exemplo cria vários lotes e chama flush() uma vez para confirmar que todos os lotes pendentes são duráveis. Use wait_for_offset() quando precisar confirmar um lote específico antes de continuar; aguardar o último deslocamento também confirma todos os deslocamentos anteriores. Para saber quando aguardar e como funciona o reconhecimento, consulte Bloqueio e reconhecimento de mensagens.

Ingestão de colunas VARIANT

O Apache Arrow não possui um tipo VARIANT nativo. Para realizar a ingestão em uma coluna VARIANT via Arrow Flight, crie os campos metadata e value de suporte da coluna como uma estrutura (struct) de duas colunas LargeBinary e, em seguida, inclua essa estrutura no seu RecordBatch. Por meio dos SDKs gRPC e REST, você passa um valor Variant como uma **strings** codificada em JSON. Consulte Tipos de dados suportados.

O exemplo de Rust a seguir cria uma coluna de struct VARIANT a partir de linhas JSON e a ingere:

Rust
fn variant_struct(json_rows: &[&str]) -> ArrayRef {
let mut metas: Vec<Vec<u8>> = Vec::new();
let mut vals: Vec<Vec<u8>> = Vec::new();
for json in json_rows {
let mut vb = VariantBuilder::new();
vb.append_json(json).expect("invalid JSON for variant");
let (metadata, value) = vb.finish();
metas.push(metadata);
vals.push(value);
}
let fields = Fields::from(vec![
Field::new("metadata", DataType::LargeBinary, false),
Field::new("value", DataType::LargeBinary, false),
]);
let meta_arr = Arc::new(LargeBinaryArray::from_iter_values(metas)) as ArrayRef;
let val_arr = Arc::new(LargeBinaryArray::from_iter_values(vals)) as ArrayRef;
Arc::new(StructArray::try_new(fields, vec![meta_arr, val_arr], None).expect("variant struct"))
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let client_id = std::env::var("DATABRICKS_CLIENT_ID")?;
let client_secret = std::env::var("DATABRICKS_CLIENT_SECRET")?;

let variant_type = DataType::Struct(Fields::from(vec![
Field::new("metadata", DataType::LargeBinary, false),
Field::new("value", DataType::LargeBinary, false),
]));
let schema = Arc::new(ArrowSchema::new(vec![
Field::new("id", DataType::Int32, true),
Field::new("payload", variant_type, true),
]));

let sdk = ZerobusSdk::builder()
.endpoint(ENDPOINT)
.unity_catalog_url(UC_URL)
.build()?;

let mut stream = sdk
.stream_builder()
.table(TABLE)
.oauth(&client_id, &client_secret)
.arrow(schema.clone())
.ipc_compression(None)
.build_arrow()
.await?;

let ids = Int32Array::from(vec![1, 2, 3]);
let payload = variant_struct(&[
r#"{"user":"alice","tags":[1,2,3]}"#,
r#""just a string""#,
r#"{"nested":{"a":true,"b":null,"c":3.14}}"#,
]);
let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(ids) as ArrayRef, payload])?;

let offset = stream.ingest_batch(batch).await?;
stream.flush().await?;
stream.close().await?;
Ok(())
}

Este exemplo está em Rust. Para uso equivalente em outras linguagens, consulte o repository do SDK Zerobus.

Compressão IPC

Por padrão, os payloads do Arrow IPC são enviados sem compressão. É possível, opcionalmente, compactá-los no tráfego utilizando um de dois codecs.

  • LZ4_FRAME: rápido, baixa sobrecarga de CPU, taxa de compressão modesta. Prefira isto quando o cliente estiver com restrições de CPU, mas ainda desejar reduzir os bytes na rede.
  • ZSTD: Maior taxa de compressão, mais CPU por lote. Recomenda-se habilitá-lo sempre que o cliente puder arcar com o custo adicional da CPU.

A compressão reduz bytes na rede, mas acarreta um custo de CPU no cliente. Cargas menores podem evitar gargalos de rede e reduzir os custos de rede.

Defina o campo ipc_compression em ArrowStreamConfigurationOptions:

Python
from zerobus.sdk.shared.arrow import IPCCompression, ArrowStreamConfigurationOptions

options = ArrowStreamConfigurationOptions(ipc_compression=IPCCompression.ZSTD)
stream = sdk.create_arrow_stream(
TABLE_NAME, schema, CLIENT_ID, CLIENT_SECRET, options=options
)

Melhores práticas

Siga estas orientações para obter o melhor desempenho e confiabilidade da ingestão de dados do Arrow Flight.

  • Reutilizar uma transmissão por muitos lotes em vez de abrir uma nova transmissão por lotes. A criação de transmissões acarreta custos indiretos significativos que podem ser amortizados reutilizando uma transmissão em vários lotes.
  • Enviar várias linhas por lote. Comece com lotes do tamanho ideal para a aplicação, e não com uma linha por chamada. Enviar uma linha de cada vez funciona, mas anula a maior parte da vantagem de desempenho de usar o Arrow.
  • Chame flush() em pontos de verificação controlados. Isso proporciona um limite de durabilidade claro para um grupo de lotes, sem bloquear cada um individualmente.
  • Habilite a compressão IPC para melhorar o throughput. ZSTD é recomendado para a maioria das cargas de trabalho quando o cliente tem CPU disponível. Use LZ4_FRAME ou nenhuma compressão se o cliente estiver com a CPU limitada.
  • Use o Arrow Flight quando seu produtor já for colunar. Se seus dados de origem forem naturalmente orientados a linhas e pequenos, usar o Zerobus Ingest com JSON ou protobuf costuma ser mais simples. Consulte Use Zerobus Ingest.

Tratamento e recuperação de erros

As transmissões do Arrow Flight usam as mesmas categorias de erro gRPC que o restante do Zerobus Ingest. Para obter informações sobre códigos de erro, orientações sobre como tentar novamente e a taxonomia completa de cliente versus servidor, consulte Tratamento de erros do Zerobus Ingest.

Ao configurar o SDK com recuperação automática (o default), ele se reconecta de forma transparente e reproduz lotes não confirmados em caso de falhas transitórias. Após o fechamento de uma transmissão com trabalho não confirmado, o SDK retém os lotes que o cliente aceitou, mas que o servidor não confirmou, incluindo lotes que podem ainda não ter sido enviados.

Após a recuperação automática ser esgotada, corrija a causa da falha e chame close() para finalizar os lotes não confirmados da transmissão. Como a transmissão já falhou, close() pode retornar o mesmo erro terminal, mesmo que a finalização seja bem-sucedida.

Você pode chamar get_unacked_batches() apenas após a transmissão ser encerrada. Ele retorna os lotes retidos para persistência ou reprodução gerenciada pelo aplicativo. Como você cria uma transmissão de substituição, persiste os lotes e os tenta novamente depende da política de recuperação do seu aplicativo.

Python
from zerobus.sdk.shared import ZerobusException

try:
stream.close()
except ZerobusException:
# The terminal error can be returned after closure is finalized.
pass

unacked_batches = stream.get_unacked_batches()

Recursos adicionais

  • Use 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 passos são compartilhados entre todos os formatos de registro.
  • Cotas do Zerobus Ingest: revise as cotas do Zerobus antes de implantar em produção. Os limites de throughput, latência e tabela particionada se aplicam ao Arrow Flight.
  • Tratamento de erros do Zerobus Ingest: consulte esta página para obter uma lista completa de códigos de erro gRPC e o comportamento recomendado de repetição e recuperação para seu cliente.