04 Reconstrução do pipeline dbt
Reconstrua o pipeline Medallion no ClickHouse com o dbt-clickhouse — modelos incrementais delete_insert, ReplacingMergeTree e views materializadas atualizáveis — e crie o dicionário de zonas.
Ponto de partida
Módulo 03 concluído: o serviço do ClickHouse Cloud está ativo e acessível, e o
setup.sh gravou .clickhouse_state no disco com CLICKHOUSE_HOST e
CLICKHOUSE_PORT. Todas as tabelas de destino e views de preparação existem —
default.trips_raw, as duas views de preparação (stg_trips, stg_taxi_zones), as seis
tabelas de analytics (fact_trips, agg_hourly_zone_trips, dim_taxi_zones,
dim_payment_type, dim_vendor, dim_date) e a view materializada atualizável
analytics.mv_live_trip_feed —, pois o dbt run do módulo 03 já criou as sete. Todas as
tabelas de analytics ainda estão vazias, exceto mv_live_trip_feed, que já contém a
única linha de snapshot daquela execução. Fora isso, somente default.trips_raw contém
dados: cerca de 50 milhões de linhas. O produtor do Snowflake continua em execução,
portanto o ClickHouse está atrasado em relação ao Snowflake aproximadamente pelo tempo
decorrido desde a janela de migração. Reserve cerca de 30 minutos.
Por quê
O módulo 03 provou que o ClickHouse consegue armazenar 50 milhões de linhas. Ele não provou que o pipeline consegue ser executado no ClickHouse — com as views de preparação, a tabela de fatos incremental, as recargas das dimensões e os testes que detectam um modelo com falha antes que qualquer parceiro o veja. É isso que este módulo reconstrói: os mesmos modelos Medallion do módulo 01, expressos com dbt-clickhouse em vez de dbt-snowflake e executados nas tabelas que o módulo 03 já criou.
Nada muda na lógica dos modelos — stg_trips continua convertendo tipos e extraindo
JSON, int_trips_enriched continua fazendo junções com as dimensões e fact_trips
continua terminando com uma linha por corrida. O que muda é a camada de materialização
subjacente: sem MERGE INTO, sem Task do Snowflake e sem cluster_by. Este módulo mostra
que a migração não é apenas uma carga pontual de dados: o pipeline que a equipe de um
parceiro executa diariamente, no mesmo cronograma e condicionado à mesma execução de
dbt test, continua funcionando quando o warehouse subjacente muda.
Conceitos — por baixo dos panos
O conjunto de modelos difere do pipeline do Snowflake de quatro maneiras. A referência
completa de configuração está em
dbt no ClickHouse — esta é a
versão resumida de que você precisa antes de executar dbt run na Etapa 1. Para conhecer
o pipeline de origem substituído por esses modelos, consulte
dbt no Snowflake.
1. delete_insert substitui MERGE. O ClickHouse não possui uma instrução
MERGE INTO. Enquanto o pipeline do Snowflake usava incremental_strategy: merge para
fazer o upsert de fact_trips e agg_hourly_zone_trips, os modelos do ClickHouse usam
incremental_strategy: delete_insert: o dbt exclui as linhas que correspondem à
unique_key do lote recebido e, em seguida, insere o lote. Para fact_trips, a
unique_key é trip_id, e o filtro incremental usa updated_at como marca-d'água, em
vez de pickup_at — uma correção de tarifa reinsere o mesmo trip_id com o mesmo
pickup_at, mas com um updated_at mais recente; por isso, usar pickup_at como
marca-d'água faria com que a correção passasse despercebida.
2. ReplacingMergeTree é a rede de proteção sob delete_insert, não um substituto
dele. Os dois modelos incrementais são declarados como ReplacingMergeTree(updated_at).
Quando uma execução de delete_insert termina normalmente, a tabela já contém uma linha
por chave e o mecanismo não tem nada para limpar. Se uma execução for interrompida no
meio — por exemplo, por uma falha depois da exclusão e antes da inserção —, as mesclagens
em segundo plano acabam desduplicando todas as linhas remanescentes e mantêm aquela com o
maior updated_at. Nunca dependa somente de ReplacingMergeTree para fazer a
desduplicação que cabe a delete_insert: as mesclagens em segundo plano são assíncronas
e, em uma tabela desse tamanho, podem ficar atrasadas por minutos ou horas.
3. Views materializadas atualizáveis substituem a tarefa agendada. O pipeline do
Snowflake usava uma Task agendada que executava um procedimento armazenado para manter
uma agregação contínua atualizada. Em vez disso, o projeto dbt do ClickHouse declara
mv_live_trip_feed com materialized = 'materialized_view' e
engine = 'ReplacingMergeTree(refreshed_at)' — criado por dbt run como uma view
materializada atualizável, em contraste com mais de 30 linhas de DDL CREATE TASK do
Snowflake para obter o mesmo efeito. Este módulo não ativa um intervalo de atualização;
isso exigiria uma instrução manual
ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ..., que o laboratório
não automatiza (veja o motivo no módulo 05).
4. mv_live_trip_feed não possui nenhum equivalente no Snowflake. Ela não é a
tradução de um modelo existente: é um novo recurso proporcionado pela migração. Uma view
materializada padrão do ClickHouse é acionada uma vez por INSERT e enxerga somente as
linhas daquele lote, portanto não consegue calcular corretamente uma agregação referente
a todo o período, como o total de corridas ou a tarifa média. Já uma view materializada
ATUALIZÁVEL executa novamente sua consulta inteira — neste caso,
SELECT ... FROM {{ ref('fact_trips') }} — de acordo com um cronograma, de modo que cada
atualização enxerga a tabela completa. A parte do Snowflake deste workshop nunca ofereceu
essa possibilidade.
Os modelos efetivamente criados pelo dbt e o momento em que cada um recebe dados:
| Modelo | Camada | Materialização | Preenchido | Observações |
|---|---|---|---|---|
stg_trips | preparação | View | Etapa 1 (toda execução) | Converte tipos; usa JSONExtract* para trip_metadata |
stg_taxi_zones | preparação | View | Etapa 1 (toda execução) | Repasse da dimensão de zonas |
int_trips_enriched | preparação | Efêmero | — (incorporado como uma CTE) | Todas as junções com dimensões; sem tabela física |
fact_trips | analytics | Incremental | Etapa 1 | delete_insert com chave trip_id e marca-d'água em updated_at; ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id) |
agg_hourly_zone_trips | analytics | Incremental | Módulo 05, após a virada | Janela contínua de 2 horas; o filtro incremental corresponde apenas às linhas do produtor ativo — veja a Etapa 1 abaixo |
dim_taxi_zones | analytics | Tabela | Etapa 1 | Recarga completa por execução; origem do dicionário de zonas na Etapa 2 |
dim_payment_type | analytics | Tabela | Etapa 1 | Recarga completa por execução |
dim_vendor | analytics | Tabela | Etapa 1 | Recarga completa por execução |
dim_date | analytics | Tabela | Etapa 1 | Estrutura de datas estática, de 2009 a 2029; recarga completa por execução |
mv_live_trip_feed | analytics | View materializada (atualizável) | dbt run do módulo 03; intervalo de atualização nunca ativado | Sem equivalente no Snowflake — veja o ponto 4 acima |
Etapa 1 — Preencher a camada de analytics
Ative o venv do dbt-clickhouse criado no módulo 00 e execute dbt run pela segunda vez.
O módulo 03 já o executou uma vez com tabelas vazias para criar os esquemas; esta execução
tem dados reais por trás — trips_raw agora contém 50 milhões de linhas.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .venv/bin/activate
source .env && source .clickhouse_state
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt runO dbt lê a conexão do ClickHouse em ~/.dbt/profiles.yml, com base em
dbt/nyc_taxi_dbt_ch/profiles.yml.example — o mesmo perfil que o módulo 03 usou para
criar os esquemas vazios.
Resultado esperado: aproximadamente 8 a 12 minutos (50 milhões de linhas processadas pelos modelos incrementais).
agg_hourly_zone_trips estará vazia depois desta execução — isso é esperado, não uma
falha. Seu filtro incremental é WHERE pickup_at >= now() - INTERVAL 2 HOUR, que
corresponde apenas às linhas gravadas pelo produtor ativo. Todas as linhas que você acabou
de migrar são históricas, portanto nenhuma está dentro de uma janela de 2 horas medida a
partir do momento atual. Essa tabela permanece vazia até que a virada do módulo 05 inicie
o produtor do ClickHouse — não perca tempo investigando isso como se fosse uma falha do
pipeline.
Em seguida, execute o conjunto de testes:
dbt testResultado esperado: todos os testes passam.
Verificação:
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265Etapa 2 — Criar o dicionário de zonas
analytics.dim_taxi_zones agora está preenchida com todas as 265 zonas da NYC TLC. Crie
analytics.taxi_zones_dict, um dicionário em memória baseado nessa tabela, para que as
consultas posteriores possam encontrar o borough de uma zona com dictGet(), em vez de
usar uma JOIN.
O que um dicionário oferece em comparação com uma junção. Um dicionário é carregado
uma vez na memória e permanece quente; na prática, uma busca nele não tem custo em todas
as consultas seguintes. Uma JOIN com dim_taxi_zones relê e volta a combinar a tabela
de dimensão toda vez que é executada. Para uma tabela de referência pequena e raramente
alterada como esta — 265 linhas, totalmente recarregadas por cada dbt run —, a vantagem
é inequívoca. A execução de benchmark do módulo 05 consulta taxi_zones_dict diretamente
com dictGet; portanto, esta etapa é uma dependência obrigatória daquele módulo, não um
complemento opcional.
Carregue os detalhes da conexão e, em seguida, aplique o DDL do dicionário por meio do
clickhouse-client ou da API HTTP — escolha a opção disponível no seu ambiente:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .clickhouse_state
# Via clickhouse-client
clickhouse-client \
--host "${CLICKHOUSE_HOST}" \
--port 9440 \
--user default \
--password "${CLICKHOUSE_PASSWORD}" \
--secure \
--multiquery \
< scripts/04_create_dictionary.sql
# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
--user "default:${CLICKHOUSE_PASSWORD}" \
--data-binary @scripts/04_create_dictionary.sqlVerificação:
-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));
-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';Como verificar se você terminou
SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 millionSELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0O valor zero está correto aqui; não indica uma falha. O filtro incremental de
agg_hourly_zone_trips (WHERE pickup_at >= now() - INTERVAL 2 HOUR) corresponde apenas
às linhas gravadas pelo produtor ativo, e todas as linhas atualmente no ClickHouse são
dados históricos transferidos pelo script de migração do módulo 03 — nenhuma tem menos
de 2 horas em relação a now(). A tabela só começa a ser preenchida depois que a virada
do módulo 05 inicia o produtor do ClickHouse. Até lá, qualquer gráfico de dashboard
baseado nessa tabela não exibirá dados, o que é esperado e não precisa ser investigado.
SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt testResultado esperado: todos os testes passam.
-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));Estado final
A camada de analytics está preenchida e testada: fact_trips contém cerca de 50 milhões
de linhas; dim_taxi_zones, dim_payment_type, dim_vendor e dim_date estão
totalmente carregadas; dbt test passa de ponta a ponta; e
analytics.taxi_zones_dict está ativo e retorna boroughs por meio de dictGet().
agg_hourly_zone_trips continua vazia — por decisão de projeto, não por defeito — e
permanece assim até a virada do módulo 05.
Dashboards, o benchmark ClickHouse versus Snowflake e a virada fazem parte do módulo 05, e não deste.
O produtor do Snowflake continua em execução, e a lacuna entre o Snowflake e o ClickHouse continua aberta. Nada neste módulo afetou o produtor ou o script de migração — o módulo 05 fecha essa lacuna de maneira deliberada, no mesmo processo controlado em duas passagens apresentado no módulo 03. Não interrompa o produtor agora.
03 Provisionamento e migração
Provisione o ClickHouse Cloud com Terraform, crie as tabelas de destino a partir do seu plano e transfira 50 milhões de linhas com um script de migração retomável em Python.
05 Benchmark e virada
Reconstrua os dashboards no ClickHouse, compare as sete consultas nos dois mecanismos, faça a virada do produtor, verifique a paridade e desmonte os ambientes.