Entender o paralelismo no Dataflow

O Dataflow foi projetado para executar grandes pipelines de processamento de dados distribuindo o trabalho em um pool gerenciado de instâncias de computação. Entender como o Dataflow paraleliza o processamento ajuda a projetar pipelines eficientes, evitar gargalos de desempenho e otimizar os custos de recursos.

Esta página explica como o Dataflow paraleliza o processamento de dados, como ele gerencia e dimensiona a execução, os fatores comuns que restringem o paralelismo e as técnicas que podem ser usadas para otimizar a capacidade de processamento do pipeline.

Modelos de paralelismo: horizontal x vertical

O Dataflow alcança o paralelismo usando duas estratégias complementares:

  • Paralelismo horizontal:os dados do pipeline são particionados e processados em várias instâncias de worker (máquinas virtuais) simultaneamente. O Dataflow pode ajustar automaticamente o tamanho do pool de workers com base na demanda de carga de trabalho usando o escalonamento automático horizontal. Por padrão, o Dataflow define um limite de recursos de 4.000 workers por job, que pode ser ajustado usando solicitações de cota.

  • Paralelismo vertical:vários núcleos de CPU e linhas de execução em uma única instância de worker processam dados de pipeline simultaneamente. Cada VM de worker executa processos de worker e linhas de execução de harness para usar os recursos de computação disponíveis. Com a escala dinâmica de linhas de execução, o Dataflow pode ajustar o número de linhas de execução ativas por worker em pipelines em lote com base na utilização da CPU e na capacidade de memória. No Dataflow Prime, o escalonamento automático vertical escalona dinamicamente a memória e a computação alocadas aos workers.

Unidades de trabalho e hierarquia de execução

Para distribuir o processamento entre workers e linhas de execução, o Dataflow divide os pipelines do Apache Beam em unidades de trabalho discretas:

  • PCollections e partições:um PCollection representa um conjunto de dados distribuído. Para dados limitados (pipelines em lote), o Dataflow divide o conjunto de dados em divisões ou fragmentos. Para dados ilimitados (pipelines de streaming), os dados chegam continuamente e são ingeridos como mensagens ou partições de stream.
  • Pacotes:o Dataflow agrupa elementos em pacotes arbitrários para processamento por um DoFn. Um pacote é a unidade de falha e nova tentativa: se o processamento de um elemento gerar uma exceção não processada, todo o pacote será tentado novamente. Operações com alto consumo de memória podem aumentar a pressão na memória do worker e causar erros de falta de memória.
  • Estágios e fusão de etapas:durante a otimização do gráfico, o Dataflow combina transformações adjacentes em estágios de execução fundidos para eliminar a sobrecarga da materialização de dados intermediários. Em um estágio fundido, os elementos são processados em um loop de execução apertado em uma única linha de execução antes de serem transmitidos para o próximo estágio ou limite de embaralhamento.

Para mais detalhes sobre a tradução de pipeline e a geração de gráficos, consulte Ciclo de vida do pipeline.

Paralelismo gerenciado e escalonamento automático

Por padrão, o Dataflow gerencia o paralelismo de pipelines automaticamente sem exigir ajuste manual de partições das seguintes maneiras:

  • Escalonamento automático horizontal:
    • Pipelines em lote:avalia o total estimado de trabalho restante, o backlog de origem e o uso da CPU para aumentar ou diminuir o pool de workers e concluir o job de forma rápida e econômica.
    • Pipelines de streaming:analisam a latência do sistema, o tamanho do backlog e a utilização da CPU para aumentar a escala dos workers durante picos de capacidade de processamento e reduzir escala vertical durante períodos de baixo tráfego. Para mais detalhes, consulte Ajustar o escalonamento automático horizontal de streaming.
  • Rebalanceamento dinâmico de trabalho (DWR, na sigla em inglês): em pipelines em lote, o Dataflow monitora o progresso das tarefas individuais do worker. Se um worker terminar antes do tempo ou outro ficar para trás devido a um desvio de dados (atrasados), o Dataflow vai dividir dinamicamente o trabalho residual não processado do worker lento e atribuí-lo a um worker ocioso. Para mais informações, consulte Rebalanceamento de trabalho dinâmico.
  • Escala dinâmica de linhas de execução:em pipelines em lote que usam o Portable Runner, ajusta automaticamente o número de linhas de execução de processamento simultâneo por worker com base na utilização da CPU e na capacidade de memória. Para mais informações, consulte escala dinâmica de linhas de execução.
  • Escalonamento automático vertical:no Dataflow Prime, o Dataflow escalona dinamicamente a memória do worker e os recursos de computação para evitar erros de falta de memória e otimizar a utilização de recursos. Para mais informações, consulte Escalonamento automático vertical.

