Todos os produtos
Search
Central de documentação

MaxCompute:Ajuste de data skew

Última atualização: Jul 03, 2026

O data skew faz com que alguns workers processem volumes de dados significativamente maiores que outros, desacelerando ou paralisando jobs distribuídos. Aprenda a identificar data skew no MaxCompute e aplique soluções para cenários comuns envolvendo JOIN, GROUP BY, COUNT(DISTINCT), ROW_NUMBER() e partições dinâmicas.

MapReduce

MapReduce é um framework de computação distribuída que divide problemas grandes em subproblemas menores, resolve-os independentemente e mescla os resultados. Ele abstrai as complexidades de um cluster distribuído — armazenamento de dados, comunicação entre nós e transferência de dados — permitindo escrever programas paralelos sem gerenciar a infraestrutura de baixo nível. Compreender o MapReduce é essencial para diagnosticar data skew.

O diagrama a seguir ilustra o fluxo de trabalho do MapReduce:MapReduce

Data skew

O data skew ocorre tipicamente na fase de Reducer. Embora os Mappers particionem os arquivos de entrada de maneira bastante uniforme, os dados resultantes podem ser distribuídos desigualmente entre os workers. Alguns workers terminam rapidamente, enquanto outros recebem uma parcela desproporcional e demoram muito mais. A maioria dos dados reais apresenta skew, frequentemente seguindo a regra 80/20 — por exemplo, 20% dos usuários de um fórum podem contribuir com 80% das postagens, ou 20% dos usuários podem responder por 80% do tráfego do site. Na escala de Big Data, esse desequilíbrio pode degradar severamente o desempenho, muitas vezes fazendo com que um job pareça travado em 99% de progresso.

Identificar data skew

No MaxCompute, use o Logview para identificar data skew. Siga este procedimento:判断数据倾斜

  1. Em Fuxi Jobs, ordene os estágios do job por Latency em ordem decrescente e selecione o estágio com o maior tempo de execução.

  2. Na lista de instâncias Fuxi do estágio selecionado, ordene as instâncias por Latency em ordem decrescente. Selecione a instância com tempo de execução significativamente superior à média, geralmente a primeira da lista. Em seguida, verifique seu log StdOut.

  3. Use o log StdOut para visualizar o grafo de execução do job.

  4. Utilize as informações-chave do grafo de execução do job para localizar o trecho SQL causador do data skew.

O exemplo a seguir demonstra esse procedimento.

  1. Encontre a URL do Logview nos logs operacionais da tarefa. Para mais informações, consulte Pontos de entrada do Logview.logview

  2. Na página do Logview, ordene as tarefas Fuxi por Latency em ordem decrescente e selecione aquela com o maior tempo de execução para identificar rapidamente o problema.Fuxi Task

  3. A tarefa R31_26_27 apresenta a maior latência. Clique em R31_26_27 para acessar a página de detalhes da instância, conforme mostra a figura a seguir.Task with the longest latencyLatency: {min:00:00:06, avg:00:00:13, max:00:26:40} indica que, para todas as instâncias da tarefa, a latência mínima é 6s, a média é 13s e a máxima é 26 minutos e 40s. Ordene as instâncias por Latency em ordem decrescente para visualizar quatro instâncias com latências relativamente altas. O MaxCompute identifica uma Fuxi Instance como instância de long-tail se sua latência for superior ao dobro da média. Isso significa que instâncias de tarefa com latência maior que 26s são classificadas como instâncias de long-tail (Long-Tails). Neste caso, 21 instâncias têm latência superior a 26s. A presença de instâncias Long-Tails não indica necessariamente data skew. Compare também os valores avg e max da latência da instância. Para tarefas em que o valor max supera muito o valor avg, indicando data skew severo, otimize essas tarefas.

  4. Clique em 输出日志 na coluna StdOut para visualizar o log StdOut, conforme mostra o exemplo a seguir.输出示例结果

  5. Após identificar o problema, na aba Job Details, clique com o botão direito em R31_26_27 e selecione expand all para expandir a tarefa. Para mais informações, consulte Usando o Logview 2.0 para visualizar informações de execução de jobs.Expand task Examine StreamLineWriter21, etapa anterior a StreamLineRead22, para localizar as chaves causadoras do data skew: new_uri_path_structure, cookie_x5check_userid e cookie_userid. Esse processo também ajuda a localizar o fragmento SQL causador do data skew.KEY

Solução de problemas de data skew

Operações comuns que causam data skew incluem:

  • join

  • group by

  • count(distinct)

  • row_number() (TopN)

  • dynamic partition

Essas operações estão classificadas por frequência da seguinte forma: join > group by > count(distinct) > row_number() > dynamic partition.

Join

