Snowflake MigrationClickHouse Workshops

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.

Fluxo de dados da migração: o produtor de corridas do Snowflake continua gravando em TRIPS_RAW enquanto um script Python de execução única transfere 50 milhões de linhas em lotes de 100 mil para o ClickHouse Cloud

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:

Lado de destino no ClickHouse: trips_raw em ReplacingMergeTree alimenta views de preparação criadas pelo dbt, tabelas de fatos e dimensões, um agregado por hora e um dicionário de zonas

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étodoComo funcionaPor que não é usado aqui
ClickPipes (origem Snowflake)Conector nativo do ClickHouse Cloud — sem ETL, interface gerenciadaO 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 S3COPY INTO @stage exporta Parquet/CSV para S3; o conector S3 do ClickPipes carrega no ClickHouseExige 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-clientA 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 → ClickHouseUm stream CDC do Snowflake alimenta um tópico Kafka; o conector Kafka do ClickPipes faz a ingestãoPipeline 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 diretamenteNenhuma 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. --resume torna 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:

  1. Interrompa o produtor do Snowflake para congelar o conjunto de dados.
  2. Execute python scripts/02_migrate_trips.py --resume: somente as linhas da diferença são transferidas (segundos, não minutos).
  3. 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ãoImplementação neste laboratórioPor quê
Mecanismo de trips_rawReplacingMergeTree(_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_tripsReplacingMergeTree(updated_at)Corridas podem ser corrigidas (ajustes de tarifa); updated_at é a coluna de versão
Mecanismo de agg_hourly_zone_tripsReplacingMergeTree(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 dbtEstraté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 = 8443

O host e a porta são salvos em .clickhouse_state. Carregue-o em qualquer terminal para obter a conexão:

source .clickhouse_state

Verificação:

curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1

Etapa 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: true

nyc_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 ClickHouse

Depois, 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 point

Resultado 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: ReplacingMergeTree

As 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.py

Saí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: 1

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

Nesta página

Acompanhar seu progresso?

Opcional. Enviaremos um link por e-mail para confirmar seu endereço; o progresso será registrado depois que você o abrir.

Use seu e-mail corporativo, não um endereço pessoal.

O acompanhamento do progresso também exige a aceitação dos Termos de Serviço atuais nas Configurações de privacidade.

PT