- Introdução
O ClickHouse apresenta desempenho excepcional em consultas de tabela única, porém quando se trata de operações JOIN, principalmente entre tabelas grandes, frequentemente enfrentamos problemas como estouro de memória e timeout.
Este cenário se torna ainda mais crítico quando ambas as tabelas A e B são tabelas distribuídas. Considere o exemplo comum de junção entre uma tabela de transações (transactions_all) e uma tabela de clientes (customers_all), ambas sendo tabelas distribuídas.
Este artigo apresenta testes de bucket no ClickHouse, onde realizamos hash no campo de identificação do cliante (customer_id) para particionamento mensal, garantindo que registros com o mesmo customer_id sejam alocados no mesmo shard, aumentando a eficiência do JOIN e reduzindo o I/O de rede.
- Estrutura das Tabelas
Tabela Local de Transações
CREATE TABLE analytics.transaction_local on cluster cluster_3shards_2replicas
(
`data_dt` String,
`categoria` String,
`tipo_transacao` String,
`estabelecimento_id` Int32,
`customer_id` Int32,
`status_pagamento` Int8,
`usuario_uid` Int64,
`timestamp_transacao` Int32,
`valor_total` Int32,
`moeda` Int32,
`categoria_item` Int32,
`direcao` Int8,
`produto_id` Int32,
`quantidade_produto` Int64,
`codigo_unico` String,
`nivel_usuario` Int32,
`assinante` Int32,
`saldo` Int32,
`indicador` Int32,
`status_pagamento_lo` Int8,
`dados_adicionais` String,
`status` String,
`sincronizado` String,
`atualizacao` String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/transaction_local', '{replica}')
PARTITION BY toYYYYMM(toDate(data_dt))
ORDER BY (categoria, estabelecimento_id, timestamp_transacao, produto_id, codigo_unico, status_pagamento, intHash64(customer_id))
SAMPLE BY intHash64(customer_id)
SETTINGS index_granularity = 8192;
Tabela Distribuída de Transações
CREATE TABLE analytics.transaction_all on cluster cluster_3shards_2replicas
as analytics.transaction_local
ENGINE = Distributed('cluster_3shards_2replicas', 'analytics', 'transaction_local', intHash64(customer_id));
Tabela Local de Clientes
CREATE TABLE analytics.customer_local on cluster cluster_3shards_2replicas
(
`data_dt` String,
`categoria` String,
`customer_id` Int32,
`identificador` String,
`ultima_atividade` Int64,
`ip_acesso` Int64,
`canal` String,
`vinculo` String,
`dispositivo` String,
`criacao_conta` Int64,
`ultimo_login` Int64,
`conta_ativa` Int32,
`grupo` Int32,
`situacao` Int8,
`hash_uuid_low` String,
`hash_uuid_high` String,
`pagamento_status` Int32,
`base_grupo` Int32,
`estado` String,
`fuso_horario` String,
`campo_x` String,
`campo_y` String
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/customer_local', '{replica}')
PARTITION BY toYYYYMM(toDate(data_dt))
ORDER BY (categoria, customer_id, intHash64(customer_id), ultima_atividade, ip_acesso, canal, vinculo, dispositivo, criacao_conta, ultimo_login, conta_ativa, grupo, situacao, hash_uuid_low, hash_uuid_high, pagamento_status, base_grupo)
SAMPLE BY intHash64(customer_id)
SETTINGS index_granularity = 8192;
Tabela Distribuída de Clientes
CREATE TABLE analytics.customer_all on cluster cluster_3shards_2replicas
as analytics.customer_local
ENGINE = Distributed('cluster_3shards_2replicas', 'analytics', 'customer_local', intHash64(customer_id));
- Resultados das Consultas
Consulta JOIN sem bucket 1Consulta SQL de teste 1: Associação por cliente com filtro de data (247 segundos)
SELECT COUNT(a.customer_id)
FROM (SELECT customer_id, data_dt FROM analytics.transaction_all) a
INNER JOIN (SELECT customer_id, data_dt FROM analytics.customer_all) b
ON a.customer_id = b.customer_id AND a.data_dt = b.data_dt
GROUP BY data_dt
ORDER BY data_dt
LIMIT 10;
Consulta JOIN com bucket 1Consulta SQL de teste 1: Associação por cliente com filtro de data (forma padrão com bucket, ainda lenta: 298 segundos)
SELECT data_dt, COUNT(a.customer_id)
FROM (SELECT customer_id, data_dt FROM analytics.transaction_all) a
INNER JOIN (SELECT customer_id, data_dt FROM analytics.customer_all) b
ON a.customer_id = b.customer_id AND a.data_dt = b.data_dt
GROUP BY data_dt
ORDER BY data_dt
LIMIT 10;
(Usando a sintaxe otimizada para bucket: 58 segundos)
SET distributed_product_mode = 'global';
SELECT a.data_dt, COUNT(a.customer_id)
FROM analytics.transaction_all a
GLOBAL INNER JOIN analytics.customer_all b USING(customer_id, data_dt)
GROUP BY a.data_dt
ORDER BY a.data_dt
LIMIT 10;
Consulta JOIN sem bucket 2Consulta SQL de teste 2: Associação entre tabela grande e pequena por cliente:
SELECT COUNT(a.customer_id)
FROM (SELECT customer_id FROM analytics.transaction_all) a
INNER JOIN (SELECT customer_id FROM analytics.customer_all) b
ON a.customer_id = b.customer_id;
(Sem bucket: 28 segundos)
Consulta JOIN com bucket 2
SELECT COUNT(a.customer_id)
FROM analytics.transaction_all a
GLOBAL INNER JOIN analytics.customer_all b USING(customer_id)
LIMIT 10;
(Com bucket: 43 segundos)
Importante: Ao executar consultas, evite envolver a tabela B com parênteses como SELECT ... FROM (SELECT ...) a JOIN (SELECT ...) b ..., pois isso anula a otimização. Isso está relacionado ao解析ador do ClickHouse.
Conclusão
Ao realizar JOINs entre duas tabelas grandes, o ganho de desempenho é significativo (especialmente ao escolhre campos frequentes para hash, como ID de cliente). Em associações entre tabelas de tamanhos muito diferentes com condições simples, o impacto pode ser menos expressivo.