O data skew em operações de join ocorre em vários cenários: junção de uma tabela grande com uma pequena, junção de uma tabela grande com uma média ou presença de valores quentes (hot values) que criam uma long tail.

  • Junção de tabela grande com tabela pequena

    • Exemplo de data skew

      No exemplo a seguir, t1 é uma tabela grande, enquanto t2 e t3 são tabelas pequenas.

      SELECT  t1.ip
              ,t1.is_anon
              ,t1.user_id
              ,t1.user_agent
              ,t1.referer
              ,t2.ssl_ciphers
              ,t3.shop_province_name
              ,t3.shop_city_name
      FROM    <viewtable> t1
      LEFT OUTER JOIN <other_viewtable> t2
      ON t1.header_eagleeye_traceid = t2.eagleeye_traceid
      LEFT OUTER JOIN (  SELECT  shop_id
                                  ,city_name AS shop_city_name
                                  ,province_name AS shop_province_name
                          FROM    <tenanttable>
                          WHERE   ds = MAX_PT('<tenanttable>')
                          AND     is_valid = 1
                      ) t3
      ON t1.shopid = t3.shop_id
    • Solução

      Use a sintaxe MAPJOIN hint, conforme mostra o exemplo a seguir.

      SELECT  /*+ mapjoin(t2,t3)*/
              t1.ip
              ,t1.is_anon
              ,t1.user_id
              ,t1.user_agent
              ,t1.referer
              ,t2.ssl_ciphers
              ,t3.shop_province_name
              ,t3.shop_city_name
      FROM    <viewtable> t1
      LEFT OUTER JOIN (<other_viewtable>) t2
      ON t1.header_eagleeye_traceid = t2.eagleeye_traceid
      LEFT OUTER JOIN (  SELECT  shop_id
                                  ,city_name AS shop_city_name
                                  ,province_name AS shop_province_name
                          FROM    <tenanttable>
                          WHERE   ds = MAX_PT('<tenanttable>')
                          AND     is_valid = 1
                      ) t3
      ON t1.shopid = t3.shop_id
      • Observações

        • Ao referenciar uma tabela pequena ou subconsulta, use obrigatoriamente seu alias.

        • O MapJoin suporta o uso de subconsultas como tabela pequena.

        • Em um MapJoin, use non-equi joins ou or para conectar múltiplas condições. Para calcular um produto cartesiano, omita a cláusula on e use a condição mapjoin on 1 = 1. Exemplo: select /*+ mapjoin(a) */ a.id from shop a join table_name b on 1=1;. No entanto, essa operação pode causar expansão significativa dos dados.

        • Em um MAPJOIN hint, separe múltiplas tabelas pequenas com vírgula (,). Exemplo: /*+ mapjoin(a,b,c)*/.

        • Em um MapJoin, o MaxCompute carrega toda a tabela especificada na memória durante a fase Map. Portanto, a tabela especificada deve ser pequena. Como o MaxCompute usa armazenamento compactado, os dados podem expandir significativamente ao serem carregados. Esse tamanho expandido na memória não pode exceder 512 MB. Aumente esse limite até no máximo 8192 MB definindo o seguinte parâmetro:

          SET odps.sql.mapjoin.memory.max=2048;
      • Limitações das operações MapJoin

        • Para um left outer join, a tabela à esquerda deve ser a grande.

        • Para um right outer join, a tabela à direita deve ser a grande.

        • full outer join não é suportado.

        • Para um inner join, tanto a tabela da esquerda quanto a da direita podem ser a grande.

        • O MapJoin suporta no máximo 128 tabelas pequenas. Exceder esse limite causa erro de sintaxe.

  • Junção de tabela grande com tabela média

    • Exemplo de data skew

      No exemplo a seguir, t0 é uma tabela grande e t1 é uma tabela média.

      SELECT  request_datetime
              ,host
              ,URI
              ,eagleeye_traceid
      FROM <viewtable>
          t0
      LEFT JOIN (
          SELECT
          traceid,
          eleme_uid,
          isLogin_is
          FROM <servicetable>
          WHERE ds = '${today}'
          AND     hh = '${hour}'
      ) t1 ON t0.eagleeye_traceid = t1.traceid
      WHERE   ds = '${today}'
      AND     hh = '${hour}'
    • Solução

      Use a sintaxe DISTRIBUTED MAPJOIN para resolver o data skew, conforme mostra o exemplo a seguir.

      SELECT  /*+distmapjoin(t1)*/
              request_datetime
              ,host
              ,URI
              ,eagleeye_traceid
      FROM <viewtable>
          t0
      LEFT JOIN (
          SELECT
          traceid,
          eleme_uid,
          isLogin_is
          FROM <servicetable>
          WHERE ds = '${today}'
          AND     hh = '${hour}'
      ) t1 ON t0.eagleeye_traceid = t1.traceid
      WHERE   ds = '${today}'
      AND     hh = '${hour}'
  • Long tail de valores quentes em um join

    • Exemplo de data skew

      Na consulta a seguir, a coluna eleme_uid possui muitos valores quentes, o que frequentemente causa data skew.

      SELECT
      eleme_uid,
      ...
      FROM (
          SELECT
          eleme_uid,
          ...
          FROM <viewtable>
      )t1
      LEFT JOIN(
          SELECT
          eleme_uid,
          ...
          FROM <customertable>
      )  t2
      ON t1.eleme_uid = t2.eleme_uid;
    • Solução

      Use um dos quatro métodos a seguir para resolver esse problema.

      ID

      Método

      Descrição

      Método 1

      Divisão manual de valores quentes

      Identifique os valores quentes. Filtre os registros contendo esses valores da tabela principal e execute um MapJoin neles. Em seguida, realize um join padrão nos registros restantes com valores não quentes. Por fim, combine os resultados de ambos os joins usando UNION ALL.

      Método 2

      Definição de parâmetros SkewJoin

      Ative o SkewJoin com o comando: set odps.sql.skewjoin=true;.

      Método 3

      SkewJoin hint

      Use o hint: /*+ skewJoin(<table_name>[(<column1_name>[,<column2_name>,...])][((<value11>,<value12>)[,(<value21>,<value22>)...])]*/. Este método adiciona uma etapa para encontrar chaves com skew, o que pode aumentar o tempo de execução da consulta. Se já conhecer as chaves com skew, economize tempo definindo diretamente os parâmetros do SkewJoin.

      Método 4

      Join com módulo e tabela multiplicadora

      Use uma tabela multiplicadora para redistribuir as chaves de junção.

      • Divisão manual de valores quentes

        Após identificar os valores quentes, filtre os registros com esses valores da tabela principal. Execute um MapJoin nesses registros e depois realize um join padrão nos registros restantes. Por fim, combine os resultados de ambos os joins usando UNION ALL. O código a seguir fornece um exemplo:

        SELECT
        /*+ MAPJOIN (t2) */
        eleme_uid,
        ...
        FROM (
            SELECT
            eleme_uid,
            ...
            FROM <viewtable>
            WHERE eleme_uid = <skewed_value>
        )t1
        LEFT JOIN(
            SELECT
            eleme_uid,
            ...
            FROM <customertable>
            WHERE eleme_uid = <skewed_value>
        )  t2
        ON t1.eleme_uid = t2.eleme_uid
        UNION ALL
        SELECT
        eleme_uid,
        ...
        FROM (
            SELECT
            eleme_uid,
            ...
            FROM <viewtable>
            WHERE eleme_uid != <skewed_value>
        )t3
        LEFT JOIN(
            SELECT
            eleme_uid,
            ...
            FROM <customertable>
            WHERE eleme_uid != <skewed_value>
        )  t4
        ON t3.eleme_uid = t4.eleme_uid
      • Definição de parâmetros SkewJoin

        Esta é uma solução comum. Embora seja possível ativar o recurso SkewJoin no MaxCompute executando o comando set odps.sql.skewjoin=true;, isso não afeta a execução da tarefa isoladamente. Para que o recurso tenha efeito, defina também o parâmetro odps.sql.skewinfo. O parâmetro odps.sql.skewinfo configura detalhes específicos para a otimização do join. O exemplo a seguir mostra a sintaxe do comando.

        SET odps.sql.skewjoin=true;
        SET odps.sql.skewinfo=skewed_src:(skewed_key)[("skewed_value")];  --skewed_src is the skewed table, and skewed_value is the hot value.

        Exemplos de uso:

        -- For a single skewed value in a single column
        SET odps.sql.skewinfo=src_skewjoin1:(key)[("0")];
        
        -- For multiple skewed values in a single column
        SET odps.sql.skewinfo=src_skewjoin1:(key)[("0")("1")];
      • SkewJoin hint

        Adicione o seguinte hint à instrução SELECT: /*+ skewJoin(<table_name>[(<column1_name>[,<column2_name>,...])][((<value11>,<value12>)[,(<value21>,<value22>)...])]*/. Neste hint, table_name é o nome da tabela com skew, column_name é o nome da coluna com skew e value é o valor da chave com skew. Os exemplos a seguir mostram como usar este hint.

        -- Method 1: Hint the table name. Note that you must hint the table's alias.
        SELECT /*+ skewjoin(a) */ * FROM T0 a JOIN T1 b ON a.c0 = b.c0 AND a.c1 = b.c1;
        
        -- Method 2: Hint the table name and the columns that you suspect are skewed. For example, columns c0 and c1 in table 'a' are skewed.
        SELECT /*+ skewjoin(a(c0, c1)) */ * FROM T0 a JOIN T1 b ON a.c0 = b.c0 AND a.c1 = b.c1 AND a.c2 = b.c2;
        
        -- Method 3: Hint the table name, columns, and provide the skewed key values. If a key is a STRING type, enclose it in quotation marks. For example, values where (a.c0=1 and a.c1="2") and (a.c0=3 and a.c1="4") are skewed.
        SELECT /*+ skewjoin(a(c0, c1)((1, "2"), (3, "4"))) */ * FROM T0 a JOIN T1 b ON a.c0 = b.c0 AND a.c1 = b.c1 AND a.c2 = b.c2;
        Nota

        Especificar valores com skew diretamente em um SkewJoin hint é mais eficiente do que dividir manualmente valores quentes ou ativar o SkewJoin sem fornecer os valores.

        Tipos de join suportados pelo SkewJoin hint:

        • Para inner join, aplique o hint em qualquer uma das tabelas da junção.

        • Para left join, semi-join ou anti-join, aplique o hint apenas na tabela da esquerda.

        • Para right join, aplique o hint apenas na tabela da direita.

        • O SkewJoin hint não é suportado para full join.

        Use este hint apenas para joins com data skew conhecido, pois ele introduz sobrecarga ao executar uma consulta de agregação.

        Os tipos de dados das chaves de junção no lado esquerdo devem corresponder aos tipos de dados das chaves de junção no lado direito. Caso contrário, o SkewJoin hint não terá efeito. Por exemplo, o tipo de dados de a.c0 deve corresponder ao tipo de dados de b.c0, e o tipo de dados de a.c1 deve corresponder ao tipo de dados de b.c1. Faça cast das chaves de junção em uma subconsulta para garantir a correspondência dos tipos de dados. Veja um exemplo:

        CREATE TABLE T0(c0 int, c1 int, c2 int, c3 int);
        CREATE TABLE T1(c0 string, c1 int, c2 int);
        
        -- Method 1:
        SELECT /*+ skewjoin(a) */ * FROM T0 a JOIN T1 b ON cast(a.c0 AS string) = cast(b.c0 AS string) AND a.c1 = b.c1;
        
        -- Method 2:
        SELECT /*+ skewjoin(b) */ * FROM (SELECT cast(a.c0 AS string) AS c00 FROM T0 a) b JOIN T1 c ON b.c00 = c.c0;

        Ao adicionar um SkewJoin hint, o otimizador executa uma operação Aggregate para recuperar os 20 principais valores quentes. 20 é o valor padrão; altere-o executando o comando set odps.optimizer.skew.join.topk.num = xx;.

        • O SkewJoin hint só pode ser aplicado a um lado do join.

        • Um Join com hint deve ter uma condição left key = right key. Joins de produto cartesiano não são suportados.

        • Não adicione um SkewJoin hint a um join que já usa um MAPJOIN hint.

      • Join com módulo e tabela multiplicadora

        Esta solução usa uma abordagem diferente. Em vez de dividir para conquistar, ela cria uma tabela multiplicadora com uma única coluna int contendo valores de 1 a N, onde N reflete o grau de skew. A tabela multiplicadora expande a tabela de comportamento do usuário N vezes, e o join passa a usar duas chaves: o ID do usuário e number. Adicionar a chave number reduz o skew causado por um único ID de usuário para 1/N do nível original, ao custo de expandir os dados N vezes.

        SELECT
        eleme_uid,
        ...
        FROM (
            SELECT
            eleme_uid,
            ...
            FROM <viewtable>
        )t1
        LEFT JOIN(
            SELECT
            /*+mapjoin(<multipletable>)*/
            eleme_uid,
            number
            ...
            FROM <customertable>
            JOIN <multipletable>
        )  t2
        ON t1.eleme_uid = t2.eleme_uid
        AND mod(t1.<value_col>,10)+1 = t2.number;

        Limite a inflação apenas aos registros de hotspot, mantendo os demais inalterados. Primeiro, identifique os valores de hotspot. Em seguida, adicione uma coluna chamada eleme_uid_join às tabelas de tráfego e de comportamento do usuário. Para IDs de usuário de hotspot, faça concat do ID com um inteiro aleatório de um intervalo predefinido (como 0 a 1000); para outros IDs, mantenha o valor original. Use eleme_uid_join como chave de junção. Isso reduz o skew para valores de hotspot sem inflar dados sem hotspot. No entanto, o SQL resultante torna-se muito mais difícil de ler e manter; portanto, este método não é recomendado.

