Resolvendo Desbalanceamento de Dados no Apache Spark com Agregação de Chaves em Duas Fases

O desbalanceamento de dados, conhecido como "data skew", é um desafio recorrente em sistemas de processamento distribuído, como o Apache Spark. Ele se manifesta quando certos valores de chave em um conjunto de dados possuem uma frequência significativamente maior do que outros, resultando em tarefas de processamento excessivamente longas em alguns nós do cluster, enquanto outros permanecem subutilizados. Essa assimetria pode levar a uma degradação severa do desempenho geral da aplicação.

Uma estratégia robusta para mitigar o desbalanceamento de dados durante operações de agregação é a técnica de agregação em duas fases (ou "double key aggregation"). Essa abordagem introduz uma etapa de pré-agregação, visando distribuir a carga de trabalho de maneira mais equitativa entre os recursos computacionais.

Princípios da Agregação em Duas Fases

O fluxo de trabalho desta estratégia é dividido em etapas distintas:

  1. Enriquecimento da Chave com Prefixo Aleatório: Para as chaves identificadas como potenciais fontes de desbalanceamento, um prefixo aleatório é adicionado a cada ocorrência. Isso efetivamente transforma uma única chave de alta frequência ("hot key") em múltiplas chaves temporárias, cada uma com um prefixo diferente (ex: "chaveX" pode se tornar "1_chaveX", "2_chaveX", "N_chaveX", dependendo do prefixo gerado).
  2. Primeira Agregação (Pré-agregação): Uma etapa inicial de agregação é executada utilizando essas chaves prefixadas. A diversificação das chaves permite que o Spark distribua o trabalho de agregação de forma mais equilibrada entre os seus executores, processando subconjuntos menores de dados para cada chave prefixada.
  3. Remoção do Prefixo: Após a pré-agregação, o prefixo aleatório é descartado, restaurando as chaves à sua forma original.
  4. Segunda Agregação (Agregação Final): Uma segunda e última operação de agregação é realizada sobre as chaves originais. Como grande parte do volume de dados e do trabalho pesado já foi processada na fase anterior de pré-agregação, esta etapa final opera com um conjunto de dados mais consolidado e, idealmente, menos desbalanceado, garantindo um processamento eficiente.

Implementação Prática com UDFs e Spark SQL

A aplicação dessa técnica no Apache Spark geralmente envolve o uso de User-Defined Functions (UDFs) para manipulação das chaves e a expressividade do Spark SQL para as operações de agregação.

Registro de Funções Personalizadas no Spark

Para iniciar, é necesário registarr as UDFs e, se aplicável, UDAFs (User-Defined Aggregate Functions) no contexto do SparkSession:


import org.apache.spark.sql.types.DataTypes;
// Assumindo 'spark' é uma instância de SparkSession
// Registra funções personalizadas para manipulação de chaves e agregação
spark.udf().register("concatenar_partes_chave", new ConcatenadorChavesUDF(), DataTypes.StringType);
spark.udf().register("gerar_prefixo_aleatorio", new GeradorPrefixoAleatorioUDF(), DataTypes.StringType);
spark.udf().register("extrair_chave_original", new RemovedorPrefixoAleatorioUDF(), DataTypes.StringType);

// Registra uma UDAF para concatenação distinta de valores agrupados (implementação de 'AgruparConcatenarDistintosUDAF' não detalhada aqui)
spark.udf().register("agrupar_e_concatenar_distintos", new AgruparConcatenarDistintosUDAF());

Exemplo de UDF: Gerador de Prefixo Aleatório

A UDF responsável por adicionar o prefixo aleatório é fundamental para a distribuição das chaves. Abaixo, um exemplo de sua implementação em Java:


package com.suaempresa.spark.udfs;

import java.util.Random;
import org.apache.spark.sql.api.java.UDF2;

/**
 * UDF que adiciona um prefixo numérico aleatório a uma chave de entrada.
 * O prefixo é gerado com base em um limite superior fornecido.
 */
public class GeradorPrefixoAleatorioUDF implements UDF2<String, Integer, String> {

    private static final long serialVersionUID = 1L;

    /**
     * Gera um prefixo aleatório e o concatena com a chave original.
     * Ex: (key, 10) -> "3_key"
     *
     * @param chaveEntrada A string da chave original.
     * @param numeroMaxPrefixo O limite superior (exclusivo) para o valor do prefixo aleatório.
     * @return A chave com o prefixo aleatório no formato "prefixo_chaveOriginal".
     */
    @Override
    public String call(String chaveEntrada, Integer numeroMaxPrefixo) throws Exception {
        Random rand = new Random();
        int prefixo = rand.nextInt(numeroMaxPrefixo);
        return prefixo + "_" + chaveEntrada;
    }
}

