As estratégias padrão de junção distribuída — shuffle join e broadcast join — movem dados entre nós de backend (BE) antes da junção, o que aumenta a latência proporcionalmente ao volume de dados. O colocation join elimina essa transferência de rede ao garantir que tabelas no mesmo Colocation Group (CG) sempre armazenem buckets correspondentes no mesmo nó BE. Assim, as junções nas colunas de bucket ocorrem inteiramente em nível local.
A replicação entre clusters (CCR) não sincroniza a propriedade de colocation de uma tabela. Se uma tabela tiver is_being_synced = true, sua propriedade de colocation será removida.
Como funciona
O ApsaraDB for SelectDB agrupa tabelas colocalizadas em um Colocation Group (CG). Todas as tabelas de um CG devem compartilhar o mesmo Colocation Group Schema (CGS). Quando essas condições são atendidas, uma sequência de buckets mapeia cada índice de bucket de forma idêntica para o mesmo nó BE em todas as tabelas do CG. Consequentemente, linhas com os mesmos valores de coluna de bucket residem sempre no mesmo nó BE, viabilizando junções locais.
Exemplo: 8 buckets mapeados para 4 nós BE (A, B, C, D)
+---+ +---+ +---+ +---+ +---+ +---+ +---+ +---+
| 0 | | 1 | | 2 | | 3 | | 4 | | 5 | | 6 | | 7 |
+---+ +---+ +---+ +---+ +---+ +---+ +---+ +---+
| A | | B | | C | | D | | A | | B | | C | | D |
+---+ +---+ +---+ +---+ +---+ +---+ +---+ +---+
Requisitos do CGS — critérios obrigatórios para todas as tabelas de um CG:
|
Requisito |
Detalhes |
|
Tipos de coluna de bucket |
Cada coluna de bucket deve ter o mesmo tipo de dados |
|
Quantidade de colunas de bucket |
A cláusula |
|
Número de buckets |
O valor de |
O que não precisa corresponder:
Quantidade de partições
Escopo ou intervalo da partição
Tipos de coluna de partição
Em tabelas de partição única, cada bucket contém um tablet. Em tabelas com múltiplas partições, cada bucket contém vários tablets — um por partição.
Conceitos principais
Colocation Group (CG): Grupo nomeado composto por uma ou mais tabelas que compartilham o mesmo CGS e a mesma distribuição de buckets. Um CG pertence a um único banco de dados e seu nome é exclusivo nesse banco. Internamente, o SelectDB armazena o CG como
dbId_groupName, mas a interação ocorre apenas pelo nome do grupo.Colocation Group Schema (CGS): Conjunto de restrições de esquema compartilhadas por todas as tabelas de um CG, incluindo tipos de coluna de bucket, quantidade de colunas de bucket e número de buckets.
Stable/Unstable: Um CG é considerado Stable quando todos os seus tablets estão totalmente replicados e nenhuma migração está em andamento. Quando um CG está Unstable (por exemplo, durante reparo de réplicas ou rebalanceamento), o colocation join é temporariamente rebaixado para uma junção comum, o que pode reduzir significativamente o desempenho das consultas.
Crie uma tabela de colocation
Adicione "colocate_with" = "group_name" à cláusula PROPERTIES ao criar uma tabela:
CREATE TABLE tbl (k1 int, v1 int sum)
DISTRIBUTED BY HASH(k1)
BUCKETS 8
PROPERTIES(
"colocate_with" = "group1"
);
Se o CG não existir, o SelectDB o criará automaticamente com esta tabela como primeiro membro.
Caso o CG já exista, o SelectDB verificará se a tabela atende ao CGS. Em caso positivo, a tabela será adicionada e seus tablets serão criados seguindo a distribuição de buckets existente.
Grupos de colocation entre bancos de dados
Para colocalizar tabelas em diferentes bancos de dados, adicione o prefixo __global__ ao nome do grupo:
CREATE TABLE tbl (k1 int, v1 int sum)
DISTRIBUTED BY HASH(k1)
BUCKETS 8
PROPERTIES(
"colocate_with" = "__global__group1"
);
Um CG global não pertence a nenhum banco de dados específico e seu nome é exclusivo em todo o cluster.
Verifique se o colocation join está ativo
Consulte tabelas de colocation da mesma forma que qualquer outra tabela. O SelectDB gera automaticamente um plano de consulta que utiliza colocation join quando as condições são atendidas.
Para verificar, execute DESC (ou EXPLAIN) em uma consulta de junção e inspecione a saída:
-- Create tbl1 and tbl2 in the same CG, then run:
DESC SELECT * FROM tbl1 INNER JOIN tbl2 ON (tbl1.k2 = tbl2.k2);
Colocation join ativo — o nó HASH JOIN exibe colocate: true:
| 2:HASH JOIN |
| | join op: INNER JOIN |
| | colocate: true |
| | `tbl1`.`k2` = `tbl2`.`k2` |
Colocation join inativo — o nó HASH JOIN mostra colocate: false com um motivo, e um nó EXCHANGE aparece:
| 2:HASH JOIN |
| | join op: INNER JOIN (BROADCAST) |
| | colocate: false, reason: group is not stable |
| | `tbl1`.`k2` = `tbl2`.`k2` |
...
| 3:EXCHANGE |
O motivo mais comum para colocate: false é o CG estar Unstable, indicando reparo de réplicas ou rebalanceamento em andamento.
Gerencie grupos de colocation
Visualize informações do CG
Execute o comando abaixo para listar todos os CGs no cluster (requer função de administrador):
SHOW PROC '/colocation_group';
Exemplo de saída:
+-------------+--------------+--------------+------------+----------------+----------+----------+
| GroupId | GroupName | TableIds | BucketsNum | ReplicationNum | DistCols | IsStable |
+-------------+--------------+--------------+------------+----------------+----------+----------+
| 10005.10008 | 10005_group1 | 10007, 10040 | 10 | 3 | int(11) | true |
+-------------+--------------+--------------+------------+----------------+----------+----------+
|
Campo |
Descrição |
|
GroupId |
Identificador exclusivo em todo o cluster no formato |
|
GroupName |
Nome completo do CG |
|
TableIds |
IDs das tabelas no CG |
|
BucketsNum |
Quantidade de buckets |
|
ReplicationNum |
Número de réplicas |
|
DistCols |
Tipos de coluna de bucket |
|
IsStable |
|
Para inspecionar o mapeamento de bucket para BE de um CG específico:
SHOW PROC '/colocation_group/10005.10008';
Exemplo de saída:
+-------------+------------+
| BucketIndex | BackendIds |
+-------------+------------+
| 0 | 10004 |
| 1 | 10003 |
| 2 | 10002 |
| 3 | 10003 |
| 4 | 10002 |
| 5 | 10003 |
| 6 | 10003 |
| 7 | 10003 |
+-------------+------------+
A execução de comandos SHOW PROC requer função de administrador. Usuários comuns não têm permissão para executá-los.
Modifique a propriedade de colocation de uma tabela
Para mover uma tabela para um CG diferente:
ALTER TABLE tbl SET ("colocate_with" = "group2");
Se a tabela ainda não estiver em nenhum CG, o SelectDB valida o esquema e a adiciona ao CG especificado (criando-o se necessário).
Caso a tabela já esteja em um CG, o SelectDB a remove do CG atual e a adiciona ao CG especificado (criando-o se necessário).
Para remover a propriedade de colocation de uma tabela:
ALTER TABLE tbl SET ("colocate_with" = "");
Exclua uma tabela de colocation
Ao executar DROP TABLE, a tabela vai para a lixeira e permanece lá por um dia antes da exclusão permanente. Após a exclusão definitiva da última tabela de um CG, o grupo é removido automaticamente.
Restrições em alterações de esquema
Ao adicionar partições (ADD PARTITION) ou alterar a contagem de réplicas de uma tabela de colocation, o SelectDB valida a alteração em relação ao CGS. Se a mudança violar as restrições do CGS, o SelectDB a rejeitará.
Configuração avançada
Itens de configuração do FE
Os itens de configuração de frontend (FE) a seguir controlam o comportamento do colocation join. Modifique dinamicamente disable_colocate_relocate e disable_colocate_balance usando ADMIN SET CONFIG. Para visualizar os valores atuais ou consultar a sintaxe de configuração, execute HELP ADMIN SHOW CONFIG; e HELP ADMIN SET CONFIG;.
|
Item de configuração |
Padrão |
Descrição |
|
|
|
Desativa o reparo automático de réplicas de colocation. Afeta apenas tabelas de colocation. |
|
|
|
Desativa o rebalanceamento automático de réplicas de colocation. Afeta apenas tabelas de colocation. |
|
|
|
Desativa completamente o recurso de colocation join. |
|
|
|
Habilita a nova lógica de agendamento de réplicas. |
API RESTful HTTP
O SelectDB expõe endpoints de API RESTful HTTP para visualizar e gerenciar CGs. Os nós FE servem todos os endpoints em fe_host:fe_http_port e exigem função de administrador.
Visualize todas as informações de colocation:
GET /api/colocate
Retorna informações de colocation no formato JSON:
{
"msg": "success",
"code": 0,
"data": {
"infos": [
["10003.12002", "10003_group1", "10037, 10043", "1", "1", "int(11)", "true"]
],
"unstableGroupIds": [],
"allGroupIds": [{"dbId": 10003, "grpId": 12002}]
},
"count": 0
}
Marcar um CG como Stable:
DELETE /api/colocate/group_stable?db_id=10005&group_id=10008
Return value: 200
Marcar um CG como Unstable:
POST /api/colocate/group_stable?db_id=10005&group_id=10008
Return value: 200
Configure manualmente a distribuição de buckets:
POST /api/colocate/bucketseq?db_id=10005&group_id=10008
Body:
[[10004],[10003],[10002],[10003],[10002],[10003],[10003],[10003],[10003],[10002]]
Return value: 200
O corpo da requisição é um array aninhado em que cada elemento lista os IDs dos nós BE para aquele índice de bucket.
Antes de chamar o endpoint bucketseq, defina tanto disable_colocate_relocate quanto disable_colocate_balance como true. Caso contrário, o reparo automático e o rebalanceamento do sistema substituirão sua configuração manual.
Rebalanceamento e reparo de réplicas
Funcionamento do reparo de réplicas
As réplicas de colocation ficam fixadas em nós BE específicos. Se um nó BE ficar inativo ou for descomissionado, o SelectDB seleciona o nó BE disponível com menor carga como substituto e inicia o reparo de todos os tablets afetados. Durante essa migração, o CG é marcado como Unstable.
Funcionamento do rebalanceamento de réplicas
Para tabelas comuns, o SelectDB equilibra as réplicas individualmente, encontrando um BE adequado para cada réplica de forma independente. Nas tabelas de colocation, o equilíbrio ocorre no nível do bucket: todas as réplicas de um bucket migram juntas para preservar a garantia de colocalização.
O algoritmo de balanceamento distribui a sequência de buckets entre os nós BE com base na contagem de réplicas, e não no tamanho real dos dados.
O algoritmo de balanceamento atual pode não distribuir a carga uniformemente em implantações heterogêneas, onde os nós BE possuem diferentes capacidades de disco, quantidades de discos ou tipos de disco (SSD vs. HDD). Nesses ambientes, nós de pequena e grande capacidade podem acabar mantendo o mesmo número de réplicas.
Se um CG se tornar Unstable e você quiser evitar que o rebalanceamento automático interrompa repetidamente as consultas, defina disable_colocate_balance como true. Altere de volta para false quando estiver pronto para permitir o rebalanceamento novamente.