Fatores que restringem o paralelismo

Um pipeline pode não atingir o paralelismo esperado devido às seguintes características de dados ou design do gráfico de pipeline:

Origens de entrada não divisíveis

Se uma fonte de entrada não puder ser dividida em intervalos independentes, o Dataflow será forçado a ler a fonte sequencialmente com uma única linha de execução de worker:

  • Compactação de arquivos não divisíveis:formatos como .gz (gzip) ou .bzip2 (sem indexação) não podem ser lidos em paralelo de deslocamentos de bytes arbitrários. A leitura de um único arquivo grande compactado restringe a etapa de ingestão a uma única linha de execução até que os dados sejam descompactados e redistribuídos.
  • Resolução:armazene os dados em formatos de arquivo divisíveis (como Parquet, Avro ou formatos compactados com Snappy) ou divida os dados de entrada em vários arquivos menores no Cloud Storage.

Fusão de etapas e alta distribuição de dados

A fusão de etapas melhora o desempenho reduzindo a sobrecarga de serialização, mas pode limitar inadvertidamente o paralelismo e aumentar a pressão na memória quando uma etapa com baixo paralelismo produz um grande número de elementos de saída (uma operação de "distribuição de dados" alta):

  • Exemplo:uma origem lê cinco arquivos e é combinada com uma transformação FlatMap que produz 1.000.000 de elementos de saída. Se a transformação FlatMap for fundida com transformações downstream, todos os 1.000.000 de elementos continuarão sendo executados em no máximo cinco linhas de execução de worker, limitando muito a capacidade de processamento downstream. Além disso, se as transformações intermediárias se expandirem significativamente na memória antes do commit, pacotes grandes poderão esgotar a memória disponível do worker.
  • Solução:insira uma transformação Redistribute (ou Reshuffle clássica) entre a etapa de alto fan-out e as transformações downstream para interromper a fusão e redistribuir o trabalho no pool de trabalhadores. Para depurar problemas relacionados à memória, consulte Resolver erros de falta de memória.

Inclinação de chave e teclas de atalho

As operações de agregação (GroupByKey, CoGroupByKey, Combine.PerKey) agrupam elementos pela chave associada.

  • Gargalo de chave quente:o Dataflow encaminha todos os elementos com a mesma chave para uma única linha de execução de worker para agregação. Se uma única chave contiver uma grande porcentagem do conjunto de dados total, esse worker se tornará um atrasado, e os workers upstream poderão sofrer contrapressão. Por exemplo, uma chave null padrão ou uma chave de categoria extremamente popular.
  • Resolução:
    1. Use combinadores (CombineFn ou Combine.PerKey) em vez de GroupByKey sempre que possível, permitindo que o Dataflow faça combinações locais parciais antes do shuffle.
    2. Adicione um prefixo ou sufixo inteiro aleatório às teclas de atalho (salga de chaves) para distribuir o espaço de chaves entre os workers, seguido por uma agregação de segunda etapa para mesclar os resultados salgados.

Limitação de coletor downstream

Ao gravar a saída do pipeline em serviços externos, como bancos de dados ou APIs de terceiros, o alto paralelismo pode saturar o sistema de destino:

  • Limitação:centenas de linhas de execução de worker emitindo chamadas de gravação simultâneas podem causar erros de limitação de taxa, tempos limite de conexão ou degradação do banco de dados.
  • Resolução:
    • Limite o paralelismo de gravação agrupando elementos com GroupByKey ou usando gravadores em lote com paralelismo controlado.
    • Implemente a espera exponencial do lado do cliente e a lógica de repetição nas implementações do coletor DoFn.

Estratégias de otimização

Para otimizar o paralelismo nos jobs do Dataflow, considere as seguintes abordagens:

A seguir