O DISTRIBUTED MAPJOIN é uma versão otimizada do MAPJOIN projetada para unir uma tabela de pequeno a médio porte (1 GB a 100 GB) com uma tabela muito grande (maior que 10 TB). Assim como o MAPJOIN, o DISTRIBUTED MAPJOIN reduz o shuffling e a ordenação na tabela grande, mas processa tabelas de junção substancialmente maiores ao fragmentar os dados em vários nós de computação.
Quando usar o DISTRIBUTED MAPJOIN
Use o DISTRIBUTED MAPJOIN quando todas as condições abaixo forem atendidas:
|
Condição |
Requisito |
|
Tamanho da tabela grande |
Maior que 10 TB |
|
Tamanho da tabela pequena |
Entre 1 GB e 100 GB |
|
Distribuição de dados da tabela pequena |
Uniforme (sem skew significativo de dados) |
|
Duração da tarefa SQL |
Mais de 20 minutos sem otimização |
Se a tabela pequena tiver menos de 1 GB, use o MAPJOIN padrão.
Se a distribuição de dados da tabela pequena for desigual, um único shard pode receber dados excessivos, causando erro de falta de memória (OOM) ou timeout de chamada de procedimento remoto (RPC).
A execução do DISTRIBUTED MAPJOIN consome recursos significativos. Não execute tarefas em um grupo de cotas pequeno.
Altere o grupo de cotas na página
Quotas
. Para mais informações, consulte
Gerenciar cotas no novo console do MaxCompute
.
Sintaxe
Para usar o DISTRIBUTED MAPJOIN, adicione o seguinte hint à instrução SELECT:
/*+distmapjoin(<table_name>(shard_count=<n>,replica_count=<m>))*/
Os parâmetros shard_count e replica_count determinam o paralelismo da tarefa. Calcule o paralelismo com a seguinte fórmula:
Parallelism = shard_count x replica_count
Parâmetros
|
Parâmetro |
Padrão |
Descrição |
|
|
-- |
Nome da tabela pequena a ser unida. |
|
|
-- |
Quantidade de shards de dados da tabela pequena. Os shards são distribuídos entre os nós de computação para processamento. Especifique manualmente este valor e defina-o como um número ímpar. |
|
|
1 |
Número de réplicas dos dados da tabela pequena. Aumentar as réplicas reduz a contenção de acesso e evita falhas na tarefa causadas pela queda de um único nó. |
Como escolher o valor de shard_count
Especifique manualmente o valor de shard_count. Estime esse valor com base no tamanho da tabela pequena para que cada shard processe entre 200 MB e 500 MB de dados.
Exemplo: Se a tabela pequena tiver 5 GB:
|
Tamanho alvo do shard |
shard_count |
|
500 MB por shard |
10 |
|
200 MB por shard |
25 |
Comece com um valor nessa faixa e ajuste conforme o desempenho observado.
Definir
shard_count
com um valor muito alto prejudica o desempenho e a estabilidade. - Definir
shard_count
com um valor muito baixo pode causar erro de falta de memória (OOM) devido ao consumo excessivo de memória em cada shard.
Quando aumentar o replica_count
Para reduzir solicitações de acesso excessivas e evitar que a falha de um único nó interrompa toda a tarefa, crie múltiplas réplicas dos dados no mesmo shard. Aumente o valor de replica_count quando:
A tarefa apresentar alto paralelismo e os nós reiniciarem frequentemente.
O desempenho do ambiente estiver instável.
Nesses casos, defina replica_count como 2 ou 3.
Exemplos de hints
-- Specify shard_count only (replica_count defaults to 1).
/*+distmapjoin(a(shard_count=5))*/
-- Specify both shard_count and replica_count.
/*+distmapjoin(a(shard_count=5,replica_count=2))*/
-- Join multiple small tables with DISTRIBUTED MAPJOIN.
/*+distmapjoin(a(shard_count=5,replica_count=2),b(shard_count=5,replica_count=2)) */
-- Combine DISTRIBUTED MAPJOIN and MAPJOIN for different tables.
/*+distmapjoin(a(shard_count=5,replica_count=2)),mapjoin(b)*/
Exemplos
O exemplo a seguir demonstra como usar o DISTRIBUTED MAPJOIN ao inserir dados em uma tabela particionada chamada tmall_dump_lasttable.
Sem DISTRIBUTED MAPJOIN (JOIN padrão)
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;
Com 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;
Neste exemplo, t2 é a tabela pequena derivada de tbcdm.dim_tb_itm. O hint distmapjoin(t2(shard_count=35)) distribui a tabela pequena em 35 shards pelos nós de computação. O replica_count assume o valor padrão 1, pois não foi especificado.