Snowflake MigrationClickHouse Workshops

ClickHouse 运维

运行目标服务:查询、监控 merge 与 part、字典,以及与 Snowflake 不同的运维习惯。

本文讲解你在本次迁移实验课中会遇到的 ClickHouse 概念,它们是什么、为什么存在,以及它们与你在第 1 部分用过的 Snowflake 构造有何不同。


1. 表引擎

ClickHouse 不是单引擎数据库。你创建的每一张表都必须声明它的引擎,引擎决定了数据如何存储在磁盘上、重复数据如何处理,以及可用的能力有哪些。选错引擎是 ClickHouse 表结构设计中最常见的错误。

MergeTree

几乎所有生产表的基础引擎。

CREATE TABLE analytics.dim_taxi_zones (
    zone_id      UInt16,
    borough      String,
    service_zone String
) ENGINE = MergeTree()
ORDER BY zone_id;

它做什么: ClickHouse 把数据存放在 part 中,磁盘上已排序、已压缩的数据块。插入数据时会写出新的 part。在后台,ClickHouse 持续把较小的 part **合并(merge)**成较大的,并保持数据按 ORDER BY 键有序。引擎的名字就来自这里。

何时使用: 任何不需要去重、且插入为追加式或批量加载的表(维度表、原始事件表、日志表)。

关键性质: 不存在主键约束。ORDER BY 取值完全相同的两行都会被存下来。如果你需要去重,请使用 ReplacingMergeTree。

ReplacingMergeTree(version_col)

去重引擎。它在 MergeTree 之上扩展了一条规则:在后台合并期间,如果两行的 ORDER BY 键相同,只保留 version_col 取值最大的那一行。

CREATE TABLE analytics.fact_trips (
    trip_id     String,
    pickup_at   DateTime,
    total_amount Float64,
    updated_at  DateTime
) ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (pickup_at, trip_id);

重要,最终一致性: 去重只在后台合并期间发生。在任意给定时刻,你的表中都可能含有重复行。这被称为最终一致性。要在查询时得到完全去重的结果,请在 SELECT 中加上 FINAL:

-- Without FINAL: may return duplicates if merges haven't run
SELECT * FROM analytics.fact_trips WHERE trip_id = 'abc';

-- With FINAL: forces deduplication at query time (slower, always correct)
SELECT * FROM analytics.fact_trips FINAL WHERE trip_id = 'abc';

dbt 如何使用它: dbt-clickhouse 适配器以 delete_insert 增量策略作为主要机制,它显式删除键匹配的行并插入新行,这总是正确的。ReplacingMergeTree 则充当安全网,清理任何漏过去的重复数据(例如来自一次失败的部分插入)。

与 Snowflake 的对应关系: 没有直接对应物。在 Snowflake 中你使用 MERGE INTO ... WHEN MATCHED THEN UPDATE。ClickHouse 没有 MERGE 语句,ReplacingMergeTree 加上 FINAL 可以达到同样的逻辑结果。

可刷新 materialized view

ClickHouse 支持两种 materialized view。

基于触发的 MV(传统方式):在每次 INSERT 时执行,只处理刚插入的那一批数据。

-- Trigger-based: only sees the rows inserted in the current batch
CREATE MATERIALIZED VIEW analytics.mv_realtime_counts
TO analytics.counts_table AS
SELECT pickup_date, count() AS trips
FROM default.trips_raw
GROUP BY pickup_date;

可刷新 MV(按计划):按计划重新运行完整查询,就像一个 cron 作业。

-- Refreshable: runs the full SELECT every 3 minutes
CREATE MATERIALIZED VIEW analytics.mv_hourly_revenue
REFRESH EVERY 180 SECOND AS
SELECT
    toStartOfHour(pickup_at) AS hour_bucket,
    pickup_borough,
    sum(total_amount)        AS revenue
FROM analytics.fact_trips FINAL
GROUP BY hour_bucket, pickup_borough;

何时用哪一种:

  • 基于触发:对插入流做实时聚合,且你只需要处理新数据
  • 可刷新:需要查询 fact_trips FINAL 的聚合(必须看到整张表才能去重),或者能容忍几分钟数据陈旧、以换取更简单逻辑的仪表板

修改刷新间隔:

ALTER TABLE analytics.mv_hourly_revenue MODIFY REFRESH EVERY 60 SECOND;