Group by

A consulta de exemplo a seguir usa uma cláusula Group By.

SELECT  shop_id
        ,sum(is_open) AS open_days
FROM    table_xxx_di
WHERE   dt BETWEEN '${bizdate_365}' AND '${bizdate}'
GROUP BY shop_id;

Para resolver o data skew, use uma das três soluções a seguir:

Solução

Descrição

Solução 1

Definir parâmetro para lidar com data skew em Group By

set MaxCompute.sql.groupby.skewindata=true;.

Solução 2

Adicionar número aleatório

Divida as chaves que causam a long tail.

Solução 3

Criar tabela rollup

Reduz I/O e uso de recursos pré-agregando dados.

  • Solução 1: Definir parâmetro para mitigar data skew em operações Group By.

    SET MaxCompute.sql.groupby.skewindata=true;
  • Solução 2: Adicionar número aleatório.

    Diferentemente da Solução 1, este método requer reescrever a consulta SQL para adicionar um número aleatório. Dividir chaves de long-tail é uma maneira altamente eficaz de resolver data skew em Group By.

    Para a consulta Select Key,Count(*) As Cnt From TableName Group By Key;, sem um Combiner, os Mappers fazem shuffle dos dados para os Reducers, que então executam a operação COUNT. Isso corresponde a um plano de execução M->R.

    Se identificou a chave de long-tail, redistribua sua carga de trabalho da seguinte forma:

    -- Assume the long-tail key is 'KEY001'
    SELECT  a.Key
            ,SUM(a.Cnt) AS Cnt
    FROM(SELECT  Key
                ,COUNT(*) AS Cnt
                FROM    <TableName>
                GROUP BY Key
                ,CASE WHEN KEY = 'KEY001' THEN Hash(Random()) % 50
                 ELSE 0
                END
            ) a
    GROUP BY a.Key;

    A consulta modificada altera o plano de execução para M->R->R. Embora isso adicione um estágio, o tempo total de execução pode diminuir porque a chave com skew é processada em dois estágios. O consumo de recursos e o desempenho são semelhantes aos da Solução 1. No entanto, em cenários reais, frequentemente há mais de uma chave com skew. Considerando o esforço para identificar chaves com skew e reescrever o SQL, a Solução 1 é mais econômica de implementar.

  • Solução 3: Criar tabela rollup.

    Para otimização de custos essenciais, o requisito principal é recuperar dados de comerciantes do último ano. Para tarefas online, ler todas as partições de T-1 a T-365 repetidamente representa um desperdício significativo de recursos. Criar uma tabela rollup reduz o número de partições lidas sem afetar a recuperação de dados do último ano. Veja um exemplo.

    Primeiro, inicialize 365 dias de dados comerciais do comerciante usando uma agregação Group By e salve-os como tabela a com uma data de atualização marcada. As tarefas online subsequentes passam então a fazer join da tabela T-2 a com a tabela table_xxx_di e realizar outra agregação Group By. Como resultado, a quantidade de dados lidos diariamente cai de 365 para 2 dias, a duplicação da chave primária shopid diminui drasticamente e o consumo de recursos também é reduzido.

    --Create a rollup table
    CREATE TABLE IF NOT EXISTS m_xxx_365_df
    (
      shop_id STRING COMMENT,
      last_update_ds COMMENT,
      365d_open_days COMMENT
    )
    PARTITIONED BY
    (
      ds STRING COMMENT 'date partition'
    )LIFECYCLE 7;
    -- Assume the 365-day period is from 2021-05-01 to 2022-05-01. First, perform a one-time initialization.
    INSERT OVERWRITE TABLE m_xxx_365_df PARTITION(ds = '20220501')
      SELECT shop_id,
             max(ds) as last_update_ds,
             sum(is_open) AS 365d_open_days
      FROM table_xxx_di
      WHERE dt BETWEEN '20210501' AND '20220501'
      GROUP BY shop_id;
    -- The subsequent online job runs the following:
    INSERT OVERWRITE TABLE m_xxx_365_df PARTITION(ds = '${bizdate}')
      SELECT aa.shop_id, 
             aa.last_update_ds, 
             365d_open_days - COALESCE(is_open, 0) AS 365d_open_days -- Prevents open days from rolling up indefinitely
      FROM (
        SELECT shop_id, 
               max(last_update_ds) AS last_update_ds, 
               sum(365d_open_days) AS 365d_open_days
        FROM (
          SELECT shop_id,
                 ds AS last_update_ds,
                 sum(is_open) AS 365d_open_days
          FROM table_xxx_di
          WHERE ds = '${bizdate}'
          GROUP BY shop_id
          UNION ALL
          SELECT shop_id,
                 last_update_ds,
                 365d_open_days
          FROM m_xxx_365_df
          WHERE dt = '${bizdate_2}' AND last_update_ds >= '${bizdate_365}'
          GROUP BY shop_id
        )
        GROUP BY shop_id
      ) AS aa
      LEFT JOIN (
        SELECT shop_id,
               is_open
        FROM table_xxx_di
        WHERE ds = '${bizdate_366}'
      ) AS bb
      ON aa.shop_id = bb.shop_id;
                                

