Nesta página, descrevemos as práticas recomendadas para desenvolver pipelines do Dataflow. O uso dessas práticas recomendadas tem os seguintes benefícios:
- Melhore a observabilidade e o desempenho do pipeline
- Melhora na produtividade dos desenvolvedores
- Aprimorar a capacidade de teste de pipelines
Os exemplos de código do Apache Beam nesta página usam Java, mas o conteúdo se aplica aos SDKs do Apache Beam Java, Python e Go.
Perguntas a serem consideradas
Ao projetar o pipeline, considere as seguintes perguntas:
- Onde os dados de entrada do pipeline são armazenados? Quantos conjuntos de dados de entrada existem?
- Como são seus dados?
- O que você quer fazer com os dados?
- Para onde devem ir os dados de saída do pipeline?
- Seu job do Dataflow usa o Assured Workloads?
Usar modelos
Para acelerar o desenvolvimento do pipeline, em vez de criar um pipeline escrevendo um código do Apache Beam, use um modelo do Dataflow quando possível. Os modelos têm os seguintes benefícios:
- Os modelos são reutilizáveis.
- Os modelos permitem personalizar cada job mudando parâmetros de pipeline específicos.
- Assim, qualquer pessoa a quem você fornecer permissões poderá usar o modelo para implantar o pipeline. Por exemplo, um desenvolvedor pode criar um job de um modelo, e um cientista de dados na organização pode implantar esse modelo posteriormente.
É possível usar um modelo fornecido pelo Google ou criar seu próprio modelo. Alguns modelos fornecidos pelo Google permitem adicionar lógica personalizada como uma etapa do pipeline. Por exemplo, a assinatura do Pub/Sub para o modelo do BigQuery fornece um parâmetro para executar uma função JavaScript definida pelo usuário (UDF, na sigla em inglês) que é armazenada no Cloud Storage.
Os modelos fornecidos pelo Google são de código aberto sob a licença Apache 2.0, para que você possa usá-los como base para novos pipelines. Os modelos também são úteis como exemplos de código. Veja o código do modelo no repositório do GitHub (link em inglês).
Assured Workloads
O Assured Workloads ajuda a aplicar os requisitos de segurança e compliance para clientes doGoogle Cloud . Por exemplo, Regiões e suporte da UE com Controles de Soberania ajudam a aplicar residência de dados e garantias de soberania de dados para clientes na UE. Para fornecer esses recursos do Dataflow, alguns deles são restritos ou limitados. Se você usa o Assured Workloads com o Dataflow, todos os recursos acessados pelo pipeline precisam estar no projeto ou pasta do Assured Workloads da organização. Esses recursos incluem:
- Buckets do Cloud Storage
- Conjuntos de dados do BigQuery
- Tópicos e assinaturas do Pub/Sub
- Conjuntos de dados do Firestore
- Conectores de E/S
No Dataflow, para os jobs de streaming criados após 7 de março de 2024, todos os dados do usuário são criptografados com CMEK.
Para jobs de streaming criados antes de 7 de março de 2024, as chaves de dados usadas em operações baseadas em chaves, como janelamento, agrupamento e mesclagem, não são protegidas pela criptografia CMEK. Para ativar essa criptografia nos jobs, drene ou cancele o job e reinicie-o. Para mais informações, consulte Criptografia de artefatos de estado do pipeline.
Compartilhar dados entre pipelines
Não há nenhum mecanismo de comunicação cruzada de pipeline específico do Dataflow para compartilhamento de dados ou processamento de contexto entre pipelines. Use um armazenamento durável, como o Cloud Storage, ou um cache na memória, como o App Engine, para compartilhar dados entre instâncias de pipelines.
Programar jobs
É possível automatizar a execução do pipeline das seguintes maneiras:
- Use os pipelines de dados do Dataflow para criar programações de jobs recorrentes e entender onde os recursos são gastos em várias execuções de jobs.
- Usar o Cloud Scheduler.
- Use o operador Dataflow do Apache Airflow, uma das várias opções de operadores doGoogle Cloud em um fluxo de trabalho do Serviço gerenciado para Apache Airflow.
- Executar processos de (cron) job personalizados no Compute Engine
Práticas recomendadas para escrever código de pipeline
As seções a seguir fornecem práticas recomendadas para criar pipelines escrevendo o código do Apache Beam.
Estruturar seu código do Apache Beam
Para criar pipelines, é comum usar a transformação genérica do Apache Beam de processamento paralelo ParDo.
Ao aplicar uma transformação ParDo, você fornece o código de usuário na forma de um
objeto DoFn. DoFn é uma classe do SDK do Apache Beam que define uma função
de processamento distribuído.
Pense no código DoFn como pequenas entidades independentes: várias instâncias
podem estar em execução em diferentes máquinas, sem que
haja conhecimento umas das outras. Dessa forma, recomendamos a criação de funções puras, que
são ideais para a natureza paralela e distribuída dos elementos DoFn.
Funções puras têm as seguintes características:
- As funções puras não dependem do estado oculto ou externo.
- Eles não têm efeitos colaterais observáveis.
- Eles são determinísticos.
O modelo de função pura não é estritamente rígido. Quando seu código não depende de
coisas que não são garantidas pelo serviço Dataflow, as informações
de estado ou os dados de inicialização externos podem ser válidos para DoFn e outros
objetos de função.
Ao estruturar as transformações ParDo e criar os elementos DoFn,
considere as seguintes diretrizes:
- Quando você usa o processamento único, o serviço Dataflow processa todos os elementos da entrada
PCollectionpor uma instânciaDoFnapenas uma vez. - O serviço Dataflow não garante quantas vezes uma
DoFné invocada. - O serviço Dataflow não garante exatamente como os elementos distribuídos são agrupados. Isso não garante quais elementos, se houver, serão processados juntos.
- O serviço do Dataflow não garante o número exato de
instâncias de
DoFncriadas ao longo de um pipeline. - O serviço Dataflow é tolerante a falhas e pode repetir o código várias vezes se os workers encontrarem problemas.
- O serviço Dataflow pode criar cópias de backup do seu código. Podem ocorrer problemas com efeitos colaterais manuais, por exemplo, se o código depender de arquivos temporários com nomes não exclusivos ou criá-los.
- O serviço do Dataflow serializa o processamento de elemento por
instância de
DoFn. Seu código não precisa ser estritamente seguro para linhas de execução, mas qualquer estado compartilhado entre várias instâncias deDoFnprecisa ser seguro para linhas de execução.
Criar bibliotecas de transformações reutilizáveis
O modelo de programação do Apache Beam permite reutilizar transformações. Ao criar uma biblioteca compartilhada de transformações comuns, é possível melhorar a reutilização, a testabilidade e a propriedade do código por diferentes equipes.
Considere os dois exemplos de código Java a seguir, que leem eventos de pagamento. Supondo que ambos os pipelines executem o mesmo processamento, eles podem usar as mesmas transformações de uma biblioteca compartilhada para as etapas de processamento restantes.
O primeiro exemplo é de uma fonte ilimitada do Pub/Sub:
PipelineOptions options = PipelineOptionsFactory.create();
Pipeline p = Pipeline.create(options)
// Initial read transform
PCollection payments =
p.apply("Read from topic",
PubSubIO.readStrings().withTimestampAttribute(...).fromTopic(...))
.apply("Parse strings into payment events",
ParDo.of(new ParsePaymentEventFn()));
O segundo exemplo é proveniente de uma fonte limitada de banco de dados relacional:
PipelineOptions options = PipelineOptionsFactory.create();
Pipeline p = Pipeline.create(options);
PCollection payments =
p.apply(
"Read from database table",
JdbcIO.<PaymentEvent>read()
.withDataSourceConfiguration(...)
.withQuery(...)
.withRowMapper(new RowMapper () {
...
}));
A implementação das práticas recomendadas de reutilização de código varia de acordo com a linguagem de programação e a ferramenta de build. Por exemplo, se você usa o Maven, é possível separar o código de transformação no próprio módulo. Em seguida, você pode incluir o módulo como um submódulo em projetos com vários módulos maiores para pipelines diferentes, conforme mostrado no exemplo de código a seguir:
// Reuse transforms across both pipelines
payments
.apply("ValidatePayments", new PaymentTransforms.ValidatePayments(...))
.apply("ProcessPayments", new PaymentTransforms.ProcessPayments(...))
...
Para mais informações, consulte as seguintes páginas de documentação do Apache Beam:
- Requisitos para escrever código de usuário para transformações do Apache Beam
- Guia de estilo
PTransform(link em inglês): um guia de estilo para escritores de novas coleçõesPTransformreutilizáveis.
Usar filas de mensagens inativas para tratamento de erros
Às vezes, o pipeline não processa elementos. Problemas de dados são uma causa comum. Por exemplo, um elemento que contém JSON mal formatado pode causar falhas de análise.
Embora você possa capturar exceções no método
DoFn.ProcessElement,
registrar o erro e descartar o elemento, essa abordagem perde os dados
e impede que sejam inspecionados posteriormente para manuseio manual ou solução de problemas.
Em vez disso, use um padrão conhecido como fila de mensagens inativas (fila de mensagens não processadas).
Capture exceções no método DoFn.ProcessElement e registre
erros. Em vez de remover o elemento com falha,
use saídas de ramificação para gravá-lo em um objeto PCollection
separado. Esses elementos são gravados em um coletor de dados para inspeção e processamento posteriores com uma transformação separada.
O exemplo de código Java a seguir mostra como implementar o padrão de fila de mensagens inativas.
TupleTag successTag = new TupleTag<>() {};
TupleTag deadLetterTag = new TupleTag<>() {};
PCollection input = /* ... */;
PCollectionTuple outputTuple =
input.apply(ParDo.of(new DoFn, Output>() {
@Override
void processElement(ProcessContext c) {
try {
c.output(process(c.element()));
} catch (Exception e) {
LOG.severe("Failed to process input {} -- adding to dead-letter file",
c.element(), e);
c.sideOutput(deadLetterTag, c.element());
}
}).withOutputTags(successTag, TupleTagList.of(deadLetterTag)));
// Write the dead-letter inputs to a BigQuery table for later analysis
outputTuple.get(deadLetterTag)
.apply(BigQueryIO.write(...));
// Retrieve the successful elements...
PCollection success = outputTuple.get(successTag);
// and continue processing ...
É possível usar o Cloud Monitoring para aplicar diferentes políticas de monitoramento e alerta na fila de mensagens inativas do pipeline. Por exemplo, é possível visualizar o número e o tamanho dos elementos processados pela transformação de mensagens inativas e configurar os alertas para serem acionados se determinadas condições de limite forem atendidas.
Gerenciar mutações de esquema
É possível processar dados que tenham esquemas inesperados (mas válidos) usando um padrão de mensagens inativas,
que grava elementos com falha em um objeto PCollection separado.
Em alguns casos, você quer processar automaticamente os elementos
que refletem um esquema modificado como elementos válidos. Por exemplo, se o esquema
de um elemento refletir uma mutação, como a adição de novos campos, será possível adaptar o
esquema do coletor de dados para acomodar mutações.
A mutação automática de esquema depende da abordagem de ramificação e da saída usada pelo padrão de mensagens inativas. No entanto, nesse caso, ele aciona uma transformação que modifica o esquema de destino sempre que esquemas aditivos forem encontrados. Para ver um exemplo dessa abordagem, consulte Como processar a modificação de esquemas JSON em um pipeline de streaming, com a Square Enix no blog do Google Cloud .
Decidir como mesclar conjuntos de dados
A junção de conjuntos de dados é um caso de uso comum para pipelines de dados. É possível usar
entradas secundárias ou a transformação CoGroupByKey para realizar mesclagens no pipeline.
Cada uma tem vantagens e desvantagens.
As entradas secundárias
fornecem uma maneira flexível de resolver problemas comuns de processamento de dados, como
enriquecimento de dados e pesquisas com chaves. Ao contrário dos objetos PCollection, as entradas secundárias são
mutáveis e podem ser determinadas no momento da execução. Por exemplo, os valores em uma
entrada secundária podem ser calculados por outra ramificação no pipeline ou determinados
chamando um serviço remoto.
O Dataflow é compatível com entradas secundárias usando a persistência de dados no armazenamento permanente (semelhante a um disco compartilhado), o que disponibiliza a entrada secundária completa para todos os workers.
Os tamanhos de entrada secundária podem ser muito grandes e não caberem na memória do worker. A leitura de uma entrada secundária grande pode causar problemas de desempenho se os workers precisarem ler constantemente do armazenamento permanente.
A transformação
CoGroupByKey
é uma
transformação principal do Apache Beam
que combina (nivela) vários objetos PCollection e elementos de grupos que
tenham uma chave comum. Diferente de uma entrada secundária, que disponibiliza todos os dados de entrada secundária
para cada worker, a CoGroupByKey executa uma operação de embaralhamento (agrupamento)
para distribuir dados entre os workers. Portanto, CoGroupByKey é ideal quando os
objetos PCollection que você quer mesclar são muito grandes e não se encaixam na memória
do worker.
Siga estas diretrizes para decidir se quer usar entradas secundárias ou
CoGroupByKey:
- Use entradas secundárias quando um dos objetos
PCollectionque você está mesclando for desproporcionalmente menor que o outro e o objetoPCollectionmenor se encaixar na memória do worker. Armazenar toda a entrada secundária em cache na memória torna a busca de elementos rápida e eficiente. - Use entradas secundárias quando tiver um objeto
PCollectionque deva ser unido várias vezes no pipeline. Em vez de usar várias transformaçõesCoGroupByKey, crie uma única entrada secundária que possa ser reutilizada por várias transformaçõesParDo. - Use
CoGroupByKeyse você precisar buscar uma grande proporção de um objetoPCollectionque exceda significativamente a memória do worker.
Para mais informações, consulte Resolver problemas de falta de memória no Dataflow.
Minimizar operações caras por elemento
Uma instância DoFn processa lotes de elementos chamados
pacotes,
(unidades atômicas de trabalho que consistem em zero ou mais
elementos). Os elementos individuais são processados pelo método
DoFn.ProcessElement,
que é executado para cada elemento. Como o método DoFn.ProcessElement
é chamado para cada elemento, qualquer operação demorada ou de computação
cara invocada por esse método
é executada para cada elemento processado.
Se você precisar executar operações caras apenas uma vez para um lote de elementos,
inclua-as nos métodos DoFn.Setup ou DoFn.StartBundle,
em vez de no elemento DoFn.ProcessElement. Os exemplos incluem as
operações a seguir:
Análise de um arquivo de configuração que controla algum aspecto do comportamento da instância
DoFn. Invoque essa ação apenas uma vez, quando a instânciaDoFnfor inicializada usando o métodoDoFn.Setup.Ao instanciar um cliente de curta duração que é reutilizado em todos os elementos de um pacote, por exemplo, quando todos os elementos do pacote são enviados por uma única conexão de rede. Invoque essa ação uma vez por pacote usando o método
DoFn.StartBundle.
Limitar os tamanhos de lotes e as chamadas simultâneas a serviços externos
Ao chamar serviços externos, você pode reduzir as sobrecargas por chamada usando a
transformação
GroupIntoBatches. Esta transformação cria lotes de elementos de um tamanho especificado.
O envio em lote envia
elementos para um serviço externo como um payload em vez de
individualmente.
Em combinação com lotes, é possível limitar o número máximo de chamadas paralelas (simultâneas) ao serviço externo escolhendo as chaves apropriadas para particionar os dados recebidos. O número de partições determina o carregamento em paralelo máximo. Por exemplo, se cada elemento receber a mesma chave, uma transformação downstream para chamar o serviço externo não será executada em paralelo.
Considere uma das seguintes abordagens para produzir chaves para elementos:
- Escolha um atributo do conjunto de dados para usar como chaves de dados, como IDs de usuário.
- Gere chaves de dados para dividir elementos de maneira aleatória em um número fixo de
partições, em que o número possível de chaves-valor determina o número
de partições. Você precisa criar partições suficientes para paralelismo.
Cada partição precisa ter elementos suficientes para que a transformação
GroupIntoBatchesseja útil.
No exemplo de código Java a seguir, mostramos como dividir aleatoriamente elementos em 10 partições:
// PII or classified data which needs redaction.
PCollection sensitiveData = ...;
int numPartitions = 10; // Number of parallel batches to create.
PCollection, Iterable >> batchedData =
sensitiveData
.apply("Assign data into partitions",
ParDo.of(new DoFn, KV, String>>() {
Random random = new Random();
@ProcessElement
public void assignRandomPartition(ProcessContext context) {
context.output(
KV.of(randomPartitionNumber(), context.element()));
}
private static int randomPartitionNumber() {
return random.nextInt(numPartitions);
}
}))
.apply("Create batches of sensitive data",
GroupIntoBatches.<Long, String>ofSize(100L));
// Use batched sensitive data to fully utilize Redaction API,
// which has a rate limit but allows large payloads.
batchedData
.apply("Call Redaction API in batches", callRedactionApiOnBatch());
Processar operações lentas por elemento simultaneamente
Quando o DoFn realiza operações lentas por elemento (como chamar um modelo de aprendizado de máquina externo ou uma API da Web), a execução síncrona dentro do DoFn.ProcessElement bloqueia as linhas de execução de trabalho até que cada operação seja concluída.
Esse bloqueio limita a capacidade de processamento e pode causar tempos limite de pacote ou altas taxas de
repetição.
Mesmo com o particionamento ideal, a serialização por chave garante que apenas um elemento seja processado por vez. No entanto, ao encapsular sua classe DoFn
com AsyncWrapper, essa restrição é violada porque os elementos são processados
simultaneamente em um pool de linhas de execução em segundo plano, e o controle é retornado ao executor
principal antes que o processamento seja concluído. Usando o estado e os timers do Beam, o AsyncWrapper
preserva o processamento único.
Ao projetar um pipeline que usa AsyncWrapper, siga estas diretrizes:
- Particione os dados recebidos com chaves: o
AsyncWrapperexige entrada de chave-valorKV. Para distribuir o trabalho entre vários workers e evitar gargalos de serialização, particione o conjunto de dados recebido em um número adequado de chaves distintas. - Limitar a simultaneidade de workers: use o parâmetro
parallelismpara definir o número máximo de operações simultâneas permitidas por nó de trabalho. - Garantir a segurança de encadeamento: se você definir
parallelismcomo maior que 1 e seuDoFnfor com estado, seuDoFntambém precisará ser thread-safe. - Evite a agregação no nível do pacote: não use lógica de agrupamento ou agregação
em
StartBundleouFinishBundle, porqueAsyncWrapperinvoca esses métodos de ciclo de vida por elemento. - Não use saídas com tags: saídas múltiplas com tags (tags de saída adicionais) não são compatíveis.
- Ajuste da execução (Java): por padrão, as tarefas são executadas na JVM compartilhada
ForkJoinPool. DefinauseThreadPool=truepara isolar a execução em um pool de linhas de execução fixo dedicado dimensionado para o paralelismo configurado. Essa abordagem evita conflitos com outros consumidores deForkJoinPoolna mesma JVM. - Execução de ajuste (Python): no SDK do Python,
AsyncWrapperoferece suporte a um modouse_asyncio=Trueadicional que executa corrotinas em um loop de eventosasyncioem vez de um pool de threads em segundo plano.
Os exemplos de código Java a seguir mostram como encapsular um DoFn com AsyncWrapper:
import org.apache.beam.sdk.transforms.AsyncWrapper;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.transforms.WithKeys;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
import org.joda.time.Duration;
// Define a synchronous DoFn that makes slow external service calls.
public class CallSlowExternalServiceFn extends DoFn, String> {
private transient ExternalServiceClient externalServiceClient;
@Setup
public void setup() {
externalServiceClient = new ExternalServiceClient();
}
@ProcessElement
public void processElement(@Element String element, OutputReceiver receiver) {
String result = externalServiceClient.call(element);
receiver.output(result);
}
}
// Partition incoming data across distinct keys to distribute worker load.
PCollection, String>> keyedData = inputData.apply(
WithKeys.of(element -> element.substring(0, 1))); // or your actual partition key
// Wrap the DoFn with AsyncWrapper.
PCollection results = keyedData.apply(
"Process Asynchronously",
ParDo.of(
new AsyncWrapper, String, String>(
new CallSlowExternalServiceFn(),
/* parallelism= */ 20,
/* timerFrequency= */ Duration.standardSeconds(5),
/* maxItemsToBuffer= */ 50,
/* timeout= */ Duration.standardSeconds(1),
/* maxWaitTime= */ Duration.millis(500),
/* idFn= */ null, // defaults to element identity for deduplication
/* useThreadPool= */ true)));
Identificar problemas de desempenho causados por etapas com a combinação inadequada
O Dataflow cria um gráfico de etapas que representa o pipeline, com base nas transformações e nos dados usados para criá-lo. Isso é chamado de gráfico de execução de pipeline.
Quando você implanta o pipeline, o Dataflow pode modificar
o gráfico de execução do pipeline para melhorar o desempenho. Por exemplo, o Dataflow
pode unir algumas operações, um processo conhecido como
otimização de fusão,
para evitar o impacto de desempenho e custo de cada gravação de objeto intermediário
PCollection no pipeline.
Em alguns casos, o Dataflow pode determinar incorretamente a melhor forma de fundir operações no pipeline, o que pode limitar a capacidade do serviço do Dataflow de usar todos os workers disponíveis. Nesses casos, é possível impedir a fusão das operações.
Considere o exemplo de código do Apache Beam a seguir. Uma transformação
GenerateSequence
cria um pequeno objeto PCollection limitado, que depois será processado
por duas transformações ParDo downstream.
A transformação Find Primes Less-than-N pode ser cara em termos de computação e provavelmente
será executada lentamente para números grandes. Em contraste, a transformação Increment Number provavelmente é concluída rapidamente.
import com.google.common.math.LongMath;
...
public class FusedStepsPipeline {
final class FindLowerPrimesFn extends DoFn, String> {
@ProcessElement
public void processElement(ProcessContext c) {
Long n = c.element();
if (n > 1) {
for (long i = 2; i < n; i++) {
if (LongMath.isPrime(i)) {
c.output(Long.toString(i));
}
}
}
}
}
public static void main(String[] args) {
Pipeline p = Pipeline.create(options);
PCollection sequence = p.apply("Generate Sequence",
GenerateSequence
.from(0)
.to(1000000));
// Pipeline branch 1
sequence.apply("Find Primes Less-than-N",
ParDo.of(new FindLowerPrimesFn()));
// Pipeline branch 2
sequence.apply("Increment Number",
MapElements.via(new SimpleFunction, Long>() {
public Long apply(Long n) {
return ++n;
}
}));
p.run().waitUntilFinish();
}
}
O diagrama a seguir mostra uma representação gráfica do pipeline na interface de monitoramento do Dataflow.

A interface de monitoramento do Dataflow mostra que a mesma taxa lenta de processamento ocorre para as duas transformações, especificamente 13 elementos por segundo. Espera-se que a transformação Increment Number processe
elementos rapidamente, mas parece que ela está vinculada à mesma taxa de
processamento que Find Primes Less-than-N.
O motivo é que o Dataflow mesclava as etapas em uma única
fase, o que as impedia de serem executadas de maneira independente. Use o comando gcloud dataflow jobs describe para encontrar mais informações:
gcloud dataflow jobs describe --full job-id --format json
Na saída resultante, as etapas combinadas são descritas no objeto
ExecutionStageSummary
na matriz
ComponentTransform
:
...
"executionPipelineStage": [
{
"componentSource": [
...
],
"componentTransform": [
{
"name": "s1",
"originalTransform": "Generate Sequence/Read(BoundedCountingSource)",
"userName": "Generate Sequence/Read(BoundedCountingSource)"
},
{
"name": "s2",
"originalTransform": "Find Primes Less-than-N",
"userName": "Find Primes Less-than-N"
},
{
"name": "s3",
"originalTransform": "Increment Number/Map",
"userName": "Increment Number/Map"
}
],
"id": "S01",
"kind": "PAR_DO_KIND",
"name": "F0"
}
...
Nesse cenário, a transformação Find Primes Less-than-N é a etapa lenta. Portanto,
a quebra da mesclagem antes dessa etapa é uma estratégia apropriada. Um método para
separar as etapas é inserir uma transformação
GroupByKey
e realizar o desagrupamento antes da etapa, conforme mostrado no exemplo de código Java
a seguir:
sequence
.apply("Map Elements", MapElements.via(new SimpleFunction, KV, Void>>() {
public KV, Void> apply(Long n) {
return KV.of(n, null);
}
}))
.apply("Group By Key", GroupByKey.<Long, Void>create())
.apply("Emit Keys", Keys.<Long>create())
.apply("Find Primes Less-than-N", ParDo.of(new FindLowerPrimesFn()));
Também é possível combinar essas etapas em uma transformação composta reutilizável.
Depois de cancelar a fusão das etapas, ao executar o pipeline, Increment Number é concluído em questão de segundos e a transformação Find Primes Less-than-N, que é muito mais longa, é executada em um estágio separado.
Este exemplo aplica uma operação de agrupamento e desagrupamento para propagar as etapas.
É possível usar outras abordagens para outras circunstâncias. Nesse caso, o processamento
da saída duplicada não é um problema, dada a saída consecutiva da
transformação
GenerateSequence.
Objetos KV com chaves duplicadas são desduplicados para uma única chave na transformação de grupo (GroupByKey) e no desagrupar (Keys). Para reter cópias após as operações de agrupar e desagrupar,
crie pares de chave-valor usando as seguintes etapas:
- Use uma chave aleatória e a entrada original como o valor.
- Agrupar usando a chave aleatória.
- Emita os valores de cada chave como saída.
Você também pode usar uma transformação Reshuffle para evitar a fusão de transformações ao redor. No entanto, os efeitos colaterais da
transformação Reshuffle não são portáteis para diferentes
executores do Apache Beam.
Para mais informações sobre paralelismo e otimização de fusão, consulte Ciclo de vida do pipeline.
Usar métricas do Apache Beam para coletar insights do pipeline
As métricas do Apache Beam são uma classe de utilitário que produz métricas para relatar as propriedades de um pipeline em execução. Quando você usa o Cloud Monitoring, as métricas do Apache Beam ficam disponíveis como métricas personalizadas do Cloud Monitoring.
O exemplo a seguir mostra as métricas Counter do Apache Beam usadas em uma subclasse DoFn.
O código de exemplo usa dois contadores. Um contador rastreia falhas de análise JSON
(malformedCounter) e o outro monitora se a mensagem JSON é
válida, mas contém um payload vazio (emptyCounter). No Cloud Monitoring,
os nomes das métricas personalizadas são custom.googleapis.com/dataflow/malformedJson e
custom.googleapis.com/dataflow/emptyPayload. Use as métricas personalizadas
para criar visualizações e políticas de alertas no Cloud Monitoring.
final TupleTag errorTag = new TupleTag (){};
final TupleTag successTag = new TupleTag (){};
final class ParseEventFn extends DoFn, MyObject> {
private final Counter malformedCounter = Metrics.counter(ParseEventFn.class, "malformedJson");
private final Counter emptyCounter = Metrics.counter(ParseEventFn.class, "emptyPayload");
private Gson gsonParser;
@Setup
public setup() {
gsonParser = new Gson();
}
@ProcessElement
public void processElement(ProcessContext c) {
try {
MyObject myObj = gsonParser.fromJson(c.element(), MyObject.class);
if (myObj.getPayload() != null) {
// Output the element if non-empty payload
c.output(successTag, myObj);
}
else {
// Increment empty payload counter
emptyCounter.inc();
}
}
catch (JsonParseException e) {
// Increment malformed JSON counter
malformedCounter.inc();
// Output the element to dead-letter queue
c.output(errorTag, c.element());
}
}
}
Proteja os buckets do Cloud Storage contra ataques de pipeline
Para proteger os buckets do Cloud Storage contra um ataque de pipeline do Dataflow, é preciso entender como uma violação pode acontecer.
Um invasor raramente ataca o Cloud Storage diretamente. Em vez disso, eles exploram uma vulnerabilidade chamada envenenamento de recursos sombra. Nesse cenário, um invasor compromete o bucket do Cloud Storage que contém seus modelos, metadados ou funções definidas pelo usuário (UDFs) personalizadas do Python do Dataflow para injetar código malicioso. Quando o Dataflow faz escalonamento automático e inicia uma nova VM de worker, ele extrai o código comprometido do Cloud Storage, o executa, rouba o token da conta de serviço do worker e extrai dados dos buckets de origem.
Para ajudar a impedir esse ponto de entrada de ataque, é necessário bloquear os buckets do Cloud Storage que alimentam o Dataflow, proteger os próprios pipelines e limitar estritamente as permissões nos buckets do Cloud Storage.
Proteger buckets do Cloud Storage
Para ajudar a proteger o ponto de injeção, proteja os buckets usados para modelos do Dataflow, arquivos de preparo e dependências de código. Ao implementar controles de acesso de gravação estritos nesses buckets, você impede que códigos não autorizados entrem no ambiente de worker.
- Isole seus buckets:separe os buckets de dados dos buckets de código operacional. O Dataflow usa as flags
--stagingLocatione--tempLocation. Mantenha esses dados em um bucket dedicado e altamente restrito, completamente separado dos dados brutos de entrada ou saída. - Aplicar o acesso uniforme no nível do bucket (UBLA): ative o UBLA em todos os buckets relacionados ao Dataflow. Isso desativa as listas de controle de acesso (ACLs) legadas e propensas a erros e garante que apenas as políticas centralizadas do Identity and Access Management (IAM) determinem quem pode acessar os arquivos.
- Bloquear permissões de gravação:conceda à sua conta de serviço de CI/CD
roles/storage.objectAdmin(gravação/exclusão) nos buckets de preparação e modelo. Conceda aos desenvolvedores e outras contas de serviço, no máximo,roles/storage.objectViewer(somente leitura). - Ative o controle de versões e a retenção de objetos:ative o controle de versões de objeto nos buckets de modelo. Se um usuário malicioso ou uma conta violada substituir um modelo, você vai receber um alerta instantâneo sobre a nova versão e poderá reverter as mudanças. Combinado com uma política de retenção de objetos, isso ajuda a impedir que usuários de má-fé excluam suas dependências de staging.
Reforçar o ambiente de execução do Dataflow
Para manter um ambiente de execução seguro e proteger contra roubo de dados, restrinja a comunicação do worker.
- Bloquear a saída da Internet (rede privada estrita): ao iniciar o Dataflow, transmita
--usePublicIps=false. Além disso, configure as regras de firewall da VPC para bloquear toda a saída para a Internet pública das sub-redes do Dataflow. Se um worker for comprometido, ele não poderá abrir um shell reverso ou exfiltrar dados para um endereço IP externo controlado por um invasor. - Aplicar o VPC Service Controls (VPC-SC): coloque os buckets de dados do Cloud Storage e o projeto do Dataflow em um perímetro de segurança estrito do VPC-SC. A VPC-SC funciona como um firewall de rede que substitui as permissões do IAM. Mesmo que um invasor roube um token de conta de serviço de worker válido, o VPC-SC bloqueia qualquer solicitação para copiar dados do bucket do Cloud Storage para um bucket fora do perímetro da sua organização.
Implementar o privilégio mínimo estrito em buckets de dados
Se um worker do Dataflow for comprometido, a área afetada será definida pelas permissões do IAM da conta de serviço dele.
- Evite usar a conta de serviço padrão do Compute Engine:essa conta tem permissões amplas de editor em todo o projeto por padrão. Use uma conta de serviço do worker dedicada e gerenciada pelo usuário com a flag
--serviceAccount. - Papéis do IAM específicos do bucket:não conceda à conta de serviço do worker do Dataflow papéis de armazenamento no nível do projeto, como
roles/storage.admin. Em vez disso, conceda direitos diretamente nos buckets específicos:- Intervalo de entrada:conceda o papel
roles/storage.objectViewer(somente leitura). - Bucket de saída:conceda a função
roles/storage.objectCreator(somente gravação, ou seja, pode criar objetos, mas não ler nem substituir dados atuais). - Bucket de preparo/temporário:conceda os papéis
roles/storage.objectAdmin(leitura, gravação e exclusão).
- Intervalo de entrada:conceda o papel
Proteja contêineres personalizados e a cadeia de suprimentos de software
Se você usar imagens personalizadas do Docker para suas dependências do Dataflow
usando --sdkContainerImage, um invasor poderá tentar contaminar o registro
de imagens em vez de um bucket do Cloud Storage.
- Verificar o Artifact Registry:armazene suas imagens de worker personalizadas no Artifact Registry e ative a verificação de vulnerabilidades.
Use resumos de imagens imutáveis:em vez de referenciar uma imagem por uma tag mutável, como
us-docker.pkg.dev/my-project/dataflow-worker:latest, especifique o resumo SHA criptográfico exato e imutável nas opções do pipeline:--sdkContainerImage=us-docker.pkg.dev/my-project/dataflow-worker@sha256:7b9c...`Isso ajuda a garantir que o código exato que você verificou e aprovou no CI/CD seja o que é executado nos workers, bloqueando totalmente ataques de substituição de imagem.
Usar o registro para detectar anomalias
As configurações seguras ajudam a interromper o ataque, mas também é recomendável configurar alertas do Cloud Logging para os seguintes gatilhos:
- Eventos no seu modelo do Dataflow ou em buckets de
armazenamento temporário do Cloud Storage originados de uma identidade não CI/CD,
como solicitações
storage.objects.updateoustorage.objects.create. - Grandes volumes de solicitações
storage.objects.getde endereços IP inesperados ou picos incomuns no volume de dados que passam pelo Dataflow, o que pode indicar uma tentativa de exfiltração ativa.
Migrar um pipeline entre projetos Google Cloud
Os jobs do Dataflow estão vinculados ao projeto Google Cloud em que você os criou. Não é possível mover um job diretamente para outro projeto. Para mover um pipeline para outro projeto, pare o job atual e recrie-o no novo projeto. Para um guia detalhado sobre esse processo, consulte Migrar jobs de pipeline para outro projeto do Google Cloud .
Saiba mais
Nas páginas a seguir, você encontra mais informações sobre como estruturar o pipeline, como escolher quais transformações aplicar aos dados e o que considerar ao escolher os métodos de entrada e saída do pipeline.
Para mais informações sobre como criar seu código de usuário, consulte os requisitos para funções fornecidas pelo usuário.