Otimização de JOIN no ClickHouse usando Tabelas em Buckets

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

  1. 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));

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

Tags: ClickHouse database-optimization Distributed-Systems join-optimization SQL

Publicado em 7-28 02:45