Snowflake MigrationClickHouse Workshops

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.


Checklist de conclusão

  • Seleção de mecanismos: concluída
  • Projeto de chaves de ordenação: concluído
  • Tradução do esquema: concluída
  • Plano de ondas de migração: concluído
  • Projeto de modelos dbt: concluído

Seção 1: resumo do perfil

MétricaValor
Total de tabelas7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE)
Total de views2 (STG_TRIPS, STG_TAXI_ZONES)
Streams1 (TRIPS_CDC_STREAM em TRIPS_RAW)
Tasks2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK)
Total de linhas em TRIPS_RAW~50.000.000
Intervalo de datasJanela contínua de 4 anos com término no momento da preparação
Colunas VARIANT1 (TRIPS_RAW.TRIP_METADATA)
Usos de QUALIFY detectados1 (consulta Q3)
Usos de MERGE INTO detectados2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK)

Seção 2: inventário de objetos

ObjetoTipoEsquemaLinhasGrau de complexidadeObservações
trips_rawTabelaraw~50 miBRMT 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_tripsView do dbtstaging—BJSONExtract para TRIP_METADATA; exige testes dos caminhos JSON
stg_taxi_zonesView do dbtstaging—ARepasse; trivial
int_trips_enrichedEfêmero do dbtstaging—ACTE; diferenças de SQL tratadas nos modelos-pai
fact_tripsIncremental do dbtanalytics~50 miCMecanismo RMT; delete_insert; reescrita de QUALIFY; FINAL obrigatória
agg_hourly_zone_tripsIncremental do dbtanalytics~140 milBRMT; janela contínua de novo cálculo de 2 h; teste o limite da partição com atenção
dim_taxi_zonesTabela do dbtanalytics265AReferência estática; recarga completa; trivial
dim_payment_typeTabela do dbtanalytics6AReferência estática; trivial
dim_vendorTabela do dbtanalytics3AReferência estática; trivial
taxi_zones_dictDicionárioanalytics265BSintaxe específica do ClickHouse; dictGet() no momento da consulta
mv_hourly_revenueMV atualizávelanalytics—BSintaxe REFRESH EVERY; verificar a substituição atômica
TRIPS_CDC_STREAM / CDC_CONSUME_TASKStream + Task do Snowflake——DSem equivalente no ClickHouse; substituídos pela virada direta do produtor na Parte 3

Seção 3: decisões sobre a seleção de mecanismos

TabelaMecanismoColuna de versãoRaciocínio
trips_rawReplacingMergeTree(_synced_at)_synced_atO 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_tripsReplacingMergeTree(updated_at)updated_atAs 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_tripsReplacingMergeTree(updated_at)updated_atO 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_zonesMergeTree()—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_typeMergeTree()—Mesmo motivo — recarga completa.
dim_vendorMergeTree()—Mesmo motivo — recarga completa.
mv_hourly_revenueMergeTree()—A MV ATUALIZÁVEL substitui atomicamente todo o conjunto de resultados em cada REFRESH. Sem upserts.

Seção 4: projeto das chaves de ordenação

TabelaORDER BYRaciocí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.

Seção 5: observações sobre a tradução do esquema

ColunaTipo no SnowflakeTipo no ClickHouseJustificativa da decisão
TRIP_METADATAVARIANTStringPreserva 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_ATTIMESTAMP_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_IDINTEGERUInt16Valores 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_IDINTEGERUInt8Valores de 1 a 3. O máximo de UInt8 é 255 — correto. 1 byte por linha.
DRIVER_RATINGFLOATNullable(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_ATTIMESTAMP_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.

Traduções de funções necessárias

Expressão do SnowflakeEquivalente 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::FLOATJSONExtractFloat(trip_metadata, 'driver', 'rating')
TRIP_METADATA:surge_multiplier::FLOATJSONExtractFloat(trip_metadata, 'surge_multiplier')
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATEIncremental delete_insert do dbt — exclui as linhas com chaves correspondentes e depois insere todas as novas linhas

Seção 6: ondas de migração

OndaObjetosDependênciasObservações
Onda 0trips_raw (esquema), dim_taxi_zones, dim_payment_type, dim_vendorNenhumaO 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 1Carga 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 2stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_tripsOnda 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 3taxi_zones_dict, mv_live_trip_feedOnda 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 4Virada 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.

Registro de riscos (objetos de grau C/D)

ObjetoRiscoMétodo de verificação
fact_tripsConsultas 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_tripsA 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 produtorA 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.

Seção 7: lacunas de dialeto conhecidas

  • QUALIFY — afeta: Q3 (queries/q03_top_trips_qualify.sql)
  • Caminho com dois-pontos de VARIANT — afeta: Q4, Q5 (acesso ao JSON de TRIP_METADATA)
  • LATERAL FLATTEN — não é usado nesta carga de trabalho; VARIANT é acessado pelo caminho com dois-pontos, não por FLATTEN
  • MERGE INTO — afeta: modelos incrementais do dbt (fact_trips, agg_hourly_zone_trips)
  • Streams do Snowflake → virada do produtor (as gravações ativas vão diretamente ao ClickHouse após a virada)
  • Diferenças nas funções de data — afetam: Q1 (DATE_TRUNC), Q3 (DATEADD), Q4 (DATEDIFF)

Seção 8: estratégia de migração

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.

Seção 9: critérios para a virada

CritérioLimiteMedido por
Paridade da contagem de linhasCorrespondência ≥ 99,9% (CH ≥ SF após a virada é esperado)scripts/01_verify_migration.sh
Paridade do checksumCorrespondência de MD5 em uma amostra de 10 mil linhasscripts/02_validate_parity.sql
Taxa de aprovação dos testes dbt100%dbt test em dbt/nyc_taxi_dbt_ch
Paridade dos resultados das consultasAs 7 consultas retornam os mesmos resultados (dentro da tolerância de ponto flutuante)Comparação manual na saída de scripts/run_benchmark.sh


Seção 10: projeto dos modelos dbt

Seleção das materializações

ModeloMaterializaçãoPor quê
stg_tripsviewLê e limpa trips_raw; esse modelo não é atualizado; custo de armazenamento zero; sempre reflete o estado atual da origem
stg_taxi_zonesviewMesmo motivo — limpeza e repasse de uma tabela de origem
int_trips_enrichedephemeralLó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_tripsincrementalAs corridas podem ser corrigidas depois do fato; somente as linhas novas e atualizadas devem ser processadas a cada execução
agg_hourly_zone_tripsincrementalO 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_zonestable265 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_typetable6 tipos estáticos; mesmo raciocínio de dim_taxi_zones
dim_vendortable3 fornecedores; mesmo raciocínio

Configuração dos mecanismos

ModeloENGINEColuna de versãoPor quê
fact_tripsReplacingMergeTree(updated_at)updated_atAs 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_tripsReplacingMergeTree(updated_at)updated_atO 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_zonesMergeTree()—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_typeMergeTree()—Mesmo motivo de dim_taxi_zones
dim_vendorMergeTree()—Mesmo motivo de dim_taxi_zones

Estratégia incremental

Modelounique_keyincremental_strategyFiltro incrementalPor que esse filtro?
fact_tripstrip_iddelete_insertWHERE 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_insertWHERE pickup_at >= now() - INTERVAL 2 HOURA 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

Posicionamento de FINAL

ModeloFINAL na cláusula FROM?Por quê
stg_tripsSim — FROM trips_raw FINALtrips_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_enrichedNãoLê de stg_trips (uma view), não de uma tabela RMT; FINAL não se aplica a views
fact_tripsNã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.

Nesta página

PT