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.
Ponto de partida
Módulo 02 concluído: migration-plan.md está preenchido e todas as caixas da lista de verificação de
conclusão estão marcadas, enquanto o produtor do Snowflake continua em execução. O setup.sh deste
módulo procura esse arquivo e avisa se ele estiver ausente ou incompleto, mas nunca bloqueia. Nada aqui
impede que você continue sem o plano; apenas sua compreensão dos próximos dois módulos será prejudicada.
Reserve cerca de 60 minutos no total, dos quais aproximadamente 40 a 50 correspondem a uma transferência
autônoma de dados que pode ficar em segundo plano. É também aqui que começam os gastos do período de
avaliação do ClickHouse Cloud: provisionar o serviço e concluir este módulo consome cerca de US$ 1 a US$ 2
em créditos de avaliação (o laboratório inteiro custa aproximadamente US$ 2 a US$ 4).
Por quê
Neste módulo, o plano se torna realidade. Todas as decisões registradas no migration-plan.md durante
o módulo 02 — mecanismo MergeTree de cada tabela, chave ORDER BY derivada da carga real de consultas
e tradução das construções exclusivas do Snowflake — são inseridas diretamente no DDL das tabelas,
e não derivadas novamente do zero. O ClickHouse não tem um índice que possa ser acrescentado depois:
se uma chave ORDER BY se mostrar errada quando 50 milhões de linhas já estiverem na tabela, a correção
será uma recarga completa, não um rápido ALTER.
É também por isso que a barreira flexível importa, embora não possa impedi-lo. Se você executar este
módulo sem um plano completo, ainda terá sucesso mecânico: dbt run criará fact_trips como
ReplacingMergeTree, e o script transferirá 50 milhões de linhas. No entanto, você não saberá por que
esse mecanismo foi escolhido em vez de um MergeTree simples, por que a chave de ordenação tem esse
formato nem como justificar os ganhos de aproximadamente 6 a 9 vezes do benchmark que o módulo 04
mostrará depois. A tabela de alinhamento de decisões abaixo associa cada escolha implementada neste
módulo à pergunta da planilha que ela responde, para que você possa conferir seu plano antes de
provisionar qualquer coisa.
Conceitos — nos bastidores
Arquitetura de destino. O Snowflake continua recebendo novas corridas do produtor enquanto um script Python de execução única transfere 50 milhões de linhas existentes para o ClickHouse. Os dois sistemas funcionam em paralelo durante toda a migração; isso ainda não é uma virada.
No lado do ClickHouse, trips_raw é a tabela de chegada na qual o script grava. Em seguida, o dbt
constrói sobre ela as views de preparação e o restante da camada analítica. Este módulo cria o esquema,
mas ainda não o preenche além de trips_raw:
Legenda das cores do diagrama:
- Verde — fontes de dados (o produtor de corridas antes e depois da virada)
- Azul — tabelas do Snowflake
- Laranja — modelos e pipeline do dbt
- Vermelho — tabelas e views materializadas do ClickHouse
- Ciano — dashboards do Apache Superset
- Setas tracejadas — fluxos após a virada
Por que usar um script Python em vez de um conector nativo. Existem vários métodos para mover dados do Snowflake para o ClickHouse. Este laboratório usa um script Python em lotes. Veja o motivo em comparação com as alternativas:
| Método | Como funciona | Por que não é usado aqui |
|---|---|---|
| ClickPipes (origem Snowflake) | Conector nativo do ClickHouse Cloud — sem ETL, interface gerenciada | O Snowflake não é uma origem aceita pelo ClickPipes. O ClickPipes aceita Kafka, S3, Kinesis, CDC de PostgreSQL e MySQL e armazenamento de objetos. |
| Exportação para S3 → ClickPipes S3 | COPY INTO @stage exporta Parquet/CSV para S3; o conector S3 do ClickPipes carrega no ClickHouse | Exige bucket S3, função do IAM, stage do Snowflake e conta AWS. Acrescenta cerca de 3 etapas de configuração antes que qualquer dado seja movido. É viável em produção, mas infraestrutura demais para um laboratório. |
Exportação para S3 → clickhouse-client | A mesma exportação para S3, carregada com INSERT INTO ... SELECT FROM s3(...) | Mesmos pré-requisitos do S3. O parceiro também precisa gerenciar manualmente a divisão de arquivos e a retomada. |
| Snowflake → Kafka → ClickHouse | Um stream CDC do Snowflake alimenta um tópico Kafka; o conector Kafka do ClickPipes faz a ingestão | Pipeline completo de streaming, adequado a requisitos de latência inferior a um minuto em produção. Um cluster Kafka é pesado demais para um laboratório. |
| Script Python (este laboratório) | snowflake-connector-python lê lotes de 100 mil linhas do cursor; clickhouse-connect insere diretamente | Nenhuma infraestrutura adicional além dos pacotes já necessários. Retomável com --resume (marca d'água max(pickup_at)). Progresso em tempo real. Cerca de 40 a 50 min para 50 milhões de linhas a aproximadamente 20 mil linhas/s, aceitável para um exercício de migração única. |
Por que o script Python é a escolha certa para este laboratório:
- Dispensa uma conta AWS. Abordagens baseadas em S3 exigem criar um bucket, políticas do IAM e um stage externo do Snowflake: três etapas sem relação com o ClickHouse.
- É autocontido. Os dois pacotes (
snowflake-connector-python,clickhouse-connect) são instalados no mesmo ambiente virtual do dbt. Nenhum serviço nem credencial novos. - Pode ser retomado.
--resumetorna seguro interromper e reiniciar o script.ReplacingMergeTree(_synced_at)garante a deduplicação automática de inserções repetidas ao tentar novamente. - É transparente. Os parceiros podem ler o script, entender o mapeamento de colunas e adaptá-lo ao próprio esquema, o que ensina mais do que percorrer um assistente na interface.
Tratamento da lacuna da migração. O produtor do Snowflake continua ativo durante os cerca de 40 a 50 minutos de execução do script. Todas as corridas gravadas no Snowflake nesse período ficam ausentes do ClickHouse. O laboratório fecha essa lacuna com uma abordagem de duas passagens durante a virada, apresentada diretamente no módulo 05:
- Interrompa o produtor do Snowflake para congelar o conjunto de dados.
- Execute
python scripts/02_migrate_trips.py --resume: somente as linhas da diferença são transferidas (segundos, não minutos). - Inicie o produtor do ClickHouse.
A mesma deduplicação de ReplacingMergeTree(_synced_at) que lida com novas tentativas da migração também
resolve este caso: se alguma linha se sobrepuser entre a execução deste módulo e a passagem posterior com
--resume, o _synced_at posterior prevalecerá.
Quando escolher S3 em produção. Se o conjunto tiver mais de 500 milhões de linhas ou se o custo da consulta de um warehouse do Snowflake que analisa a tabela inteira for significativo, prefira exportar para S3: o Snowflake exporta Parquet compactado em paralelo (muito mais rápido que um único cursor), e o ClickHouse também pode carregar do S3 em paralelo. A abordagem com script Python funciona bem na escala do laboratório.
Alinhamento das decisões. A tabela abaixo apresenta a mesma lista de decisões das planilhas 1
(seleção de mecanismo), 2 (chaves de ordenação) e 3 (tradução de esquema) no migration-plan.md, conferida
com o que o laboratório realmente constrói. Compare-a com seu próprio plano antes de provisionar qualquer coisa:
| Decisão | Implementação neste laboratório | Por quê |
|---|---|---|
Mecanismo de trips_raw | ReplacingMergeTree(_synced_at) | O script de migração Python usa INSERTs em lotes que podem ser tentados novamente após uma interrupção. _synced_at DateTime DEFAULT now() é definido em cada INSERT; assim, uma linha repetida chega depois e tem um _synced_at maior, prevalecendo durante a deduplicação do RMT e tornando as tentativas idempotentes. Novas tentativas do produtor após a virada também são seguras pelo mesmo motivo. stg_trips consulta com FINAL para garantir uma linha por corrida. |
Mecanismo de fact_trips | ReplacingMergeTree(updated_at) | Corridas podem ser corrigidas (ajustes de tarifa); updated_at é a coluna de versão |
Mecanismo de agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | Recálculo em janela móvel = padrão de upsert |
Mecanismo das tabelas dim_* | MergeTree() | Recarga completa em cada execução do dbt; sem upserts |
ORDER BY de fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Todas as consultas de Q1 a Q7 filtram por pickup_at; trip_id garante exclusividade no último nível |
ORDER BY de agg_hourly_zone_trips | (hour_bucket, zone_id) | As duas colunas aparecem em todas as consultas de agregação |
| VARIANT → | String + JSONExtract* | Preserva o JSON bruto; a extração ocorre durante a consulta |
| QUALIFY → | Subconsulta que envolve ROW_NUMBER() | O ClickHouse tem uma cláusula QUALIFY nativa desde a v24.5, mas a forma com subconsulta é ensinada por ser portável para versões do ClickHouse e mecanismos SQL anteriores ou sem QUALIFY |
| MERGE INTO → | Incremental delete_insert no dbt | Estratégia idiomática de upsert do dbt-clickhouse; evita reescrever toda a tabela |
Etapa 1 — provisionar o cluster ClickHouse
setup.sh faz uma única coisa: executa terraform apply e grava os dados de conexão em
.clickhouse_state. Ele também verifica novamente a existência de migration-plan.md antes de
provisionar qualquer coisa — consulte a seção Por quê —, mas apenas avisa e nunca bloqueia.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
# Configure credentials
cp .env.example .env
vim .env
# Fill in: CLICKHOUSE_ORG_ID, CLICKHOUSE_TOKEN_KEY, CLICKHOUSE_TOKEN_SECRET, CLICKHOUSE_PASSWORD
# Provision
source .env && ./setup.sh.env é ignorado pelo Git; nunca faça commit dele.
Saída esperada: o Terraform cria 2 recursos (serviço + lista de acesso por IP) em aproximadamente 2 a 3 minutos:
Apply complete! Resources: 2 added, 0 changed, 0 destroyed.
Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443O host e a porta são salvos em .clickhouse_state. Carregue-o em qualquer terminal para obter a conexão:
source .clickhouse_stateVerificação:
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1Etapa 2 — criar as tabelas vazias
Primeiro, crie manualmente trips_raw com o mecanismo correto. O script de migração carrega dados
nessa tabela na etapa 3. Ela já precisa existir com ReplacingMergeTree para que a coluna de versão
esteja configurada antes da chegada de qualquer linha.
-- Run in the ClickHouse SQL console (cloud.clickhouse.com -> SQL console)
CREATE TABLE IF NOT EXISTS default.trips_raw (
trip_id String,
vendor_id UInt8,
pickup_at DateTime64(3, 'UTC'),
dropoff_at DateTime64(3, 'UTC'),
passenger_count UInt8,
trip_distance_miles Float32,
pickup_location_id UInt16,
dropoff_location_id UInt16,
payment_type_id UInt8,
rate_code_id UInt8,
store_fwd_flag String,
fare_amount_usd Float32,
extra_amount_usd Float32,
mta_tax_usd Float32,
tip_amount_usd Float32,
tolls_amount_usd Float32,
total_amount_usd Float32,
ingested_at DateTime64(3, 'UTC'),
trip_metadata String,
_synced_at DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id);_synced_at é definido automaticamente em cada INSERT. Se o script de migração for interrompido e
executado novamente com --resume, poderão existir brevemente linhas duplicadas para o mesmo trip_id.
O RMT mantém a linha posterior (com _synced_at maior). stg_trips consulta trips_raw FINAL para
forçar a deduplicação antes que qualquer modelo posterior veja os dados.
Em seguida, carregue os dados de referência das zonas. São dados estáticos (265 zonas do NYC TLC) lidos
como origem por stg_taxi_zones no dbt.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
clickhouse-client --host "${CLICKHOUSE_HOST}" --port 9440 --secure \
--user default --password "${CLICKHOUSE_PASSWORD}" \
--multiquery < scripts/00_seed_zones.sql(Como alternativa, cole o conteúdo de scripts/00_seed_zones.sql diretamente no console SQL do ClickHouse.)
Configure o perfil dbt. O dbt_project.yml deste projeto declara profile: 'nyc_taxi_ch'. Sem um perfil
correspondente em ~/.dbt/profiles.yml, dbt run falhará imediatamente com
Could not find profile named 'nyc_taxi_ch'. O modelo fica em
workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example.
O módulo 01 já gravou ~/.dbt/profiles.yml com um perfil nyc_taxi: para o Snowflake, e o ciclo de
atualização da etapa 4 continua consultando esse perfil enquanto o produtor do Snowflake estiver ativo.
Não substitua esse arquivo pelo modelo do ClickHouse: sobrescrevê-lo com profiles.yml.example apagaria
o perfil nyc_taxi: e interromperia o ciclo de atualização do módulo 01. Em vez disso, abra o modelo e
adicione o bloco nyc_taxi_ch: ao ~/.dbt/profiles.yml existente como um segundo perfil de primeiro nível,
ao lado de nyc_taxi::
nyc_taxi: # from module 01 — leave this one alone
target: dev
outputs:
dev:
type: snowflake
# ...
nyc_taxi_ch: # add this block
target: dev
outputs:
dev:
type: clickhouse
schema: nyc_taxi_ch
host: "{{ env_var('CLICKHOUSE_HOST') }}"
port: 8443
user: "{{ env_var('CLICKHOUSE_USER', 'default') }}"
password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
secure: truenyc_taxi_ch: lê CLICKHOUSE_HOST, CLICKHOUSE_USER e CLICKHOUSE_PASSWORD do ambiente por meio de
env_var(). Portanto, .env e .clickhouse_state precisam ser carregados antes de qualquer comando dbt
deste módulo; o dbt run abaixo já faz isso. Assim como o perfil do Snowflake, ~/.dbt/profiles.yml contém
credenciais e é ignorado pelo Git. Nunca faça commit dele. Essa combinação acrescenta um segundo conjunto
de credenciais a um arquivo que já continha o primeiro.
Verificação:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt debug
# Expected: "All checks passed!" — confirms dbt found the nyc_taxi_ch profile and
# connected to ClickHouseDepois, execute dbt run para criar as tabelas analíticas e as views de preparação:
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"
dbt deps # install packages (first run only)
dbt run # creates analytics tables and staging views; all empty at this pointResultado esperado: cerca de 8 modelos criados em menos de 2 minutos (todas as tabelas vazias).
Verificação:
# Check analytics tables were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+analytics" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: agg_hourly_zone_trips, dim_date, dim_payment_type, dim_vendor, dim_taxi_zones, fact_trips
# Check staging views were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+staging" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: stg_trips, stg_taxi_zones
# Check trips_raw exists with the correct engine
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+engine+FROM+system.tables+WHERE+database%3D%27default%27+AND+name%3D%27trips_raw%27" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: ReplacingMergeTreeAs seis tabelas analíticas e as duas views de preparação agora existem, mas todas continuam vazias:
dbt run apenas criou seus esquemas. A única tabela que terá dados após esta etapa é trips_raw, mas
ela também ainda está vazia; os dados virão a seguir.
Etapa 3 — migrar os dados
Carregue todas as linhas de NYC_TAXI_DB.RAW.TRIPS_RAW no Snowflake para default.trips_raw no ClickHouse
usando o script de migração em lotes escrito em Python.
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate
python scripts/02_migrate_trips.pySaída esperada (aproximadamente 40 a 50 minutos para 50 milhões de linhas):
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
NYC Taxi Migration: Snowflake -> ClickHouse
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
Rows to migrate: 50,000,000
Batch size: 100,000
Rows inserted Elapsed ETA Rate
-------------------- ------------ ---------------------- ---------------
100,000 0m 07s 56m 14s remaining 13,945 rows/s
200,000 0m 14s 55m 28s remaining 14,021 rows/s
...O produtor do Snowflake continua gravando em TRIPS_RAW durante toda a execução do script, portanto
o ClickHouse fica para trás por aproximadamente o tempo desta transferência. Essa lacuna é esperada
e será tratada no módulo 05, não aqui.
Se o script for interrompido, execute-o novamente com --resume para continuar a partir do último ponto:
python scripts/02_migrate_trips.py --resume--resume lê max(pickup_at) no ClickHouse e ignora as linhas já carregadas. Por isso, é sempre seguro
interromper e reiniciar o script; você nunca ficará com uma carga parcial impossível de recuperar.
Como verificar se você terminou
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
# Row count in trips_raw
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+count()+FROM+default.trips_raw" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: approximately 50000000
# .clickhouse_state was written by setup.sh
ls -la .clickhouse_state
# Service is reachable
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1Estado final
O serviço ClickHouse Cloud está ativo e acessível, e .clickhouse_state foi gravado no disco por
setup.sh 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 — sete objetos no total no
esquema analytics, criados pelo dbt run deste módulo. Todas as tabelas de analytics ainda estão vazias,
com exceção de mv_live_trip_feed, que já contém a linha de snapshot produzida quando dbt run criou a view. Fora isso, somente default.trips_raw contém dados: aproximadamente 50 milhões de
linhas transferidas pelo script Python na etapa 3.
O produtor do Snowflake continua em execução. Ele não foi interrompido neste módulo e não deve parar
aqui. Cada corrida gravada em TRIPS_RAW no Snowflake depois do último lote do script é uma linha ausente
do ClickHouse. Portanto, o ClickHouse agora está atrasado em aproximadamente a duração da janela de migração
(cerca de 40 a 50 minutos, além do tempo de configuração deste módulo). A lacuna é real e continua crescendo
enquanto o produtor estiver ativo. Não a feche neste módulo. A virada do módulo 05 a fecha deliberadamente,
em uma etapa controlada de duas passagens que mede seu tamanho antes de eliminá-la. Interromper o produtor ou
executar o script novamente agora removeria exatamente o fenômeno que o módulo 05 foi criado para demonstrar.
Planilha 5: projeto dos modelos dbt
Configure materialização, mecanismo e estratégia incremental para cada modelo dbt, com feedback imediato para cada resposta.
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.