UDF para Concatenação de Chaves

Para criar a chave composta que será prefixada (e.g., combinação de área e rodovia):


package com.suaempresa.spark.udfs;

import org.apache.spark.sql.api.java.UDF3;

/**
 * UDF para concatenar múltiplas partes de uma string com um separador.
 * Ex: (parte1, parte2, ":") -> "parte1:parte2"
 */
public class ConcatenadorChavesUDF implements UDF3<String, String, String, String> {
    private static final long serialVersionUID = 1L;

    @Override
    public String call(String parte1, String parte2, String separador) throws Exception {
        return parte1 + separador + parte2;
    }
}

UDF para Remoção de Prefixo

Para extrair a chave original após a pré-agregação:


package com.suaempresa.spark.udfs;

import org.apache.spark.sql.api.java.UDF1;

/**
 * UDF que remove o prefixo aleatório de uma chave.
 * Espera o formato "prefixo_chaveOriginal".
 */
public class RemovedorPrefixoAleatorioUDF implements UDF1<String, String> {
    private static final long serialVersionUID = 1L;

    @Override
    public String call(String chaveComPrefixo) throws Exception {
        int indicePrimeiroSublinhado = chaveComPrefixo.indexOf('_');
        if (indicePrimeiroSublinhado > -1 && indicePrimeiroSublinhado < chaveComPrefixo.length() - 1) {
            return chaveComPrefixo.substring(indicePrimeiroSublinhado + 1);
        }
        return chaveComPrefixo; // Retorna a chave original se o formato não for o esperado
    }
}

Executando a Agregação em Duas Fases com Spark SQL

Considerando uma tabela temporária registros_trafego_bruto com colunas como nome_area, id_rodovia, id_monitor e placa_veiculo, a consulta SQL para aplicar a agregação em duas fases seria estruturada da seguinte forma:


-- Consulta SQL para aplicar a agregação em duas fases e mitigar desbalanceamento de dados
WITH ChavesEnriquecidas AS (
    -- Etapa 1: Concatena chaves e adiciona um prefixo aleatório para redistribuição
    SELECT
        id_monitor,
        placa_veiculo,
        -- Combina 'nome_area' e 'id_rodovia' e aplica um prefixo aleatório (e.g., 10 prefixos possíveis)
        gerar_prefixo_aleatorio(concatenar_partes_chave(nome_area, id_rodovia, ':'), 10) AS chave_prefixada_destino
    FROM registros_trafego_bruto
),
PrimeiraAgregacaoParcial AS (
    -- Etapa 2: Realiza a pré-agregação nas chaves prefixadas
    SELECT
        chave_prefixada_destino,
        COUNT(placa_veiculo) AS contagem_veiculos_parcial,
        -- Agrupa informações de monitores distintos para cada chave prefixada
        agrupar_e_concatenar_distintos(id_monitor) AS detalhes_monitores_parcial
    FROM ChavesEnriquecidas
    GROUP BY chave_prefixada_destino
),
ChavesOriginalizadas AS (
    -- Etapa 3: Remove o prefixo aleatório para obter as chaves originais
    SELECT
        extrair_chave_original(chave_prefixada_destino) AS chave_area_rodovia,
        contagem_veiculos_parcial,
        detalhes_monitores_parcial
    FROM PrimeiraAgregacaoParcial
)
-- Etapa 4: Agregação final sobre as chaves originais
SELECT
    chave_area_rodovia,
    SUM(contagem_veiculos_parcial) AS total_veiculos,
    -- Agrupa novamente as informações dos monitores para a chave original
    agrupar_e_concatenar_distintos(detalhes_monitores_parcial) AS detalhes_monitoramento_final
FROM ChavesOriginalizadas
GROUP BY chave_area_rodovia;

O resultado desta consulta SQL será um Dataset<Row>, que pode ser processado ou salvo. Por exemplo, para disponibliizá-lo para outras consultas Spark SQL, pode-se registrá-lo como uma view temporária:


import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
// Assumindo 'spark' é uma instância de SparkSession e 'sqlConsultaAgregacaoDupla' contém a consulta SQL acima
Dataset<Row> resultadoAgregacao = spark.sql(sqlConsultaAgregacaoDupla);
resultadoAgregacao.createOrReplaceTempView("view_analise_fluxo_trafego");

Tags: Apache Spark Data Skew Spark SQL UDF agregação

Publicado em 7-28 08:06