Projetar para flexibilidade e eficiência no Dataflow

Este documento explica o design para flexibilidade e eficiência (DFE, na sigla em inglês), um conjunto de práticas recomendadas de arquitetura para criar pipelines do Dataflow resilientes.

A transição de restrições rígidas de infraestrutura para definições flexíveis de recursos ajuda você a:

  • Maximizar a disponibilidade de recursos de computação.
  • Garantir o escalonamento automático contínuo durante períodos de alta demanda regional.
  • Evite atrasos no lançamento do pipeline e elimine gargalos de capacidade.

Por exemplo, em vez de restringir o pipeline a um tipo de máquina específico em uma zona, como exigir workers n1-standard-4 em us-central1-a, você pode definir necessidades mínimas de recursos (como 4 vCPUs e 16 GB de RAM). Se us-central1-a ou a série de máquinas N1 tiver restrições temporárias de capacidade, o Dataflow poderá provisionar automaticamente VMs de worker compatíveis em outras zonas e famílias de máquinas (como E2, N2 ou N2D). Essa flexibilidade ajuda a garantir que seu pipeline seja iniciado e escalonado sem esperar por um único pool de hardware restrito.

Este documento é destinado a engenheiros de dados, arquitetos de nuvem e administradores de plataforma que gerenciam cargas de trabalho do Dataflow e querem otimizar a confiabilidade, a capacidade de processamento e a disponibilidade da infraestrutura do pipeline.

Visão geral do DFE

O Dataflow é um serviço de processamento de dados totalmente gerenciado e sem servidor que provisiona dinamicamente instâncias de máquina virtual (VM) do Compute Engine para executar pipelines do Apache Beam. Em pipelines de streaming e processamento em lote de grande escala, os pools de workers geralmente escalonar verticalmente para dezenas ou centenas de instâncias de VM.

Pipelines configurados com restrições de infraestrutura rígidas estão sujeitos a atrasos no provisionamento durante períodos de alta demanda. Exemplos de restrições rígidas:

  • Codificar um único tipo de máquina, como n1-standard-4.
  • Fixar o pipeline em uma zona específica do Compute Engine.

Se esse tipo de máquina ou zona específica tiver uma alta demanda temporária, o Dataflow não poderá alocar recursos de computação. Isso pode causar atrasos no provisionamento ou erros como ZONE_RESOURCE_POOL_EXHAUSTED ou RESOURCE_POOL_EXHAUSTED.

Usar os princípios do DFE ajuda a mudar a arquitetura do pipeline de declarações de infraestrutura rígidas e estáticas para definições de recursos flexíveis e baseadas em requisitos. Essa flexibilidade permite que o Dataflow distribua dinamicamente a computação em diversos pools de hardware disponíveis no Google Cloud, ajudando você a maximizar a capacidade de computação e minimizar a sobrecarga operacional.

Práticas recomendadas para o DFE

Adote as práticas recomendadas a seguir para maximizar a capacidade de obtenção de computação, melhorar a capacidade de resposta do escalonamento automático e criar pipelines resilientes.

Ativar a seleção automática de VMs

Em vez de codificar um tipo de máquina estático com a opção de pipeline tipo de máquina do worker, use a seleção automática de VM com dicas de recursos do Apache Beam. Quando você especifica requisitos mínimos de recursos (min_ram ou cpu_count), o Dataflow ativa automaticamente a flexibilidade de instâncias e provisiona workers de uma lista de tipos de máquinas compatíveis.

Suporte a cargas de trabalho:

  • Pipelines em lote:o ajuste correto e a seleção automática de VM são ativados automaticamente quando você especifica dicas de recursos.
  • Pipelines de streaming:o ajuste adequado exige a definição da opção de pipeline --experiments=enable_streaming_rightfitting, além do escalonamento automático horizontal (ativado por padrão) e do Streaming Engine (--enable_streaming_engine).

