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:
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:
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.
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.
Use o log StdOut para visualizar o grafo de execução do job.
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.
Encontre a URL do Logview nos logs operacionais da tarefa. Para mais informações, consulte Pontos de entrada do Logview.

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.

A tarefa
R31_26_27apresenta a maior latência. Clique emR31_26_27para acessar a página de detalhes da instância, conforme mostra a figura a seguir.
Latency: {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 é13se a máxima é26 minutos e 40s. Ordene as instâncias porLatencyem 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 que26ssão classificadas como instâncias de long-tail (Long-Tails). Neste caso, 21 instâncias têm latência superior a26s. A presença de instâncias Long-Tails não indica necessariamente data skew. Compare também os valoresavgemaxda latência da instância. Para tarefas em que o valormaxsupera muito o valoravg, indicando data skew severo, otimize essas tarefas.Clique em
na coluna StdOut para visualizar o log StdOut, conforme mostra o exemplo a seguir.
Após identificar o problema, na aba Job Details, clique com o botão direito em
R31_26_27e 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.
Examine StreamLineWriter21, etapa anterior aStreamLineRead22, para localizar as chaves causadoras do data skew:new_uri_path_structure,cookie_x5check_useridecookie_userid. Esse processo também ajuda a localizar o fragmento SQL causador do data skew.
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, enquantot2et3sã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
orpara conectar múltiplas condições. Para calcular um produto cartesiano, omita a cláusulaone use a condiçãomapjoin 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 joinnã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 et1é 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_uidpossui 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âmetroodps.sql.skewinfo. O parâmetroodps.sql.skewinfoconfigura 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 evalueé 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;NotaEspecificar 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.c0deve corresponder ao tipo de dados deb.c0, e o tipo de dados dea.c1deve corresponder ao tipo de dados deb.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 comandoset 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
intcontendo 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 enumber. Adicionar a chavenumberreduz o skew causado por um único ID de usuário para1/Ndo 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çaconcatdo ID com um inteiro aleatório de um intervalo predefinido (como 0 a 1000); para outros IDs, mantenha o valor original. Useeleme_uid_joincomo 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:
|
Nº |
Solução |
Descrição |
|
Solução 1 |
Definir parâmetro para lidar com data skew em |
|
|
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çãoCOUNT. Isso corresponde a um plano de execuçãoM->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-1aT-365repetidamente 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
acom uma data de atualização marcada. As tarefas online subsequentes passam então a fazer join da tabelaT-2acom a tabelatable_xxx_die 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áriashopiddiminui 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:
|
Nº |
Solução |
Descrição |
|
Solução 1 |
Ajuste de parâmetros |
|
|
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 |
-
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_idtiverem 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:
|
Nº |
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
iteratee mesclando apenas N elementos na fasemerge. 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 EXPLODEpara 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.
|
Nº |
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=nEm casos extremos, isso pode gerar
K*Narquivos 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 comset 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.dynamicptcomofalseeconomiza 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.dynamicptnão está definido comofalse.
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_levelindica a urgência da seguinte forma:Last_Fuxi_Inst_Timemaior que 30 minutos:Diag_Level=4 ('Severe').Last_Fuxi_Inst_Timeentre 20 e 30 minutos:Diag_Level=3 ('High').Last_Fuxi_Inst_Timeentre 10 e 20 minutos:Diag_Level=2 ('Medium').Last_Fuxi_Inst_Timemenor 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:
-
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
-
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;
-