Todos os produtos
Search
Central de documentação

MaxCompute:DISTRIBUTED MAPJOIN

Última atualização: Aug 20, 2026

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.

    Nota

    Altere 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.

      Nota
      • Especifique 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.

      Nota

      Crie 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;