Todos os produtos
Search
Central de documentação

ApsaraDB for SelectDB:Usar Bucket Shuffle Join

Última atualização: Jun 29, 2026

O Bucket Shuffle Join otimiza consultas de Join distribuídas no ApsaraDB for SelectDB ao rotear os dados da tabela da direita diretamente para os nós de armazenamento onde os dados da tabela da esquerda já residem, em vez de redistribuir ambas as tabelas pelo cluster. Em comparação com Broadcast Join e Shuffle Join, essa abordagem reduz a sobrecarga de rede e de memória ao tamanho da tabela da direita, independentemente da escala do cluster.

Para mais detalhes sobre o contexto de design, consulte ISSUE 4394.

Funcionamento

O ApsaraDB for SelectDB oferece suporte a duas estratégias comuns de Join distribuído:

  • Broadcast Join: envia a tabela da direita completa para cada HashJoinNode que contém dados da tabela da esquerda. Com três HashJoinNodes, isso representa 3x o volume de dados da tabela da direita, tanto em rede quanto em memória.

  • Shuffle Join: aplica hash em ambas as tabelas e redistribui todas as linhas pelos nós do cluster. A sobrecarga de rede equivale ao tamanho combinado das duas tabelas; a sobrecarga de memória corresponde ao tamanho da tabela da direita.

image.png

O nó frontend (FE) armazena as informações de distribuição de dados de cada tabela. Quando a condição de igualdade de um Join corresponde à coluna de bucket da tabela da esquerda, o FE utiliza esse mapa de distribuição para enviar os dados da tabela da direita diretamente aos nós de armazenamento da tabela da esquerda, sem necessidade de redistribuir a tabela da esquerda.

Como resultado, a sobrecarga de rede e de memória equivale apenas ao tamanho da tabela da direita. Essa vantagem aumenta quando o FE também consegue aplicar partition pruning ou bucket pruning na tabela da esquerda, o que reduz ainda mais o volume de dados que cada nó precisa processar.

Diferença entre Bucket Shuffle Join e Colocate Join: O Bucket Shuffle Join não exige conhecimento prévio nem alterações na distribuição de dados de nenhuma das tabelas. Além disso, evita data skew, pois não impõe regras de distribuição correspondentes entre as duas tabelas. Ele também oferece mais caminhos para o Join Reorder explorar durante o planejamento da consulta.

Conceitos principais

Termo

Definição

Tabela da esquerda

Tabela posicionada à esquerda em uma consulta de Join, utilizada na operação de probe. O Join Reorder pode alterar qual tabela fica à esquerda.

Tabela da direita

Tabela posicionada à direita em uma consulta de Join, utilizada na operação de build. O Join Reorder pode alterar qual tabela fica à direita.

Coluna de bucket

Coluna pela qual as linhas de uma tabela são distribuídas entre buckets no momento da criação da tabela.

Nó FE

Nó frontend que armazena as informações de distribuição de dados de cada tabela e gera os planos de execução de consultas.

Ative o Bucket Shuffle Join

Defina a variável de sessão enable_bucket_shuffle_join como true. O FE selecionará automaticamente o Bucket Shuffle Join para qualquer consulta elegível.

set enable_bucket_shuffle_join = true;

Ao planejar uma consulta distribuída, o FE segue esta ordem de prioridade:

Colocate Join > Bucket Shuffle Join > Broadcast Join > Shuffle Join

Para substituir a seleção automática, utilize um hint:

SELECT * FROM test JOIN [shuffle] baseall ON test.k1 = baseall.k1;

Verifique o uso do Bucket Shuffle Join

Execute explain em uma consulta para inspecionar seu plano de execução. Procure por BUCKET_SHUFFLE no campo join op.

Bucket Shuffle Join aplicado:

|   2:HASH JOIN                                                       |
|   |  join op: INNER JOIN (BUCKET_SHUFFLE)                          |
|   |  hash predicates:                                              |
|   |  colocate: false, reason: table not in the same group          |
|   |  equal join conjunct: `test`.`k1` = `baseall`.`k1`            |

Retorno a outro método de Join (Bucket Shuffle Join não foi aplicado):

|   2:HASH JOIN                                                       |
|   |  join op: INNER JOIN (BROADCAST)                               |
|   |  hash predicates:                                              |
|   |  colocate: false, reason: table not in the same group          |
|   |  equal join conjunct: `test`.`k1` = `baseall`.`k1`            |

Caso o plano exiba BROADCAST ou SHUFFLE em vez de BUCKET_SHUFFLE, revise as regras de planejamento abaixo para identificar o motivo pelo qual a otimização não foi aplicada.

Regras de planejamento

O Bucket Shuffle Join roteia dados com base em atribuições de bucket calculadas por hash. Todos os requisitos de planejamento derivam dessa restrição central: o FE deve ser capaz de determinar, no momento do planejamento, exatamente quais nós de armazenamento contêm cada linha da tabela da esquerda.

Requisito

Detalhe

Condição de join de igualdade

O Join deve usar um predicado de igualdade (=). Condições de intervalo ou de não igualdade não permitem roteamento baseado em hash.

Coluna de bucket da tabela da esquerda na condição de join

O predicado de igualdade deve incluir a coluna de bucket da tabela da esquerda.

Tipos de dados correspondentes

A coluna de bucket da tabela da esquerda e a coluna correspondente da tabela da direita devem ter o mesmo tipo de dados. Tipos incompatíveis impedem o planejamento.

Apenas tabelas OLAP nativas

A tabela da esquerda deve ser uma tabela nativa de processamento analítico online (OLAP). Tabelas externas — como tabelas Open Database Connectivity (ODBC), MySQL ou Elasticsearch — não podem atuar como tabela da esquerda.

Partição única na tabela da esquerda

Para tabelas particionadas, o Bucket Shuffle Join se aplica somente quando a consulta acessa exatamente uma partição. Utilize uma cláusula WHERE para habilitar partition pruning e atender a esse requisito.

Melhor desempenho com tabelas Colocate

Se a tabela da esquerda for uma tabela Colocate, sua distribuição de dados por partição é fixa e totalmente conhecida pelo FE, fornecendo ao Bucket Shuffle Join as informações de roteamento mais precisas possíveis.