Para configurar a seleção automática de VMs, especifique os requisitos mínimos de recursos (min_ram ou cpu_count) no nível do pipeline usando opções de linha de comando, opções de pipeline do SDK ou parâmetros de execução de modelo flexível. Para instruções detalhadas de configuração e exemplos de código em Java e Python, consulte Usar dicas de recursos.

Use o posicionamento regional de workers (evite a fixação zonal)

Configure o Dataflow para programar dinamicamente VMs de worker em qualquer zona íntegra na região escolhida.

Especifique a opção de pipeline --region e omita --zone e --worker_zone. Exemplo:

--region=us-central1

Desacoplar estado e embaralhamento usando serviços gerenciados

Os pipelines que não usam serviços de back-end gerenciados executam operações de dados de redistribuição e armazenamento de estado de streaming diretamente nos discos e na memória da VM de worker. Esse acoplamento estreito exige discos de worker maiores e vincula a sobrevivência da carga de trabalho a instâncias de VM específicas, dificultando a substituição do worker durante restrições de capacidade.

  • Para jobs em lote, use o Dataflow Shuffle: o Dataflow Shuffle é ativado por padrão para pipelines em lote executados em tipos de máquinas worker compatíveis e descarrega operações de embaralhamento das VMs de worker para um serviço de back-end dedicado gerenciado pelo Google.
  • Para jobs de streaming, use o Streaming Engine: o Streaming Engine descarrega o armazenamento de estado da janela e o gerenciamento de timers das VMs de worker para uma infraestrutura de back-end especializada e altamente responsiva. Para pipelines que usam o SDK do Apache Beam 2.30.0 ou posterior, o Streaming Engine é ativado por padrão. Para ativar explicitamente, transmita a opção de pipeline --enable_streaming_engine.

Usar a programação flexível de recursos (FlexRS) para pipelines em lote

Para cargas de trabalho em lote não urgentes, como ETL noturno, ingestão de data lake ou rollups diários, use a Programação flexível de recursos (FlexRS).

Para ativar a FlexRS, defina a opção de pipeline de meta da FlexRS:

  • Para pipelines do Python: --flexrs_goal=COST_OPTIMIZED
  • Para pipelines Java: --flexRSGoal=COST_OPTIMIZED

Configurar tipos de VM de inicialização flexível para modelos Flex

Ao iniciar pipelines usando modelos flexíveis, a VM do inicializador de pipeline usa e2-standard-2 por padrão. A VM padrão funciona na maioria dos casos, mas, se você tiver restrições de capacidade, personalize a configuração usando a opção --launcher-machine-type ao executar o comando gcloud dataflow flex-template run:

gcloud dataflow flex-template run my-job \
    --template-file-gcs-location="gs://my-bucket/template.json" \
    --region="us-central1" \
    --launcher-machine-type="n2-standard-2"

Considerações e compensações operacionais

Embora a adoção das práticas recomendadas do DFE melhore significativamente a capacidade de obtenção de computação, a capacidade de resposta do escalonamento automático e a confiabilidade operacional, considere os seguintes fatores operacionais e compensações ao projetar sua arquitetura:

Considerações sobre a seleção automática de VMs

  • Confiabilidade x desempenho máximo:a seleção automática de VMs prioriza a confiabilidade do lançamento de jobs e a capacidade de computação em vez do desempenho máximo de execução. Como o Dataflow provisiona de várias famílias de máquinas candidatas (como E2, N2, N4 e N2D), o desempenho e a capacidade de processamento do tempo de execução podem variar um pouco dependendo da família de máquinas provisionada. Para cargas de trabalho com uso intenso de computação e SLAs de execução rigorosos, teste seu pipeline com a seleção automática de VM para estabelecer um comparativo de performance antes de implantá-lo de forma geral. Se uma carga de trabalho exigir uma plataforma de hardware ou velocidade do clock específica e você puder tolerar restrições de capacidade, continue definindo um tipo de máquina específico.
  • Cota do Compute Engine em famílias candidatas:como a seleção automática de VM pode provisionar workers de várias famílias de máquinas candidatas, verifique se o projeto Google Cloud tem cota suficiente de vCPU e memória do Compute Engine para cada família candidata na região de destino. Se houver uma escassez de capacidade na família principal e seu projeto não tiver cota para a família de substituição, o provisionamento de workers vai falhar com um erro QUOTA_EXCEEDED.
  • Pré-requisito do pipeline de streaming:para pipelines de streaming, o ajuste adequado e a seleção automática de VM não são ativados por padrão. É preciso especificar explicitamente --experiments=enable_streaming_rightfitting e verificar se o Streaming Engine (--enable_streaming_engine) e o escalonamento automático horizontal estão ativos.
  • Exclusões de configuração:a seleção automática de VM é ignorada ou indisponível se você configurar qualquer um dos recursos ou opções na tabela a seguir:

    Recurso Opção de flag ou configuração Observações
    Tipos de máquina explícitos --worker_machine_type ou --machine_type (Python)
    --workerMachineType (Java)
    A seleção automática de VM é ignorada em favor do tipo de máquina especificado.
    Tipos de disco personalizados, IOPS provisionadas ou capacidade de processamento --disk_type, --disk_provisioned_iops ou --disk_provisioned_throughput_mibps A seleção automática de VM é ignorada. É possível definir um tamanho de disco personalizado com --disk_size_gb.
    Plataformas mínimas de CPU --min_cpu_platform (Python)
    --minCpuPlatform (Java)
    Definir uma plataforma de CPU mínima ignora a seleção automática de VM.
    VM confidencial --experiments=enable_confidential_compute As instâncias de VM confidenciais não são compatíveis com a seleção automática de VM.
    Aceleradores de GPU ou TPU Dica de recurso --dataflow_service_options=worker_accelerator=... ou accelerator A seleção automática de VM se aplica apenas a cargas de trabalho sem aceleradores.
    Dataflow Prime --dataflow_service_options=enable_prime O Dataflow Prime usa o escalonamento automático vertical e o ajuste dinâmico em vez da seleção automática de VM.
    Programação flexível de recursos (FlexRS, na sigla em inglês) --flexrs_goal=COST_OPTIMIZED (Python)
    --flexRSGoal=COST_OPTIMIZED (Java)
    A FlexRS gerencia o próprio pool de workers e buffer de programação.

Compensações da programação flexível de recursos (FlexRS)

  • Janela de atraso no agendamento:o FlexRS pode introduzir um buffer de agendamento de até 6 horas antes do início da execução do job. Não use a FlexRS para pipelines com SLAs de tempo até a conclusão rigorosos ou dependências downstream apertadas.

Posicionamento regional e localidade de dados

  • Pré-requisito do serviço gerenciado:a veiculação regional de workers é compatível apenas com jobs que usam o Dataflow Shuffle para lote ou o Streaming Engine para streaming. Os jobs que não usam esses serviços de back-end gerenciados usam o posicionamento automático de zona, que seleciona uma única melhor zona na região.
  • Localidade de dados e saída entre regiões:o posicionamento regional distribui workers nas zonas disponíveis na região escolhida. Para minimizar a latência da rede e evitar cobranças de saída de rede entre regiões, verifique se todas as origens e destinos de dados (como buckets do Cloud Storage, conjuntos de dados do BigQuery e tópicos do Pub/Sub) estão na mesma região do job do Dataflow.

Reservas do Compute Engine

  • Afinidade de reserva:os jobs do Dataflow sob demanda consomem automaticamente as reservas correspondentes do Compute Engine que usam a afinidade de reserva ANY. No entanto, a seleção automática de VM não oferece suporte ao consumo de instâncias de reservas específicas nomeadas.
  • Adequação para cargas de trabalho temporárias:as reservas do Compute Engine geralmente não são recomendadas para cargas de trabalho em lote de curta duração, com escalonamento automático ou propensas a picos. Além disso, a criação de novas reservas durante uma escassez zonal ativa falha com as mesmas restrições de capacidade da criação de VMs sob demanda.

A seguir