Count(distinct)

Considere uma tabela com a seguinte distribuição de dados.

Ds (partição)

Cnt (contagem de registros)

20220416

73025514

20220415

2292806

20220417

2319160

A consulta a seguir é suscetível a data skew:

SELECT  ds
        ,COUNT(DISTINCT shop_id) AS cnt
FROM    demo_data0
GROUP BY ds;

As soluções a seguir resolvem esse problema:

Solução

Descrição

Solução 1

Ajuste de parâmetros

SET odps.sql.groupby.skewindata=true;

Solução 2

Agregação genérica em dois estágios

Acrescente um número aleatório ao valor do campo de partição.

Solução 3

Agregação manual em dois estágios

Primeiro, aplique GROUP BY em ambos os campos de agrupamento (ds, shop_id) e, em seguida, use count(distinct).

  • Solução 1: Ajuste de parâmetros

    Defina o seguinte parâmetro.

    SET odps.sql.groupby.skewindata=true;
  • Solução 2: Agregação genérica em dois estágios

    Se os dados no campo shop_id tiverem skew, a Solução 1 pode não ser eficaz. Um método mais geral é acrescentar um número aleatório ao valor do campo de partição.

    -- Method 1: Append a random number, for example, CONCAT(ROUND(RAND(),1)*10,'_', ds) AS rand_ds
    SELECT  SPLIT_PART(rand_ds, '_',2) ds
            ,COUNT(*) id_cnt
      FROM (
            SELECT  rand_ds
                    ,shop_id
            FROM    demo_data0
            GROUP BY rand_ds,shop_id
            )
    GROUP BY SPLIT_PART(rand_ds, '_',2);
    
    -- Method 2: Add a new random number column, for example, ROUND(RAND(),1)*10 AS randint10
    SELECT  ds
            ,COUNT(*) id_cnt
    FROM    (SELECT  ds
                     ,randint10
                     ,shop_id
               FROM  demo_data0
            GROUP BY ds,randint10,shop_id
            )
    GROUP BY ds;
  • Solução 3: Agregação manual em dois estágios

    Se os dados estiverem distribuídos uniformemente entre os campos ds e shop_id, otimize a consulta agrupando primeiro por ambos os campos e depois usando a função count(distinct).

    SELECT  ds
            ,COUNT(*) AS cnt
    FROM(SELECT  ds
                ,shop_id
                FROM    demo_data0
                GROUP BY ds ,shop_id
        )
    GROUP BY ds;

