Esta página descreve as características de desempenho, a partir da versão 2.75.0 do Apache Beam, para jobs de streaming do Dataflow que leem tabelas do Apache Kafka e gravam em tabelas do Apache Iceberg. Ela avalia as diferenças de desempenho entre gravações diretas do Apache Iceberg e gravações roteadas pela API BigQuery gerenciada e compara esses resultados com os comparativos de referência de pipelines do Kafka para o BigQuery. Como as otimizações para E/S do Apache Iceberg estão em andamento, essas métricas de desempenho estão sujeitas a mudanças.
As comparações de referência estão disponíveis em três configurações de mapeamento sem estado (ou seja, elas leem da origem, convertem a mensagem em um registro e gravam no coletor sem rastrear o estado ou aplicar uma lógica de negócios complexa; referidas como map_only ou mapping em comparativos):
- Kafka para BigQuery (
map_only) (referência de desempenho do Kafka para o BigQuery) - Kafka para Iceberg Direct (
map_only,autosharding=false) - Kafka para Iceberg usando a API BigQuery gerenciada (
map_only)
Além disso, este guia avalia padrões de streaming direto do Apache Iceberg, como o lote com estado usando groupbykey, e detalha considerações downstream importantes sobre distribuições de tamanho de arquivo, comportamento de fragmentação automática e latência de consulta do lado de leitura.
Metodologia de teste
Os comparativos foram realizados usando os seguintes recursos:
- Cluster do serviço gerenciado para Apache Kafka:o tráfego foi gerado usando o modelo do gerador de dados de streaming do Dataflow.
- Capacidade de entrada:1 GBps
- Taxa de mensagens:cerca de 1.000.000 de mensagens por segundo
- Formato da mensagem:texto JSON com um esquema fixo (cerca de 1 KB por mensagem)
- Partições:1.000 partições do Kafka
- Coletores de destino:
- BigQuery:tabela padrão (não particionada) gravada usando a API BigQuery Storage Write.
- Apache Iceberg:catálogo com suporte do Cloud Storage. O coletor direto é particionado usando
bucket(id, 64)(em 64 fragmentos na chave primária) e usa o modo de distribuiçãohash.
Depois que o escalonamento automático horizontal foi estabilizado, cada configuração de pipeline foi executada em estado estável por 24 horas. Os comparativos de cada caso de pipeline foram executados três vezes separadas, e todos os valores informados representam as médias calculadas nessas execuções para garantir métricas de desempenho confiáveis e sustentadas.
Desempenho de ingestão: cargas de trabalho de mapeamento
Os pipelines de mapeamento sem estado leem da origem, convertem o formato da mensagem em um registro e gravam no coletor sem rastrear o estado entre os registros. As seções a seguir analisam arquiteturas de referência executadas a 1 GBps.
Configuração do job
| Configuração | Kafka para BigQuery (map_only) |
Kafka para Iceberg Direct (autosharding=false) |
Kafka para Iceberg usando a API BigQuery gerenciada |
|---|---|---|---|
| Tipo de máquina do worker | e2-standard-2 |
e2-standard-4 |
e2-standard-4 |
| vCPUs por worker | 2 | 4 | 4 |
| RAM por worker | 8 GB | 16 GB | 16 GB |
| Streaming Engine | Ativado | Ativado | Ativado |
| Escalonamento automático horizontal | Ativado | Ativado | Ativado |
| Frequência de acionamento | 5 segundos | 60 segundos | 60 segundos |
Capacidade de processamento e uso de recursos
A gravação direta em arquivos Parquet físicos no armazenamento de objetos incorre em uma sobrecarga de E/S maior do que a ingestão de streaming do BigQuery. Em comparação com as gravações diretas do Iceberg, o roteamento de gravações pela API BigQuery gerenciada melhora a utilização da CPU do worker (cerca de 70% em vez de cerca de 60%) e reduz modestamente o consumo do Streaming Engine (cerca de 180 SECU/h em vez de cerca de 200 SECU/h), embora os requisitos gerais de recursos computacionais do worker permaneçam semelhantes (cerca de 440 vCPUs em vez de cerca de 450 vCPUs).
| Métrica | Kafka para BigQuery (map_only) |
Kafka para Iceberg Direct (autosharding=false) |
Kafka para Iceberg usando a API BigQuery gerenciada |
|---|---|---|---|
| Capacidade média de entrada por worker | Cerca de 15 MBps | Cerca de 9 MBps | Cerca de 9 MBps |
| Uso médio da CPU | Cerca de 70% | Cerca de 60% | Cerca de 70% |
| vCPUs estimadas para entrada de 1 GBps | Cerca de 126 vCPUs | Cerca de 450 vCPUs | Cerca de 440 vCPUs |
| Workers estimados para entrada de 1 GBps | Cerca de 63 workers | Cerca de 110 workers | Cerca de 110 workers |
| SECU estimado por hora para 1 GBps | Cerca de 58 SECU/h | Cerca de 200 SECU/h | Cerca de 180 SECU/h |
Perfil de latência de gravação
As gravações diretas do Iceberg exibem latência de cauda grave (P99) devido a restrições de confirmação de metadados de armazenamento de objetos. O uso da API BigQuery gerenciada elimina picos de latência de cauda, mantendo a latência mediana baixa.
| Latência de gravação de ponta a ponta | Kafka para BigQuery | Kafka para Iceberg Direct (autosharding=false) |
Kafka para Iceberg usando a API BigQuery gerenciada |
|---|---|---|---|
| P50 (mediana) | Cerca de 1.200 ms | Cerca de 1.000 ms | Cerca de 1.000 ms |
| P95 | Cerca de 3.000 ms | Cerca de 7.400 ms | Cerca de 1.900 ms |
| P99 (cauda) | Cerca de 5.400 ms | Cerca de 14.000 ms | Cerca de 2.700 ms |
Considerações sobre fragmentação automática e escolhas de design
Esta seção discute as implicações da fragmentação automática nos tamanhos de arquivo e na latência do pipeline ao gravar no Apache Iceberg.
Por que autosharding=false foi escolhido como a referência
Nos testes iniciais, a ativação da fragmentação automática fez com que os tamanhos de arquivo fossem reduzidos a pequenos blocos e flutuassem arbitrariamente devido à divisão dinâmica de fragmentos acionada por picos de carga localizados no nível da linha de execução, mesmo sob uma carga de entrada agregada constante.
Para manter layouts de arquivo Parquet estáveis e previsíveis (cerca de 800 KB em média) e garantir uma referência justa sem liberações prematuras, autosharding=false foi escolhido para a configuração do coletor direto.
O que acontece se você desativar a fragmentação automática em vez de mantê-la ativada?
- Com
autosharding=false(referência): você consegue tamanhos de arquivo iniciais maiores (cerca de 800 KB em média) em comparação com a fragmentação automática. Embora ainda seja pequeno em comparação com os tamanhos de arquivo ideais do Iceberg (128 a 512 MB), ele exige uma compactação downstream significativamente menor. No entanto, a desvantagem é uma alta latência de cauda de gravação (P99 atingindo cerca de 14,0 s) devido a gargalos de metadados de armazenamento de objetos. - Se a fragmentação automática estiver ativada:o Dataflow escalona dinamicamente as linhas de execução do gravador para absorver picos de capacidade de processamento local, o que reduz a latência de cauda de gravação. No entanto, ele compromete a camada de armazenamento, produzindo um grande volume de arquivos Parquet pequenos e fragmentados (cerca de 100 KB ou menos). Esses tamanhos de arquivo exibem alta variância e flutuam arbitrariamente entre as execuções (variando de cerca de 39 KB a cerca de 100 KB em média), aumentando a necessidade de manutenção agressiva da compactação downstream.
Ajuste e recomendações de partição
Durante nossa avaliação, testamos vários valores de partição fixos para a tabela de destino para encontrar um equilíbrio ideal. Descobrimos que o uso de 64 buckets (por exemplo, bucket(id, 64)) para o particionamento da tabela de destino gerou os tamanhos de arquivo desejados, mantendo a utilização e a capacidade de processamento adequadas. Essa abordagem nos ajudou a combinar os benefícios de desempenho da fragmentação automática, evitando os problemas de fragmentação de tamanho de arquivo arbitrários vinculados ao escalonamento totalmente dinâmico.
Recomendação para profissionais:os clientes são incentivados a realizar testes preliminares semelhantes com configurações de partição direcionadas para localizar o ponto ideal que maximiza o paralelismo do pipeline sem comprometer os tamanhos de arquivo Parquet.
Implicações de leitura downstream: tamanhos de arquivo e compactação
Embora as métricas do lado de gravação favoreçam a API BigQuery gerenciada para ingestão do Iceberg, a eficiência geral do pipeline depende muito do desempenho de leitura downstream:
- Geração de arquivos pequenos na API BigQuery Gerenciada:a API BigQuery Gerenciada libera dados com frequência para garantir baixa latência de gravação. Esse comportamento resulta em um grande volume de arquivos Parquet pequenos gravados no catálogo de destino do Iceberg.
- Impacto da latência de consulta de leitura:mecanismos de consulta (por exemplo, Starburst/Trino, Apache Spark, BigQuery, Dremio) que leem tabelas com milhões de arquivos Parquet pequenos incorrem em uma grande sobrecarga de análise de metadados e penalidades de verificação de partição.
- Requisitos de compactação:para evitar a degradação do desempenho de leitura ao usar a API BigQuery gerenciada (ou se a fragmentação automática estiver ativada em gravações diretas), execute jobs de manutenção de compactação do Iceberg regulares (por exemplo,
REWRITE DATA FILES). A sobrecarga de recursos computacionais para compactação precisa ser considerada no design geral da arquitetura. - Distribuição de arquivos de gravação direta (
autosharding=false):as gravações diretas do Iceberg com fragmentação fixa produzem arquivos Parquet médios maiores (cerca de 800 KB), resultando em um layout menos fragmentado para acesso imediato à consulta sem demandas de compactação imediata (embora ainda abaixo do intervalo ideal).
Pipelines diretos com estado do Iceberg (groupbykey)
Para avaliar estratégias de lote manual, o agrupamento de chaves com estado (groupbykey) foi testado em relação ao pipeline de Kafka para Iceberg Direct (map_only, autosharding=false) de referência. As duas configurações gravam arquivos Parquet diretamente no armazenamento de objetos.
Comparativo de mercado
| Métrica / recurso | Referência do coletor direto (autosharding=false) |
Coletor direto com estado (groupbykey) |
Impacto no desempenho |
|---|---|---|---|
| vCPUs estimadas para 1 GBps | Cerca de 450 vCPUs | Cerca de 520 vCPUs | Cerca de +16% de computação necessária |
| Uso médio da CPU | Cerca de 60% | Cerca de 50% | Cerca de -17% de eficiência do worker |
| SECU/h estimado para 1 GBps | Cerca de 200 SECU/h | Cerca de 300 SECU/h | Cerca de +50% de carga do Streaming Engine |
| Tamanho médio do arquivo | Cerca de 800 KB | Cerca de 100 KB | Gera lotes de arquivos menores |
| Latência P50 | Cerca de 1.000 ms | Cerca de 1.200 ms | Cerca de +20% de mediana mais lenta |
| Latência P95 | Cerca de 7.400 ms | Cerca de 5.500 ms | Cerca de -26% de latência menor |
| Latência P99 | Cerca de 14.000 ms | Cerca de 13.000 ms | Mudança marginal na latência de cauda |
análise de compensação
- Sobrecarga do Streaming Engine:a adição de uma etapa
groupbykeycom estado exige que o Beam armazene o estado intermediário entre os limites da janela. Isso aumenta o consumo de unidades de computação do Streaming Engine em cerca de 50% (de cerca de 200 SECU/h para cerca de 300 SECU/h). - Latência de buffer:a agregação manual de chaves introduz o buffer de janela obrigatório, aumentando a latência de gravação mediana (P50) para cerca de 1.200 ms e a latência P95 para cerca de 5,5 s.
Pipelines inversos: streaming do Iceberg para o Kafka
Para avaliar os recursos bidirecionais do lakehouse, também foram realizados comparativos de dados de streaming que fluem na direção inversa: leitura de fluxos somente de anexação de uma tabela do Apache Iceberg e publicação de volta no Apache Kafka.
Configuração e eficiência do job
Ao contrário dos pipelines de ingestão que precisam lidar com gravações de arquivos de armazenamento de objetos pesados ou gargalos de confirmação de metadados, a leitura e o streaming de mudanças do Iceberg operam com alta eficiência:
| Métrica | Iceberg para Kafka (somente anexação, exatamente uma vez) |
|---|---|
| Tipo de máquina do worker | e2-standard-4 |
| vCPUs estimadas para entrada de 1 GBps | Cerca de 30 vCPUs |
| Workers estimados para entrada de 1 GBps | Cerca de 7 workers |
| SECU estimado por hora para 1 GBps | Cerca de 0,2 SECU/h |
Principais aprendizados para pipelines inversos
- Sobrecarga de computação significativamente menor:a leitura e a projeção de fluxos de CDC do Iceberg exigem substancialmente menos recursos de computação (cerca de 30 vCPUs em vez de cerca de 450 vCPUs para gravações diretas), porque evitam o trabalho pesado de particionamento, codificação e confirmação de grandes volumes de arquivos Parquet para armazenamento de objetos.
- Eficiência de recursos:o consumo ou a replicação downstream orientados a eventos de um formato de lakehouse de volta para camadas de streaming são altamente eficientes em comparação com os caminhos de ingestão de entrada.
Resumo das recomendações arquitetônicas
| Padrão de arquitetura | Latência de gravação P99 | Layout do arquivo | Considerações de leitura downstream |
|---|---|---|---|
Kafka para BigQuery (map_only) |
Cerca de 5,4 s | N/A | Ideal (mecanismo de armazenamento do BigQuery gerenciado) |
| Kafka para Iceberg usando a API BigQuery Gerenciada | Cerca de 2,7 s | Arquivos arbitrariamente pequenos | Exige compactação periódica para leituras de alto volume |
Kafka para Iceberg Direct (autosharding=false) |
Cerca de 14,0 s | Cerca de 800 KB | Bom (tamanhos de arquivo iniciais maiores, menor demanda de compactação) |
Kafka para Iceberg Direct (groupbykey) |
Cerca de 13,0 s | Cerca de 100 KB | Moderado (maior sobrecarga de computação e estado) |
Estimar custos
É possível estimar o custo de referência do seu próprio pipeline comparável com o faturamento baseado em recursos usando a Google Cloud calculadora de preços, da seguinte maneira:
- Abra a calculadora de preços.
- Clique em Adicionar à estimativa.
- Selecione Dataflow.
- Em Tipo de serviço, selecione "Dataflow clássico".
- Selecione Configurações avançadas para mostrar o conjunto completo de opções.
- Escolha o local em que o job é executado.
- Em Tipo de serviço, selecione "Streaming".
- Selecione Ativar o Streaming Engine.
- Insira informações sobre as horas de execução do job, nós de worker, máquinas de worker e armazenamento de Persistent Disk.
- Insira o número estimado de unidades de computação do Streaming Engine.
O uso de recursos e o custo são escalonados de forma aproximadamente linear com a capacidade de processamento de entrada, embora, para jobs pequenos com apenas alguns workers, o custo total seja dominado por custos fixos. Como ponto de partida, é possível extrapolar o número de nós de worker e o consumo de recursos dos resultados de referência.
Por exemplo, suponha que você execute um pipeline usando a arquitetura Kafka para Iceberg Direct (autosharding=false), com uma taxa de dados de entrada de 100 MBps. Com base nos resultados de referência para um pipeline de 1 GBps, é possível estimar os requisitos de recursos da seguinte maneira:
- Fator de escalonamento: (100 MBps) / (1024 MBps) = cerca de 0,1
- Nós de worker projetados: 110 workers × 0,1 = cerca de 11 workers
- Número projetado de unidades de computação do Streaming Engine por hora: 200 × 0,1 = cerca de 20 unidades por hora
Esse valor só pode ser usado como uma estimativa inicial. A capacidade de processamento e o custo reais podem variar significativamente, com base em fatores como tipo de máquina, distribuição de tamanho da mensagem, código do usuário, tipo de agregação, paralelismo de chave e tamanho da janela. Para mais informações, consulte Práticas recomendadas para otimização de custos do Dataflow.
Executar um pipeline de teste
Para implantar um job de streaming do Apache Iceberg usando o modelo Flex
do Dataflow, use o
gcloud dataflow flex-template run
comando.
gcloud dataflow flex-template run JOB_NAME \
--project=PROJECT_ID \
--region=REGION \
--template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
--enable-streaming-engine \
--parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'
Substitua:
JOB_NAME: o nome do job do DataflowPROJECT_ID: o ID do seu Google Cloud projetoREGION: a Google Cloud região em que o job é executado (por exemplo,us-central1)KAFKA_BOOTSTRAP_ADDRESS: o endereço de inicialização do cluster do Apache KafkaKAFKA_TOPIC: o nome do tópico do KafkaICEBERG_TABLE_IDENTIFIER: o identificador da tabela de destino do IcebergCATALOG_NAME: o nome do catálogo do IcebergCATALOG_TYPE: o tipo de catálogo a ser usado (por exemplo,hadoopoubigquery)BUCKET_NAME: o nome do bucket do Cloud Storage para o local do data warehouseSCHEMA_DEFINITION: a definição de esquema para os dados do tópico do Kafka (por exemplo,{"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})