2. 排序键(ORDER BY)

在 Snowflake 中,你用 CLUSTER BY 作为给优化器的提示。在 ClickHouse 中,ORDER BY 就是主索引,它决定数据在磁盘上的物理排序顺序,并驱动所有范围扫描。

它如何工作

ClickHouse 维护一个稀疏主索引:每约 8,192 行(一个数据 granule)一条索引项。当你按 ORDER BY 中的列过滤时,ClickHouse 可以整块跳过 granule 而完全不去读取它们。这就是 ClickHouse 每秒能扫描数十亿行的原因,大部分数据根本没离开磁盘。

基数顺序很重要

始终把低基数列放在前面,高基数列放在最后。这样索引在常见场景下才具备最强的跳过能力。

-- Good: low cardinality (borough, ~6 values) first, then high cardinality (trip_id)
ORDER BY (pickup_borough, toStartOfMonth(pickup_at), trip_id)

-- Bad: high cardinality first — the index can't skip anything useful
ORDER BY (trip_id, pickup_borough, pickup_at)

Snowflake CLUSTER BY 与 ClickHouse ORDER BY

特性Snowflake CLUSTER BYClickHouse ORDER BY
用途查询性能提示物理排序顺序(必需)
强制程度后台重新聚簇(异步)插入时始终强制执行
作用范围微分区数据 granule(约 8K 行)
是否必填否是,每张 MergeTree 表都必须有一个

示例:匹配已有的 Snowflake 聚簇键

-- Snowflake
CLUSTER BY (DATE_TRUNC('month', PICKUP_AT), PICKUP_LOCATION_ID)

-- ClickHouse equivalent
ORDER BY (toStartOfMonth(pickup_at), pickup_location_id, trip_id)
-- Note: trip_id added as tiebreaker to ensure unique sort order

跳数索引(简要说明)

对于不在 ORDER BY 键中的列,ClickHouse 支持跳数索引(bloom filter、minmax、set),它们为每个 granule 存储列级元数据。当需要过滤那些在排序键中排在高基数列之后的低基数列时,这很有用。

-- Add a bloom filter skip index on payment_type
ALTER TABLE analytics.fact_trips
ADD INDEX idx_payment_type payment_type TYPE bloom_filter GRANULARITY 4;

3. JSON 处理

Snowflake 的 VARIANT 列类型支持用冒号路径记法遍历嵌套 JSON。ClickHouse 则改用显式的 JSONExtract* 函数。

并排对照表

SnowflakeClickHouse说明
col:key::FLOATJSONExtractFloat(col, 'key')顶层浮点字段
col:driver.rating::FLOATJSONExtractFloat(col, 'driver', 'rating')嵌套浮点字段
col:app.surge_multiplier::FLOATJSONExtractFloat(col, 'app', 'surge_multiplier')嵌套浮点
col:route.waypoints[0]::STRINGJSONExtractString(col, 'route', 'waypoints', 0)按下标取数组元素
col:driver.id::INTJSONExtractInt(col, 'driver', 'id')整数字段

函数变体

-- Float (returns 0.0 if key missing or wrong type)
JSONExtractFloat(trip_metadata, 'driver', 'rating')

-- String (returns '' if missing)
JSONExtractString(trip_metadata, 'app', 'version')

-- Integer (returns 0 if missing)
JSONExtractInt(trip_metadata, 'driver', 'id')

-- Bool (returns 0/1)
JSONExtractBool(trip_metadata, 'app', 'is_shared')

-- Raw value as string (preserves JSON sub-object)
JSONExtractRaw(trip_metadata, 'route')

性能提示

如果你反复查询同一个 JSON 列,考虑在 staging 模型层(在 stg_trips.sql 中)把字段提取为带类型的列,而不是在每个下游查询里都调用 JSONExtractFloat。本实验课的 dbt 模型就是这么做的。


4. 日期/时间函数

Snowflake 和 ClickHouse 的日期/时间能力相似,但语法不同。最常见的对应转换:

并排对照表

