O DISTRIBUTED MAPJOIN é uma versão otimizada do MAPJOIN. Use o DISTRIBUTED MAPJOIN ao unir uma tabela pequena a uma tabela grande. Tanto o DISTRIBUTED MAPJOIN quanto o MAPJOIN reduzem o shuffling e a ordenação na tabela maior.
Observações de uso
As tabelas a serem unidas devem ter tamanhos diferentes. A tabela grande precisa ter mais de 10 TB, e a tabela pequena deve estar no intervalo de [1 GB, 100 GB].
Os dados da tabela pequena precisam estar distribuídos uniformemente. Se essa tabela apresentar long tails, um único shard concentrará volume excessivo de dados, o que pode causar erros de out of memory (OOM) ou timeouts em chamadas de procedimento remoto (RPC).
Se uma tarefa SQL executar por mais de 20 minutos, use o DISTRIBUTED MAPJOIN para otimização.
-
A execução da tarefa consome muitos recursos. Portanto, não a execute em grupos de cota pequenos.
NotaAltere o grupo de cotas na página Quotas. Para mais informações, consulte Gerencie quotas in the new MaxCompute console.
Uso do DISTRIBUTED MAPJOIN
Para usar DISTRIBUTED MAPJOIN em uma instrução SELECT, adicione a hint /*+distmapjoin(<table_name>(shard_count=<n>,replica_count=<m>))*/ à instrução. Os parâmetros shard_count e replica_count definem o paralelismo das tarefas. Calcule o paralelismo com a fórmula: Parallelism = shard_count × replica_count.
-
Parâmetros
table_name: nome da tabela pequena a ser unida.
-
shard_count=<n>: número de shards de dados da tabela pequena envolvida na união. Esses shards são distribuídos entre os nós de computação para processamento. O valor n indica o número de shards e, geralmente, deve ser um número ímpar.
NotaEspecifique manualmente o parâmetro shard_count. Estime o valor de shard_count com base no tamanho da tabela pequena, considerando que cada nó de shard processe entre 200 MB e 500 MB de dados.
Valores muito altos para shard_count prejudicam o desempenho e a estabilidade do processamento. Valores muito baixos para shard_count podem gerar erros devido ao consumo excessivo de memória.
-
replica_count=<m>: número de réplicas da tabela pequena. O valor m define essa quantidade. Valor padrão: 1.
NotaCrie múltiplas réplicas dos dados no mesmo shard para reduzir o excesso de requisições de acesso e evitar que a falha de um único nó interrompa toda a tarefa. Se houver reinicializações frequentes de nós devido a alto paralelismo ou instabilidade do ambiente, aumente o valor de replica_count. Defina este parâmetro como 2 ou 3.
-
Sintaxe
-- Recommended: Specify shard_count. replica_count defaults to 1. /*+distmapjoin(a(shard_count=5))*/ -- Recommended: Specify both shard_count and replica_count. /*+distmapjoin(a(shard_count=5,replica_count=2))*/ -- Use DISTRIBUTED MAPJOIN for multiple small tables. /*+distmapjoin(a(shard_count=5,replica_count=2),b(shard_count=5,replica_count=2)) */ -- Use DISTRIBUTED MAPJOIN and MAPJOIN together. /*+distmapjoin(a(shard_count=5,replica_count=2)),mapjoin(b)*/
Exemplos
Este exemplo demonstra como usar o DISTRIBUTED MAPJOIN ao inserir dados na tabela particionada tmall_dump_lasttable.
-
Sintaxe JOIN
insert OVERWRITE table tmall_dump_lasttable partition(ds='20211130') select t1.* from ( select nid, doc,type from search_ods.dump_lasttable where ds='20211203' )t1 join ( select distinct item_id from tbcdm.dim_tb_itm where ds='20211130' and bc_type='B' and is_online='Y' )t2 on t1.nid=t2.item_id; -
Sintaxe DISTRIBUTED MAPJOIN
insert OVERWRITE table tmall_dump_lasttable partition (ds='20211130') select /*+ distmapjoin(t2(shard_count=35)) */ t1.* from ( select nid, doc, type from search_ods.dump_lasttable where ds='20211203' )t1 join ( select distinct item_id from tbcdm.dim_tb_itm where ds='20211130' and bc_type='B' and is_online='Y' )t2 on t1.nid=t2.item_id;