Exemplo resolvido: um plano concluído
Um plano de migração preenchido para a carga de trabalho NYC Taxi, que você pode comparar com o seu depois de escrevê-lo.
Esta é a resposta totalmente desenvolvida das cinco planilhas
(1,
2,
3,
4,
5) aplicada à carga
de trabalho NYC Taxi. Use-a para:
- Conferir as respostas das planilhas após concluir cada seção
- Entender o raciocínio por trás das decisões implementadas na Parte 3
- Comparar suas escolhas com a tabela de alinhamento de decisões da Parte 3, caso sejam diferentes
Este é o gabarito — não o preencha como se fosse seu plano. Em vez disso, preencha
migration-plan.md.
| Métrica | Valor |
|---|
| Total de tabelas | 7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE) |
| Total de views | 2 (STG_TRIPS, STG_TAXI_ZONES) |
| Streams | 1 (TRIPS_CDC_STREAM em TRIPS_RAW) |
| Tasks | 2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK) |
| Total de linhas em TRIPS_RAW | ~50.000.000 |
| Intervalo de datas | Janela contínua de 4 anos com término no momento da preparação |
| Colunas VARIANT | 1 (TRIPS_RAW.TRIP_METADATA) |
| Usos de QUALIFY detectados | 1 (consulta Q3) |
| Usos de MERGE INTO detectados | 2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK) |
| Objeto | Tipo | Esquema | Linhas | Grau de complexidade | Observações |
|---|
trips_raw | Tabela | raw | ~50 mi | B | RMT com coluna de versão _synced_at; a sobreposição da carga em massa e do CDC exige desduplicação; stg_trips deve usar FINAL |
stg_trips | View do dbt | staging | — | B | JSONExtract para TRIP_METADATA; exige testes dos caminhos JSON |
stg_taxi_zones | View do dbt | staging | — | A | Repasse; trivial |
int_trips_enriched | Efêmero do dbt | staging | — | A | CTE; diferenças de SQL tratadas nos modelos-pai |
fact_trips | Incremental do dbt | analytics | ~50 mi | C | Mecanismo RMT; delete_insert; reescrita de QUALIFY; FINAL obrigatória |
agg_hourly_zone_trips | Incremental do dbt | analytics | ~140 mil | B | RMT; janela contínua de novo cálculo de 2 h; teste o limite da partição com atenção |
dim_taxi_zones | Tabela do dbt | analytics | 265 | A | Referência estática; recarga completa; trivial |
dim_payment_type | Tabela do dbt | analytics | 6 | A | Referência estática; trivial |
dim_vendor | Tabela do dbt | analytics | 3 | A | Referência estática; trivial |
taxi_zones_dict | Dicionário | analytics | 265 | B | Sintaxe específica do ClickHouse; dictGet() no momento da consulta |
mv_hourly_revenue | MV atualizável | analytics | — | B | Sintaxe REFRESH EVERY; verificar a substituição atômica |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Stream + Task do Snowflake | — | — | D | Sem equivalente no ClickHouse; substituídos pela virada direta do produtor na Parte 3 |
| Tabela | Mecanismo | Coluna de versão | Raciocínio |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | O script de migração em Python (scripts/02_migrate_trips.py) pode repetir um lote e inserir novamente o mesmo trip_id. Após a virada, o produtor ativo também pode repetir uma operação em caso de falhas temporárias. _synced_at DateTime DEFAULT now() é definido automaticamente no INSERT — uma repetição posterior tem um carimbo de data/hora maior, por isso o RMT mantém a gravação mais recente. stg_trips consulta com FINAL para impor a desduplicação antes da execução de qualquer modelo posterior. |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | As corridas podem ser corrigidas (ajustes de tarifa e mudanças de status). O mesmo trip_id é inserido novamente com valores atualizados. updated_at cresce monotonicamente a cada correção — o maior valor prevalece durante a desduplicação do RMT. Sempre consulte com FINAL. |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | O dbt recalcula as 2 horas anteriores e faz uma nova inserção. Sem RMT, as agregações antigas e novas se acumulam e são contadas duas vezes. updated_at definido como now() em cada execução do dbt garante que os valores mais recentes prevaleçam. |
dim_taxi_zones | MergeTree() | — | Recarga completa pelo dbt (troca atômica de tabelas, com reconstrução completa). As duplicatas não podem se acumular. Nenhuma desduplicação é necessária. |
dim_payment_type | MergeTree() | — | Mesmo motivo — recarga completa. |
dim_vendor | MergeTree() | — | Mesmo motivo — recarga completa. |
mv_hourly_revenue | MergeTree() | — | A MV ATUALIZÁVEL substitui atomicamente todo o conjunto de resultados em cada REFRESH. Sem upserts. |
| Tabela | ORDER BY | Raciocínio |
|---|
trips_raw | (pickup_at, trip_id) | As varreduras por intervalo de tempo filtram primeiro por pickup_at. trip_id é a chave de desduplicação do RMT — ela precisa estar em ORDER BY para o RMT identificar quais linhas são duplicadas. pickup_at vem primeiro porque predominam as varreduras analíticas por intervalo; trip_id vem por último porque tem alta cardinalidade e serve apenas para distinguir valores exclusivos. |
fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | As 7 consultas analíticas filtram por pickup_at. O prefixo do mês agrupa em blocos adjacentes os dados de cada mês do calendário — isso permite uma eliminação ampla de blocos em agregações mensais sem adicionar PARTITION BY. trip_id vem por último para garantir a exclusividade no RMT sem prejudicar a eliminação de blocos. |
agg_hourly_zone_trips | (hour_bucket, zone_id) | Q6 (e todas as consultas de agregação) filtra por hour_bucket e zone_id. hour_bucket tem cerca de 35 mil valores distintos; zone_id tem 265. hour_bucket vem primeiro porque as varreduras por intervalo de tempo são o principal padrão de acesso. zone_id vem depois para a filtragem secundária. |
dim_taxi_zones | (location_id) | 265 linhas = um grânulo. ORDER BY é irrelevante para o desempenho. Usar location_id, a chave de junção, é convencional e facilita a leitura. |
| Coluna | Tipo no Snowflake | Tipo no ClickHouse | Justificativa da decisão |
|---|
TRIP_METADATA | VARIANT | String | Preserva o JSON bruto exatamente. JSONExtract* trata caminhos arbitrários no momento da consulta. Map(String,String) perde estruturas aninhadas; Tuple exige um esquema fixo. String é a escolha segura para JSON arbitrário. |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | A precisão de milissegundos é suficiente para os carimbos de data/hora das corridas. Nanossegundos (9) são excessivos. 'UTC' torna explícito o fuso horário e evita surpresas relacionadas ao horário de verão nas agregações por intervalo de tempo. |
PICKUP_LOCATION_ID | INTEGER | UInt16 | Valores de 1 a 265. O máximo de UInt8 é 255 (pequeno demais). O máximo de UInt16 é 65535 (correto). 2 bytes em vez de 4 para Int32 — economiza cerca de 95 MB sem compactação por coluna em 50 milhões de linhas. |
VENDOR_ID | INTEGER | UInt8 | Valores de 1 a 3. O máximo de UInt8 é 255 — correto. 1 byte por linha. |
DRIVER_RATING | FLOAT | Nullable(Float32) | É NULL com frequência (nem todas as corridas têm avaliação). Nullable preserva a semântica correta de valores nulos. Float32 é suficiente para o intervalo de 1,0 a 5,0. Float64 desperdiçaria armazenamento sem acrescentar uma precisão significativa. |
UPDATED_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | Coluna de versão do ReplacingMergeTree. Deve usar DateTime64, não DateTime — duas correções no mesmo segundo seriam não determinísticas com precisão de segundos. A precisão de milissegundos garante a ordenação correta da desduplicação. |
| Expressão do Snowflake | Equivalente no ClickHouse |
|---|
DATE_TRUNC('hour', pickup_at) | toStartOfHour(pickup_at) |
DATEADD('day', -7, CURRENT_DATE) | today() - 7 |
DATEDIFF('minute', pickup_at, dropoff_at) | dateDiff('minute', pickup_at, dropoff_at) |
TRIP_METADATA:driver.rating::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | Incremental delete_insert do dbt — exclui as linhas com chaves correspondentes e depois insere todas as novas linhas |
| Onda | Objetos | Dependências | Observações |
|---|
| Onda 0 | trips_raw (esquema), dim_taxi_zones, dim_payment_type, dim_vendor | Nenhuma | O dbt cria tabelas vazias. As tabelas de dimensão são preenchidas imediatamente a partir dos dados de referência estáticos (sem dependência das corridas). Execute: dbt run --select trips_raw dim_* |
| Onda 1 | Carga em massa com Python (scripts/02_migrate_trips.py) | Onda 0 (o esquema de trips_raw deve existir) | 50 milhões de linhas de TRIPS_RAW no Snowflake. Pode ser retomada com --resume. Verifique a contagem de linhas com scripts/01_verify_migration.sh. Cerca de 40 a 50 min. |
| Onda 2 | stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_trips | Onda 1 concluída (trips_raw preenchida) + Onda 0 (tabelas de dimensão existentes) | dbt run completo. stg_trips lê trips_raw; int_trips_enriched faz junções com as dimensões; fact_trips e agg_hourly_zone_trips são construídas sobre essas tabelas. |
| Onda 3 | taxi_zones_dict, mv_live_trip_feed | Onda 2 (dim_taxi_zones preenchida para o dicionário; fact_trips preenchida para a MV) | Dicionário criado por meio de scripts/04_create_dictionary.sql. MV atualizável criada pelo modelo dbt; ativar seu intervalo de atualização depois disso é uma etapa manual com ALTER TABLE ... MODIFY REFRESH, não algo que o dbt executa automaticamente. |
| Onda 4 | Virada do produtor (scripts/03_cutover.sh) | Onda 1 concluída (carga em massa verificada) + Onda 2 concluída (camada de analytics criada) | Interrompa o produtor do Snowflake; inicie o produtor do ClickHouse, que grava diretamente no ClickHouse Cloud; execute dbt run para preencher agg_hourly_zone_trips com dados ativos. |
| Objeto | Risco | Método de verificação |
|---|
fact_trips | Consultas sem FINAL fazem contagem excessiva durante o atraso das mesclagens. O intervalo de partições de delete_insert deve se limitar ao prefixo de ORDER BY para não excluir partições fora do destino. | SELECT COUNT(*) FINAL corresponde ao Snowflake ± o atraso do CDC. Execute dbt test. Compare os resultados de Q3 entre os sistemas. Procure trip_ids duplicados: SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10. |
agg_hourly_zone_trips | A janela contínua de novo cálculo de 2 h precisa delimitar corretamente o intervalo de exclusão. Se for ampla demais, as agregações antigas serão excluídas; se for estreita demais, as agregações obsoletas permanecerão. | Confira algumas tuplas (hour_bucket, zone_id) com o Snowflake. Verifique se o total de trip_count em todas as zonas corresponde a AGG_HOURLY_ZONE_TRIPS do Snowflake no mesmo período. |
| Virada do produtor | A interrupção do script de migração durante a execução deixa uma lacuna na contagem de linhas; execute-o novamente com --resume para preenchê-la. Uma repetição do produtor após a virada pode reinserir corridas que já estão no ClickHouse. | scripts/01_verify_migration.sh — verifica a paridade da contagem de linhas entre o Snowflake e o ClickHouse. ReplacingMergeTree(_synced_at) trata as inserções duplicadas de maneira idempotente. |
Movimentação dos dados: script de migração em Python (scripts/02_migrate_trips.py)
Por que usar o script em Python em vez da intermediação por armazenamento de objetos ou
ClickPipes?
remoteSecure() serve para transferir dados entre instâncias do ClickHouse — não se aplica aqui.
- A intermediação por armazenamento de objetos (Snowflake → S3 → função de tabela S3 do ClickHouse) funcionaria, mas aumentaria a complexidade: exige provisionar um bucket S3, funções de IAM e COPY INTO no Snowflake — uma sobrecarga desnecessária para um laboratório.
- O ClickPipes não aceita o Snowflake como origem. Suas origens compatíveis são Kafka, S3, Kinesis, CDC do PostgreSQL e CDC do MySQL.
- O script em Python usa
snowflake-connector-python e clickhouse-connect — pacotes já instalados para o laboratório. Ele mostra o progresso em tempo real, permite retomadas com --resume em caso de interrupção e disponibiliza todo o código para inspeção.
Estratégia incremental (dbt): delete_insert
Por que escolher delete_insert em vez das estratégias append ou merge?
append insere novas linhas sem alterar as existentes. Para fact_trips, cujas linhas podem ser atualizadas, isso cria duplicatas. Incorreto.
merge (quando disponível) seria a opção mais próxima do MERGE INTO do Snowflake, mas a estratégia merge do dbt-clickhouse tem limitações com o ReplacingMergeTree e não é a abordagem recomendada.
delete_insert exclui as linhas no intervalo de chaves do lote recebido e depois insere todas as novas linhas. Isso é idempotente (repetir a execução produz o mesmo resultado), trata inserções e atualizações e funciona corretamente com o ReplacingMergeTree. É a recomendação padrão da comunidade dbt-clickhouse para padrões de upsert.
| Critério | Limite | Medido por |
|---|
| Paridade da contagem de linhas | Correspondência ≥ 99,9% (CH ≥ SF após a virada é esperado) | scripts/01_verify_migration.sh |
| Paridade do checksum | Correspondência de MD5 em uma amostra de 10 mil linhas | scripts/02_validate_parity.sql |
| Taxa de aprovação dos testes dbt | 100% | dbt test em dbt/nyc_taxi_dbt_ch |
| Paridade dos resultados das consultas | As 7 consultas retornam os mesmos resultados (dentro da tolerância de ponto flutuante) | Comparação manual na saída de scripts/run_benchmark.sh |
| Modelo | Materialização | Por quê |
|---|
stg_trips | view | Lê e limpa trips_raw; esse modelo não é atualizado; custo de armazenamento zero; sempre reflete o estado atual da origem |
stg_taxi_zones | view | Mesmo motivo — limpeza e repasse de uma tabela de origem |
int_trips_enriched | ephemeral | Lógica pura de junção usada somente por fact_trips; incorporá-la como uma CTE evita uma tabela física redundante; nenhum modelo a consulta diretamente |
fact_trips | incremental | As corridas podem ser corrigidas depois do fato; somente as linhas novas e atualizadas devem ser processadas a cada execução |
agg_hourly_zone_trips | incremental | O novo cálculo contínuo de 2 horas é um padrão incremental — processa as linhas recentes, não todas as 50 milhões |
dim_taxi_zones | table | 265 zonas estáticas; reconstrução completa em cada execução do dbt por meio de uma troca atômica de tabelas; sem atualizações parciais |
dim_payment_type | table | 6 tipos estáticos; mesmo raciocínio de dim_taxi_zones |
dim_vendor | table | 3 fornecedores; mesmo raciocínio |
| Modelo | ENGINE | Coluna de versão | Por quê |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | As corridas podem ser corrigidas; updated_at definido como now() em cada inserção significa que a versão mais recente prevalece na desduplicação em segundo plano do RMT; delete_insert é o principal caminho para a exatidão, e o RMT é a rede de proteção |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | O novo cálculo contínuo reinsere agregações para os mesmos pares (hour_bucket, zone_id); o RMT garante que as agregações obsoletas sejam removidas na mesclagem em segundo plano |
dim_taxi_zones | MergeTree() | — | A recarga completa pelo dbt significa uma troca atômica de tabelas (reconstrução completa) em cada execução; duplicatas não podem se acumular; nenhuma desduplicação necessária |
dim_payment_type | MergeTree() | — | Mesmo motivo de dim_taxi_zones |
dim_vendor | MergeTree() | — | Mesmo motivo de dim_taxi_zones |
| Modelo | unique_key | incremental_strategy | Filtro incremental | Por que esse filtro? |
|---|
fact_trips | trip_id | delete_insert | WHERE updated_at > (SELECT max(updated_at) FROM {{ this }}) | A marca-d'água superior em updated_at captura tanto as novas corridas quanto as corridas corrigidas (ajustes de tarifa reinserem o mesmo trip_id com o mesmo pickup_at, mas um updated_at mais recente); uma marca-d'água em pickup_at deixaria as correções passarem silenciosamente |
agg_hourly_zone_trips | [hour_bucket, zone_id] | delete_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | A janela contínua de 2 horas força uma nova agregação das horas no limite, para que as contagens de horas parciais sejam sempre corrigidas; uma marca-d'água superior em max(pickup_at) deixaria a hora no limite permanentemente subestimada |
| Modelo | FINAL na cláusula FROM? | Por quê |
|---|
stg_trips | Sim — FROM trips_raw FINAL | trips_raw é ReplacingMergeTree; ela pode ter linhas trip_id duplicadas devido a repetições do script de migração ou do produtor após a virada. stg_trips é o único ponto de imposição: desduplique aqui para que todos os modelos posteriores (int_trips_enriched, fact_trips, agg_hourly_zone_trips) recebam dados limpos |
int_trips_enriched | Não | Lê de stg_trips (uma view), não de uma tabela RMT; FINAL não se aplica a views |
fact_trips | Não (no corpo do modelo) | delete_insert mantém fact_trips limpa após cada execução concluída; adicionar FINAL dentro do modelo a aplicaria desnecessariamente à subconsulta is_incremental() que lê max(updated_at) de {{ this }}. Dashboards e testes do dbt usam FINAL externamente ao consultar diretamente fact_trips |
Este é o exemplo concluído. Seu migration-plan.md deve corresponder às principais
decisões apresentadas aqui — ou documentar explicitamente por que suas escolhas foram
diferentes.