SnowflakeClickHouse说明
DATE_TRUNC('hour', col)toStartOfHour(col)截断到小时
DATE_TRUNC('day', col)toStartOfDay(col) 或 toDate(col)截断到天
DATE_TRUNC('month', col)toStartOfMonth(col)截断到月
CURRENT_TIMESTAMP()now()当前日期时间
CURRENT_DATE()today()当前日期
DATEADD('day', -7, CURRENT_DATE())today() - INTERVAL 7 DAY日期运算
DATEDIFF('day', a, b)dateDiff('day', a, b)两个日期之间的天数

ClickHouse 独有的便捷函数

yesterday()              -- today() - 1 day
toStartOfWeek(col)       -- Monday of the containing week
toStartOfQuarter(col)    -- first day of the quarter
toYear(col)              -- extract year as integer
toMonth(col)             -- extract month as integer (1-12)
toDayOfWeek(col)         -- 1=Monday, 7=Sunday

INTERVAL 语法

-- ClickHouse
now() - INTERVAL 7 DAY
now() - INTERVAL 1 HOUR
now() - INTERVAL 30 MINUTE
pickup_at + INTERVAL 90 SECOND

-- Snowflake equivalent
DATEADD('day', -7, CURRENT_TIMESTAMP())
DATEADD('hour', -1, CURRENT_TIMESTAMP())

5. 近似函数

ClickHouse 是为分析型工作负载而生的,在这类场景中,对数十亿行求精确答案比给出一个对仪表板来说足够准确的近似答案要慢得多。ClickHouse 内置了若干近似聚合函数。

去重计数

函数精确度速度何时使用
uniqExact(col)精确最慢合规报表、开票
uniq(col)约 2% 误差快仪表板、探索分析
uniqHLL12(col)约 1.6% 误差最快,固定 2.5KB 内存高基数、内存受限场景
-- Exact (like Snowflake COUNT(DISTINCT ...))
SELECT uniqExact(trip_id) FROM analytics.fact_trips FINAL;

-- Approximate — good for "how many unique passengers today?"
SELECT uniq(passenger_id) FROM analytics.fact_trips FINAL;

分位数

函数说明
quantile(level)(col)精确分位数,内存开销大
quantileTDigest(level)(col)使用 t-digest 的近似算法,内存固定
quantileTDigestWeighted(level)(col, weight)加权 t-digest
-- P95 trip duration — approximate but uses O(1) memory
SELECT quantileTDigest(0.95)(duration_minutes)
FROM analytics.fact_trips FINAL;

-- Multiple percentiles in one pass
SELECT quantileTDigestMerge(0.5)(state), quantileTDigestMerge(0.95)(state)
FROM analytics.fact_trips FINAL;

经验法则: 交互式仪表板使用 uniq 和 quantileTDigest。只有在计费、SLA 或合规场景需要精确值时才使用 uniqExact 和 quantile。


6. 字典

字典是 ClickHouse 常驻内存并在查询时预先关联好的内存查找表。它们相当于 ClickHouse 中的小型维度表,你希望关联它,但不想承担完整 JOIN 的代价。

它们是什么

字典由某个数据源支撑(一张 ClickHouse 表、一个文件,或一个外部数据库),并在服务启动时或你执行 SYSTEM RELOAD DICTIONARIES 时被加载进内存。查找通过键进行,返回一个或多个属性。

CREATE DICTIONARY 语法

-- From 04_create_dictionary.sql
CREATE DICTIONARY analytics.taxi_zones_dict (
    zone_id      UInt16,
    borough      String,
    service_zone String
)
PRIMARY KEY zone_id
SOURCE(CLICKHOUSE(
    TABLE 'dim_taxi_zones'
    DB    'analytics'
))
LIFETIME(MIN 300 MAX 600)   -- refresh every 5-10 minutes
LAYOUT(FLAT());             -- hash map, best for < 1M rows

LAYOUT 选项:

  • FLAT():以整数键为下标的数组,最快,要求键是连续整数
  • HASHED():哈希表,适用于任意整数键
  • COMPLEX_KEY_HASHED():哈希表,支持复合键或字符串键

dictGet 用法

-- Instead of: JOIN analytics.dim_taxi_zones USING (zone_id)
SELECT
    trip_id,
    dictGet('analytics.taxi_zones_dict', 'borough', toUInt64(pickup_location_id)) AS pickup_borough,
    dictGet('analytics.taxi_zones_dict', 'borough', toUInt64(dropoff_location_id)) AS dropoff_borough
FROM analytics.fact_trips FINAL;

何时用字典、何时用 JOIN