ROW_NUMBER(): TopN

A consulta a seguir recupera os 10 principais registros.

SELECT  main_id
        ,type
FROM    (SELECT  main_id
                 ,type
                 ,ROW_NUMBER() OVER(PARTITION BY main_id ORDER BY type DESC ) rn
            FROM <data_demo2>
        ) A
WHERE   A.rn <= 10;

Para resolver o data skew, use uma das soluções a seguir:

Solução

Descrição

Solução 1

Agregação em dois estágios em SQL.

Adicione uma coluna aleatória e use-a como chave de partição adicional.

Solução 2

Agregação em dois estágios com UDAF.

Otimize a consulta usando uma UDAF que implementa uma fila de prioridade min-heap.

  • Solução 1: Agregação em dois estágios em SQL

    Para garantir distribuição uniforme dos dados entre os grupos na fase map, adicione uma coluna aleatória e use-a como chave de partição adicional.

    SELECT  main_id
            ,type
      FROM  (SELECT  main_id
                     ,type
                     ,ROW_NUMBER() OVER(PARTITION BY main_id ORDER BY type DESC ) rn
                FROM (SELECT  main_id
                              ,type
                            FROM (SELECT  main_id
                                          ,type
                                          ,ROW_NUMBER() OVER(PARTITION BY main_id,src_pt ORDER BY type DESC ) rn
                                     FROM (SELECT  main_id
                                                   ,type
                                                   ,ceil(110 * rand()) % 11 AS src_pt
                                             FROM  data_demo2
                                          )
                                    ) B
                            WHERE   B.rn <= 10
                        )
            ) A
    WHERE   A.rn <= 10;
    -- 2. An alternative implementation:
    SELECT  main_id
            ,type
      FROM  (SELECT  main_id
                     ,type
                     ,ROW_NUMBER() OVER(PARTITION BY main_id ORDER BY type DESC ) rn
                FROM (SELECT  main_id
                              ,type
                        FROM(SELECT  main_id
                                     ,type
                                     ,ROW_NUMBER() OVER(PARTITION BY main_id,src_pt ORDER BY type DESC ) rn
                               FROM  (SELECT  main_id
                                              ,type
                                              ,ceil(10 * rand()) AS src_pt
                                              FROM    data_demo2
                                      )
                                    ) B
                            WHERE B.rn <= 10
                        )
            ) A
    WHERE   A.rn <= 10;
  • Solução 2: Agregação em dois estágios com UDAF

    A abordagem SQL pode ser verbosa e difícil de manter. Uma alternativa mais limpa é uma UDAF que implementa uma fila de prioridade min-heap, retendo apenas os elementos TopN na fase iterate e mesclando apenas N elementos na fase merge. O processo é o seguinte.

    • iterate: Insere os primeiros K elementos. Em seguida, compara cada elemento subsequente com o menor elemento no topo do heap e o troca por um elemento no heap.

    • merge: Mescla dois heaps in-place e retorna os K principais elementos.

    • terminate: Retorna o heap como um array.

    • A consulta SQL usa então LATERAL VIEW EXPLODE para desaninhar o array em linhas separadas.

    @annotate('* -> array<string>')
    class GetTopN(BaseUDAF):
        def new_buffer(self):
            return [[], None]
        def iterate(self, buffer, order_column_val, k):
            # heapq.heappush(buffer, order_column_val)
            # buffer = [heapq.nlargest(k, buffer), k]
            if not buffer[1]:
                buffer[1] = k
            if len(buffer[0]) < k:
                heapq.heappush(buffer[0], order_column_val)
            else:
                heapq.heappushpop(buffer[0], order_column_val)
        def merge(self, buffer, pbuffer):
            first_buffer, first_k = buffer
            second_buffer, second_k = pbuffer
            k = first_k or second_k
            merged_heap = first_buffer + second_buffer
            merged_heap.sort(reverse=True)
            merged_heap = merged_heap[0: k] if len(merged_heap) > k else merged_heap
            buffer[0] = merged_heap
            buffer[1] = k
        def terminate(self, buffer):
            return buffer[0]
    
    SET odps.sql.python.version=cp37;
    SELECT main_id,type_val FROM (
      SELECT  main_id ,get_topn(type, 10) AS type_array
      FROM data_demo2
      GROUP BY main_id
    )
    LATERAL VIEW EXPLODE(type_array)type_ar AS type_val;

Dynamic partition

O particionamento dinâmico permite inserir dados em uma tabela particionada especificando um nome de coluna de partição sem fornecer seu valor. O valor é derivado da coluna correspondente na cláusula SELECT. Consequentemente, as partições específicas a serem criadas permanecem desconhecidas até a execução da instrução SQL. O sistema cria novas partições somente após a conclusão da instrução, com base nos valores gerados para a coluna de partição. Para mais informações, consulte Inserir ou sobrescrever dados em partições dinâmicas (DYNAMIC PARTITION). Veja um exemplo.

CREATE TABLE total_revenues (revenue bigint) partitioned BY (region string);

INSERT overwrite TABLE total_revenues PARTITION(region)
SELECT total_price AS revenue,region
FROM sale_detail;

Tabelas com partições dinâmicas são comuns e propensas a data skew. Quando ocorrer data skew, use as soluções a seguir.

Solução

Descrição

Solução 1

Ajuste de parâmetros

Configure parâmetros para otimizar o desempenho.

Solução 2

Otimização por pruning

Identifique partições com alta contagem de registros, remova-as do job principal e insira seus dados separadamente para resolver o problema.

  • Solução 1: Ajuste de parâmetros

    O particionamento dinâmico distribui dados em diferentes partições com base nos valores das colunas, eliminando a necessidade de múltiplas instruções INSERT OVERWRITE. Isso simplifica o código ao lidar com muitas partições, mas também pode causar o problema de arquivos pequenos.

    • Exemplo de data skew

      Considere a seguinte instrução SQL simples como exemplo.

      INSERT INTO TABLE part_test PARTITION(ds) SELECT * FROM  part_test;

      Suponha que o job tenha K instâncias map e N partições de destino.

                                  ds=1
      cfile1                      ds=2
      ...             X           ds=3
      ...
                                  ds=n

      Em casos extremos, isso pode gerar K*N arquivos pequenos, impondo uma carga significativa de gerenciamento ao sistema de arquivos. Para resolver isso, o MaxCompute introduz uma tarefa reduce adicional que direciona dados da mesma partição de destino para uma única (ou poucas) instâncias reduce. Isso evita excesso de arquivos pequenos e é sempre a última tarefa reduce no job. O recurso vem ativado por padrão:

      SET odps.sql.reshuffle.dynamicpt=true;

      Embora essa configuração padrão evite o problema de arquivos pequenos e previna falhas de job devido a excesso de arquivos por instância, ela pode introduzir data skew. A tarefa reduce extra também consome recursos computacionais adicionais; portanto, avalie cuidadosamente essas compensações.

    • Solução

      Ativar set odps.sql.reshuffle.dynamicpt=true; resolve o problema de arquivos pequenos. No entanto, se o número de partições de destino for pequeno e não houver risco de excesso de arquivos pequenos, essa configuração desperdiça recursos computacionais e degrada o desempenho. Desativá-la com set odps.sql.reshuffle.dynamicpt=false; pode melhorar significativamente o desempenho, conforme mostra o exemplo a seguir.

      INSERT overwrite TABLE ads_tb_cornucopia_pool_d PARTITION (ds, lv, tp)
      SELECT /*+ mapjoin(t2) */
          '20150503' AS ds,
          t1.lv AS lv,
          t1.type AS tp
      FROM
          (SELECT  ...
          FROM tbbi.ads_tb_cornucopia_user_d
          WHERE ds = '20150503'
          AND lv IN ('flat', '3rd')
          AND tp = 'T'
          AND pref_cat2_id > 0
          ) t1
      JOIN
          (SELECT ...
          FROM tbbi.ads_tb_cornucopia_auct_d
          WHERE ds = '20150503'
          AND tp = 'T'
          AND is_all = 'N'
          AND cat2_id > 0
          ) t2
      ON t1.pref_cat2_id = t2.cat2_id;

      Ao executar o código anterior com o parâmetro padrão, o job leva cerca de 1 hora e 30 minutos para ser concluído. A tarefa reduce final leva cerca de 1 hora e 20 minutos, representando aproximadamente 90% do tempo total de execução. A tarefa reduce adicional distribui os dados de forma desigual, causando data skew do tipo long-tail.

    Dados históricos mostram que este job gera apenas cerca de duas partições dinâmicas por dia; portanto, defina com segurança set odps.sql.reshuffle.dynamicpt=false;. Com essa alteração, o job é concluído em apenas 9 minutos — uma melhoria drástica no desempenho e no uso de recursos mediante uma simples mudança de configuração.

    Essa otimização aplica-se a qualquer job com poucas partições dinâmicas. Definir o parâmetro odps.sql.reshuffle.dynamicpt como false economiza recursos e melhora o desempenho.

    Otimize qualquer nó que atenda a todas as condições a seguir, independentemente do tempo de execução:

    • O nó usa particionamento dinâmico.

    • O número de partições dinâmicas é menor ou igual a 50.

    • O parâmetro odps.sql.reshuffle.dynamicpt não está definido como false.

    Determine a urgência da otimização para um nó com base no tempo de execução de sua última Fuxi Instance. O campo diag_level indica a urgência da seguinte forma:

    • Last_Fuxi_Inst_Time maior que 30 minutos:Diag_Level=4 ('Severe').

    • Last_Fuxi_Inst_Time entre 20 e 30 minutos:Diag_Level=3 ('High').

    • Last_Fuxi_Inst_Time entre 10 e 20 minutos:Diag_Level=2 ('Medium').

    • Last_Fuxi_Inst_Time menor que 10 minutos:Diag_Level=1 ('Low').

  • Solução 2: Otimização por pruning

    Para resolver data skew ocorrido na fase map ao inserir dados com partições dinâmicas, modifique os parâmetros da fase map conforme mostra o exemplo a seguir:

    SET odps.sql.mapper.split.size=128;
    INSERT OVERWRITE TABLE data_demo3 partition(ds,hh)
    SELECT  *
    FROM    dwd_alsc_ent_shop_info_hi;

    Os resultados mostram que o job realiza uma varredura completa da tabela. Para otimizar ainda mais, desative a tarefa reduce introduzida pelo sistema, conforme mostra o exemplo a seguir:

    SET odps.sql.reshuffle.dynamicpt=false ;
    INSERT OVERWRITE TABLE data_demo3 partition(ds,hh)
    SELECT *
    FROM dwd_alsc_ent_shop_info_hi;

    Para resolver data skew ocorrido na fase map ao usar partições dinâmicas para inserir dados, identifique partições com grande número de registros, remova-as do job principal e insira-as separadamente. Siga estas etapas:

    1. Execute o comando a seguir para encontrar as partições com as maiores contagens de registros.

      SELECT  ds
              ,hh
              ,COUNT(*) AS cnt
      FROM    dwd_alsc_ent_shop_info_hi
      GROUP BY ds
               ,hh
      ORDER BY cnt DESC;

      A tabela a seguir mostra algumas das partições resultantes:

      ds

      hh

      cnt

      20200928

      17

      1052800

      20191017

      17

      1041234

      20210928

      17

      1034332

      20190328

      17

      1000321

      20210504

      1

      19

      20191003

      20

      18

      20200522

      1

      18

      20220504

      1

      18

    2. Use os comandos a seguir para inserir primeiro os dados das partições sem skew e, em seguida, inserir os dados das partições com skew em uma instrução separada.

      SET odps.sql.reshuffle.dynamicpt=false ;
      INSERT OVERWRITE TABLE data_demo3 partition(ds,hh)
      SELECT  *
      FROM    dwd_alsc_ent_shop_info_hi
      WHERE   CONCAT(ds,hh) NOT IN ('2020092817','2019101717','2021092817','2019032817');
      
      set odps.sql.reshuffle.dynamicpt=false ;
      INSERT OVERWRITE TABLE data_demo3 partition(ds,hh)
      SELECT  *
      FROM    dwd_alsc_ent_shop_info_hi
      WHERE   CONCAT(ds,hh) IN ('2020092817','2019101717','2021092817','2019032817');
      
      SELECT  ds
        ,hh,COUNT(*) AS cnt
       FROM dwd_alsc_ent_shop_info_hi
       GROUP BY ds,hh ORDER BY cnt desc;