场景使用
小而稳定的参照表(< 100 万行,很少变动)字典
大维度表或频繁更新的数据JOIN
反复运行且使用相同查找的仪表板查询字典(首次加载后查找几乎零成本)
一次性的分析查询JOIN

7. SAMPLE 子句

ClickHouse 支持在查询语法中直接进行行级采样。采样会读取数据中一个确定性的比例,当你不需要精确结果时,这对探索性分析很有用。

语法

-- Read approximately 10% of rows
SELECT count(), avg(total_amount)
FROM analytics.fact_trips SAMPLE 0.1;

-- Read a specific number of rows (approximately)
SELECT trip_id, pickup_at, total_amount
FROM analytics.fact_trips SAMPLE 1000000;

结果缩放

采样时,把聚合结果乘以 1 / sample_rate 来估算全表数值:

-- Estimate total revenue from 10% sample
SELECT sum(total_amount) * 10 AS estimated_total_revenue
FROM analytics.fact_trips SAMPLE 0.1;

何时使用 SAMPLE

  • 在整表运行之前做探索性分析(“我的查询逻辑对不对?”)
  • 可以接受近似值的仪表板卡片
  • 在有代表性的子集上训练 ML 模型

注意: SAMPLE 要求表的 ORDER BY 键以采样列开头,否则你必须在 CREATE TABLE 语句中加上 SAMPLE BY 子句。本实验课的 trips_raw 表正是为此而以 SAMPLE BY cityHash64(trip_id) 创建的。


8. dbt-clickhouse 适配器说明

dbt-clickhouse 适配器(dbt-clickhouse>=1.8)支持大多数标准 dbt 特性,但有一些 ClickHouse 特有的行为你需要了解。

delete_insert 增量策略

ClickHouse 没有 MERGE INTO。dbt-clickhouse 适配器的 delete_insert 策略对它进行模拟:

  1. 从目标表中删除键列与来批数据匹配的行
  2. 插入全部来批行
-- What dbt generates for incremental models
DELETE FROM analytics.fact_trips WHERE trip_id IN (SELECT trip_id FROM __dbt_tmp);
INSERT INTO analytics.fact_trips SELECT * FROM __dbt_tmp;

在你的模型中这样配置:

{{
    config(
        materialized='incremental',
        incremental_strategy='delete_insert',
        unique_key='trip_id',
        engine='ReplacingMergeTree(updated_at)',
        order_by='(pickup_at, trip_id)'
    )
}}

engine 与 order_by 配置

每张 MergeTree 表都需要一个引擎和一个 ORDER BY。在 dbt 模型配置中同时指定两者:

{{
    config(
        engine='MergeTree()',
        order_by='(zone_id)'
    )
}}

用 order_by 而不是 cluster_by

在 Snowflake 的 dbt 模型中你可能用过 cluster_by。在 dbt-clickhouse 中请改用 order_by。ClickHouse 中没有与 Snowflake cluster_by 等价的东西,ORDER BY 始终就是物理排序。

面向 ClickHouse Cloud 的 profiles.yml

ClickHouse Cloud 要求使用 TLS。请设置 secure: true:

# ~/.dbt/profiles.yml
nyc_taxi_ch:
  target: dev
  outputs:
    dev:
      type: clickhouse
      host: "{{ env_var('CLICKHOUSE_HOST') }}"
      port: 8443
      user: default
      password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
      schema: analytics       # default database for models without a custom schema
      secure: true
      threads: 4

schema 命名与独立数据库

Snowflake 中称为 schema 的东西,ClickHouse 中称为 database。dbt-clickhouse 适配器把 dbt 的 schema 映射为 ClickHouse 的 database。本项目中的 generate_schema_name 宏覆盖了 dbt 的默认行为,使带有 +schema: analytics 的模型落到 analytics 数据库,而不是 staging_analytics。

-- macros/generate_schema_name.sql
{% macro generate_schema_name(custom_schema_name, node) -%}
  {%- if custom_schema_name is none -%}
    {{ target.schema }}
  {%- else -%}
    {{ custom_schema_name }}
  {%- endif -%}
{%- endmacro %}

这与 Snowflake dbt 项目(第 1 部分)中使用的模式完全相同,该宏被刻意保持一致,以便两个适配器下的 schema 命名行为一致。

本页内容

ZH