首页
/ Ghost 的 Tinybird 实战指南:Datafile 格式、SQL 规则与性能优化最佳实践

Ghost 的 Tinybird 实战指南:Datafile 格式、SQL 规则与性能优化最佳实践

2026-09-07 17:36:28作者:邓越浪Henry

本篇技术指南以 Ghost 仓库中的 Tinybird 技能规范(.agents/skills/tinybird/ 及其 13 个规则文件)为主体,系统讲解 Tinybird 项目的目录结构、.datasource / .pipe / .connection 三类 datafile 的书写格式、SELECT-only 的 SQL 与参数模板规则、物化视图与去重的优化模式,以及 tb build / tb deploy 的构建部署流程。读完本文后,你将能够按 Ghost 的规范创建和重构 Tinybird 数据文件,并能对照仓库中 Ghost 流量分析(Traffic Analytics)的真实 Tinybird 工程验证每一条规则。

为什么 Tinybird 是 Ghost 流量分析的核心

Ghost 的流量分析系统采用双数据源架构:MySQL 存储内容、会员与订阅数据,Tinybird 负责处理页面浏览与访问会话的实时分析。两者通过 UUID 关联——例如 posts 表和 members 表中的 uuid 字段用于与 Tinybird 事件做关联(见 架构说明)。

Ghost 仓库中的 Tinybird 工程位于 tinybird 数据目录,目录组织正好体现了下文将讲解的分层结构:

ghost/core/core/server/data/tinybird/
├── datasources/     # 落地数据源(analytics_events 等 5 个 .datasource)
├── pipes/           # 物化管道:mv_hits、mv_session_data、filtered_sessions 等
├── endpoints/       # 30 个 API 端点:api_kpis、api_top_pages、api_active_visitors...
├── fixtures/       # 测试用样本数据(analytics_events.ndjson)
├── tests/           # 与端点同名的 .yaml 测试文件
└── scripts/

这套工程与技能规范中推荐的 datasources/endpoints/materializations/fixtures/tests/ 目录布局一致(Ghost 额外用 pipes/ 承载中间物化层)。

项目结构与构建部署规则

项目根与默认目录约定

规范(project-files 规则)要求:

  • 默认在项目根目录创建 tinybird/ 文件夹,Tinybird 资源目录嵌套其下;
  • .tinyb 凭证文件必须与执行 CLI 命令的目录同级,CLI 依赖它定位上下文;
  • tinybird.config.json 是构建/部署行为的唯一事实来源(source of truth)。

各类型资源的默认存放位置:

资源类型 目录
Endpoints /endpoints
Materialized pipes /materializations
Sink pipes /sinks
Copy pipes /copies
Connections /connections
Datasources /datasources
Fixtures /fixtures
Tests /tests

当项目规模增大时,规范建议按“域”或“消费方”拆分端点目录,例如 endpoints_dashboard/endpoints_public/,并用 tinybird.config.json 中的 include 字段控制哪些目录参与构建。

tb info:排障 CLI 上下文的第一选择

遇到凭证问题时,先运行 tb info 确认 CLI 上下文。它同时报告 Local 与 Cloud 两套环境的信息:.tinyb 文件的加载位置、当前登录的 Workspace、API URL、UI URL 以及 ClickHouse HTTP 接口地址。

构建与部署的目标环境

build-deploy 规则 定义了完整的目标环境矩阵:

新项目用 tb init 初始化,随后以 tinybird.config.json 为准:

{
  "dev_mode": "branch",
  "include": ["tinybird"]
}

构建/部署流程

  1. tinybird.config.json 读取 dev_mode
  2. 对配置的开发目标执行 tb build
  3. 只有在明确要求部署到云端生产时,才执行 tb deploy

tb build 的目标由 dev_mode 决定:

  • dev_mode: "local"tb build 指向 Tinybird Local;
  • dev_mode: "branch"tb build 指向 Tinybird Cloud 的某个 branch。

tb deploy(等价于 tb --cloud deploy)指向 Tinybird Cloud 生产环境:它会创建 staging 部署、迁移数据并晋级到 live。规则特别强调:不要把 tb build 当作生产部署。CI 中推荐用 tb --cloud deploy --check 做不落地应用的部署校验;需要显式确认时,先 tb --cloud deployment create --wait,再 tb --cloud deployment promote

非构建命令的环境指定

tb sqltb logs 等命令默认指向 local,需要显式覆盖才能指向其他环境:

tb sql "SELECT 1"                          # 默认 local
tb sql --cloud "SELECT 1"                  # 指向 Cloud
tb sql --branch=feature_metrics "SELECT 1"  # 指向指定 branch

tb logs
tb logs --cloud
tb logs --branch=feature_metrics

数据分层架构

对复杂项目,规范给出了六层数据分层模型:

用途 示例
Landing 外部源原始落地数据 raw_eventss3_import_logs
Cleaned 去重或转换后的数据 events_dedupnormalized_logs
Dimensions 查找/参照表 dim_organizationsdim_users
Aggregation 预计算指标的物化视图 mv_events_dailymv_usage_hourly
API 提供最终查询的端点管道 kpistop_pagesuser_activity
Export 向外部系统发送数据的 Sink 管道 sink_to_s3sink_to_kafka

并非每个项目都需要全部层数,从简单开始、随复杂度增长再加分层。Ghost 的流量分析工程正是这一分层的实例:Landing 层是 analytics_events 数据源,Aggregation 层是 mv_hitsmv_session_data 等物化管道,API 层是 30 个 api_* 端点。

规范还要求文档与描述中使用一致的术语大小写:Data Source(不是 datasource/data source)、PipeEndpointAPI EndpointMaterialized ViewTokenWorkspaceSinkCopy PipeConnection;datafile 中的指令名用大写,如 FORWARD_QUERYENGINE_SORTING_KEYENGINE_PARTITION_KEYCOPY_SCHEDULECOPY_MODETYPESCHEMADESCRIPTION

Datasource 文件格式:Ghost 的落地层怎么写

datasource-files 规则 的核心要求:

  • 内容不能为空;数据源名必须唯一;
  • 属性名(DESCRIPTIONSCHEMAENGINE 等)不缩进;
  • 默认使用 MergeTree 引擎;物化目标(materialized target)使用 AggregatingMergeTree
  • SCHEMA 一律使用 JSON 路径映射,如 `user_id` String `json:$.user_id`
  • 数组列写法:`items` Array(String) `json:$.items[:]`
  • DateTime64 必须带精度,如 DateTime64(3)
  • ENGINE_PARTITION_KEYENGINE_PRIMARY_KEY 仅在明确要求时才写;
  • 导入配置:S3/GCS 设置 IMPORT_CONNECTION_NAMEIMPORT_BUCKET_URIIMPORT_SCHEDULE(GCS 仅支持 @on-demand,S3 支持 @auto);Kafka 设置 KAFKA_CONNECTION_NAMEKAFKA_TOPICKAFKA_GROUP_ID
  • 从没有指定 schema 的 .ndjson 文件创建的落地数据源,使用 SCHEMA > 加单一列 `data` String `json:$`

标准 datafile 形态:

DESCRIPTION >
    Some meaningful description of the datasource

SCHEMA >
    `column_name_1` Type `json:$.column_name_1`,
    `column_name_2` Type `json:$.column_name_2`

ENGINE "MergeTree"
ENGINE_PARTITION_KEY "partition_key"
ENGINE_SORTING_KEY "sorting_key_1, sorting_key_2"

Ghost 实例:analytics_events 数据源

analytics_events.datasource 是该规范的完整落地:

TOKEN "tracker" APPEND
TOKEN "analytics-service" APPEND

SCHEMA >
    `timestamp` DateTime `json:$.timestamp`,
    `session_id` String `json:$.session_id`,
    `action` LowCardinality(String) `json:$.action`,
    `version` LowCardinality(String) `json:$.version`,
    `payload` String `json:$.payload`,
    `site_uuid` LowCardinality(String) `json:$.payload.site_uuid`,
    `inserted_at` DateTime64(3) DEFAULT now64() `json:$.inserted_at`

ENGINE "MergeTree"
ENGINE_PARTITION_KEY "toYYYYMM(timestamp)"
ENGINE_SORTING_KEY "site_uuid, timestamp"

几个值得注意的点:actionversionsite_uuid 使用 LowCardinality(String)——这正对应后文优化规则中“低基数字符串用 LowCardinality”的结构化要求;inserted_atDateTime64(3) 满足“DateTime64 必须带精度”的规则;排序键以 site_uuid 打头,符合“多租户场景下避免时间戳作为首键”的建议。

FORWARD_QUERY:云端 schema 演进

当 schema 变更与已部署的 Cloud 数据源不兼容时,规则要求在 datafile 中加 FORWARD_QUERY 把存量数据在读取时转换到新 schema——它只是 SELECT 列表(无 FROM/WHERE),在下一次部署压缩(compact)存量数据前一直生效。适用场景:为存量行补默认值的新列、列类型变更(如 String→UUID、Int32→Int64)、列重命名、删列(从 SELECT 中省略即可)。

# 新增带默认值的列
FORWARD_QUERY >
    SELECT *, 'unknown' as source

# 列类型变更
FORWARD_QUERY >
    SELECT timestamp, accurateCastOrDefault(session_id, 'UUID') as session_id, action, version, payload

# 列重命名
FORWARD_QUERY >
    SELECT old_name as new_name, other_column

Ghost 的 analytics_events.datasource 本身就带有一个 FORWARD_QUERY

FORWARD_QUERY >
    SELECT timestamp, session_id, action, version, toString(payload) as payload, site_uuid, defaultValueOfTypeName('DateTime64(3)') AS inserted_at

即把 payload 统一为 String 并为新增的 inserted_at 列补默认值——典型的“新增列 + 类型收敛”迁移。规则同时提醒:FORWARD_QUERY 完成使命后应在后续部署中移除,长期保留过期的 FORWARD_QUERY 只会增加复杂度。

TTL 与分区键对齐

当数据源同时设置 ENGINE_TTLENGINE_PARTITION_KEY 时,分区粒度必须等于或细于 TTL 窗口(如按天 TTL 配按天分区),这样分区才能作为整体过期删除,而不是被 ClickHouse 反复重写。分区粗于 TTL 窗口(如年分区配 65 天 TTL)的分区永远不会整体过期,只能持续重写。规则建议尽量在 ENGINE_SETTINGS 中设置 ttl_only_drop_parts=1,让 ClickHouse 只整块删除过期 part:

# 反例:年分区 + 65 天 TTL——分区永远不会整体过期
ENGINE_PARTITION_KEY "toYYYY(timestamp)"
ENGINE_TTL "toDateTime(timestamp) + toIntervalDay(65)"

# 正例:日分区与 TTL 窗口对齐——旧分区整体删除
ENGINE_PARTITION_KEY "toDate(timestamp)"
ENGINE_TTL "toDate(timestamp) + toIntervalDay(65)"
ENGINE_SETTINGS "ttl_only_drop_parts=1"

共享数据源

SHARED_WITH > 列出目标 workspace 即可共享数据源,但有三个限制:共享的数据源是只读的;不能把已共享的数据源再共享出去;不能基于共享数据源建物化视图。

Pipe 文件格式:四种 TYPE 的通用规则

pipe-files 规则 是所有 .pipe 文件的通用约束:

  • Pipe 名唯一;节点名不能与 pipe 名或其他资源名相同;
  • 属性名(DESCRIPTIONNODESQLTYPE 等)不缩进;
  • TYPE 的合法取值:endpointcopymaterializedsink
  • 输出节点要么写在 TYPE 段,要么写在最后一个节点。

最小形态:

DESCRIPTION >
    Some meaningful description of the pipe

NODE node_1
SQL >
    SELECT ...
TYPE endpoint

Ghost 实例:物化管道与端点

mv_session_data_v2.pipe 展示了物化管道的标准写法(文件头部注释说明:_v2 系列是性能实验产物,当前未启用,后续版本会移除):

NODE mv_session_data_v2_0
SQL >
    SELECT
        site_uuid,
        session_id,
        countState() as pageviews,
        minState(timestamp) as first_pageview,
        maxState(timestamp) as last_pageview,
        argMinState(source, timestamp) as source,
        argMinState(device, timestamp) as device,
        argMinState(utm_source, timestamp) as utm_source,
        ...
    FROM _mv_hits
    GROUP BY site_uuid, session_id

TYPE materialized
DATASOURCE _mv_session_data_v2

注意这里的 countState / argMinState 等 State 聚合——这正是后文物化视图规则的要点:管道中用 State 修饰符,目标数据源用 AggregateFunction 列。

端点文件(TYPE endpoint,放在 /endpoints)在 SQL 上多了模板参数。api_top_pages.pipe 是一个典型例子,它完整体现了参数声明、条件模板与分页:

%
select
    case when post_uuid = 'undefined' then '' else post_uuid end as post_uuid,
    pathname,
    uniqExact(session_id) as visits
from _mv_hits h
inner join filtered_sessions fs
    on fs.session_id = h.session_id
where
    site_uuid = {{String(site_uuid, 'mock_site_uuid', description="Tenant ID", required=True)}}
    {% if defined(member_status) %}
        and member_status IN (
            select arrayJoin(
                {{ Array(member_status, "'undefined', 'free', 'paid'", ...) }}
                || if('paid' IN {{ Array(member_status) }}, ['comped', 'gift'], [])
            )
        )
    {% end %}
    ...
group by post_uuid, pathname
order by visits desc
limit {{ Int32(skip, 0) }},{{ Int32(limit, 50) }}

文件开头还声明了只读 Token:TOKEN "stats_page" READ / TOKEN "axis" READ。第一行的 % 是参数化查询的强制标记,{% if defined(...) %} 是 Tornado 条件模板,required=True/False 与硬编码默认值都符合下文 SQL 规则。

端点测试与 URL 规则

endpoint-files 规则 补充了端点专属要求:

  • tb endpoint data 而非 tb pipe data 测试端点——前者像真实消费者一样调用端点,包含参数校验与输出格式化:
tb endpoint data my_endpoint
tb endpoint data my_endpoint --start_date 2024-01-01 --end_date 2024-01-31
  • tb endpoint ls 列出所有端点及其 URL;拼 URL 时按需带上动态参数,日期格式为:
    • DateTime64YYYY-MM-DD HH:MM:SS.MMM
    • DateTimeYYYY-MM-DD HH:MM:SS
    • DateYYYYMMDD
  • 取全部端点的 OpenAPI 定义:curl <api_base_url>/v0/pipes/openapi.json?token=<token>

SQL 规则:SELECT-only 与严格的参数处理

sql 规则 是全部查询的底线约束。

核心四原则:尽早过滤、读尽可能少的数据;只选需要的列;把复杂计算尽量放到流水线后段;优先使用 ClickHouse 函数(只允许受支持的函数)。

查询要求

  • SQL 必须是合法的 ClickHouse SQL 加 Tinybird 模板(Tornado);
  • 只允许 SELECT 语句;
  • 避免 CTE,改用节点(node)或子查询;
  • 禁用 system 表(system.tablessystem.datasourcesinformation_schema.tables);
  • 禁用 CREATE/INSERT/DELETE/TRUNCATEcurrentDatabase()

参数与模板规则

  • 只要用了参数,查询必须以单独一行的 % 开头;
  • 参数函数白名单:StringDateTimeDateFloat32Float64IntIntegerUInt8/16/32/64/128/256Int8/16/32/64/128/256
  • 参数名不能与列名相同;
  • 默认值必须硬编码;
  • 参数永远不加引号;defined() 检查中也不要给参数名加引号。
# 错误写法(缺少 % 声明行)
SELECT * FROM events WHERE session_id={{String(my_param, "default")}}

# 正确写法
%
SELECT * FROM events WHERE session_id={{String(my_param, "default")}}

JOIN 与聚合规则:JOIN 和 GROUP BY 之前先过滤;避免在未过滤的情况下 JOIN 超过 100 万行的表;避免嵌套聚合,用子查询替代;AggregateFunction 列搭配 -Merge 组合器使用。

操作顺序:WHERE 过滤 → 选择需要的列 → JOIN → GROUP BY/聚合 → ORDER BY → LIMIT。

外部表:支持 Iceberg 与 Postgres 直读,密钥走 tb_secret,且主机与端口不要拆成多个 secret:

FROM iceberg('s3://bucket/path/to/table', {{tb_secret('aws_access_key_id')}}, {{tb_secret('aws_secret_access_key')}})
FROM postgresql({{ tb_secret("db_host_port") }}, 'database', 'table', {{tb_secret('db_username')}}, {{tb_secret('db_password')}}, 'schema_optional')

端点优化清单:先取证,再动手

endpoint-optimization 规则 是一套“证据驱动”的优化流程,分三步:收集运行时证据、套用结构化规则、按运行时阈值处理性能规则。

第一步:收集运行时证据

优化前先取证,来源有三处:

  1. 端点源码:工作区内的 SQL、数据源、物化视图与管道;
  2. pipe_stats_rtSELECT * FROM tinybird.pipe_stats_rt WHERE pipe_name = 'endpoint_name',检查执行时长分位数(p50/p90/p95/p99)、read_bytesrows_read 与错误数;
  3. 查询计划:带 ?explain=true 调用端点(例如 https://$TB_HOST/v0/pipes/endpoint_name?explain=true),检查 JOIN 策略、聚合阶段、索引使用与分区裁剪。

规则明确:数据少于 1 万行或不足 50 MB 的数据源直接忽略。若 tinybird.pipe_stats_rttinybird.pipe_stats 中没有需要的信息,再查 system.query_log

第二步:结构化规则(无需运行时证据,检出即改)

问题 修复
SELECT * 或未使用的列 显式只选需要的列(减少 I/O、解压与缓存压力)
过大的数据类型 用最小安全类型;低基数字符串用 LowCardinality;用默认值代替 Nullable
不必要的 Nullable 列永不为空时把 Nullable(T) 改为 T
ORDER BY 键序不当 以低基数/时间列开头;多租户场景避免时间戳做首键
冗余类型转换 去掉同类型 cast;在入库阶段修类型(cast 浪费 CPU 且可能挡住分区裁剪)
字符串过度物化 入库时就把需要的字符串属性提取为类型化列
过滤晚于 JOIN/聚合 把过滤尽可能前置

Ghost 的 analytics_events 数据源actionsite_uuid 等列使用 LowCardinality(String),正是这条规则在数据模型层面的落实。

第三步:运行时规则(超过阈值才处理)

每条规则都给出触发条件与修复手段,摘录关键的几条:

  • 查询时聚合:p95 > 5s、或聚合主导 EXPLAIN、或内存 > 60%、或 OOM/超时 → 用物化视图预计算;
  • 查询时 JOIN:p95 > 5s、或 JOIN 主导 EXPLAIN、或内存飙升 → 把 JOIN 移到入库时(物化视图),或反范式化;
  • 排序键缺失/不当:读取超过 10% granule 且 p95 > 3s,或 rows_read/rows_returned > 100x → 重建数据源,ORDER BY 对齐高选择性过滤列;
  • PREWHERE 提前过滤:rows_read/rows_returned > 50x,或 p95 > 3s,或 EXPLAIN 显示过滤过晚 → 把高选择性过滤推进 PREWHERE
  • 数据跳过索引:过滤落在非主键列且 rows_read/rows_returned > 100x,或 p95 > 3s → 加 skip index 并用 EXPLAIN 验证;
  • 大 GROUP BY:p95 > 5s、聚合内存 > 50% 或 OOM → 入库时预聚合;
  • 查询时正则:p95 > 3s 或 CPU > 70% → 移到入库阶段;
  • 无 TTL 的无界历史:p95 逐周上升、同一查询 rows_read 持续增长 → 建带 TTL 的数据源;
  • 分区裁剪失效:扫描超过 20% 分区、或 p95 > 3s → 重建对齐的分区键,并在过滤条件中包含分区键列;
  • ORDER BY + LIMIT 未下推:排序行数 >> LIMIT(>100x)→ 重构查询或入库时预物化 top-k;
  • DISTINCT 代替 GROUP BY:p95 > 5s 或内存 > 50% → 换入库时聚合或 GROUP BY
  • FINAL 滥用:p95 > 3s 或 rows_read >> rows_returned → 通过入库时保证正确性(lambda 架构)去掉 FINAL
  • 查询时昂贵 JSON 提取:p95 > 3s 或 CPU > 70% → 入库时提取 JSON 字段为类型化列;
  • 大 IN 列表:p95 > 3s 或查询规划时间长 → 换 lookup 数据源或入库时物化;
  • 精确去重计数:精确 COUNT(DISTINCT) 且 p95 > 5s、内存 > 50% 或 OOM → 在可接受时用 uniqHLL12 等近似函数。

验证标准:持续跟踪 tinybird.pipe_stats_rttinybird.pipe_stats;成功指标是延迟下降、read_bytes 下降、read_bytes/write_bytes 比改善。

规则自带的两个模板

物化视图模板(预计算指标的标准形态):

NODE materialized_view_name
SQL >
  SELECT toDate(timestamp) as date, customer_id, countState(*) as event_count
  FROM source_table
  GROUP BY date, customer_id

TYPE materialized
DATASOURCE mv_datasource_name
ENGINE "AggregatingMergeTree"
ENGINE_PARTITION_KEY "toYYYYMM(date)"
ENGINE_SORTING_KEY "customer_id, date"

优化后的端点查询模板(注意 %、类型化参数与默认值):

NODE endpoint_query
SQL >
  %
  SELECT date, sum(amount) as daily_total
  FROM events
  WHERE customer_id = {{ String(customer_id) }}
    AND date >= {{ Date(start_date) }}
    AND date <= {{ Date(end_date) }}
  GROUP BY date
  ORDER BY date DESC

物化视图规则:State/AggregateFunction 配对与 JSON 单次解析

materialized-files 规则 的要点:

  • 默认不创建,除非明确要求;创建在 /materializations 下;
  • TYPE MATERIALIZED 并通过 DATASOURCE 指定目标数据源;
  • 管道中用 State 修饰符,目标数据源用 AggregateFunction;读取 AggregateFunction 列时用 Merge 修饰符;
  • ENGINE_SORTING_KEY 放全部维度,按基数从低到高排序。

最小示例(规则原文):

NODE daily_sales
SQL >
    SELECT toStartOfDay(starting_date) day, country, sumState(sales) as total_sales
    FROM teams
    GROUP BY day, country

TYPE MATERIALIZED
DATASOURCE sales_by_hour

目标数据源:

SCHEMA >
    `total_sales` AggregateFunction(sum, Float64),
    `sales_count` AggregateFunction(count, UInt64),
    `dimension_1` String,
    `dimension_2` String,
    `date` DateTime

ENGINE "AggregatingMergeTree"
ENGINE_PARTITION_KEY "toYYYYMM(date)"
ENGINE_SORTING_KEY "date, dimension_1, dimension_2"

Ghost 的 mv_session_data_v2.pipe 是同一模式的真实案例:countState() 统计页面浏览数、minState/maxState(timestamp) 记录会话首末访问、argMinState(source, timestamp) 取会话内时间最早一条的 UTM 来源——目标数据源 _mv_session_data_v2 中对应 AggregateFunction 列。

JSON 提取:解析一次,而不是一字段一次

适用场景:同一条查询对同一 JSON/字符串列多次调用 JSONExtractString / JSONExtractInt / JSONExtractBool / JSONExtractFloat / simpleJSONExtractString / visitParam*。每次调用都会从头重新解析原始 JSON——N 个字段就是 N 次完整解析,每行如此。在物化视图中最昂贵,因为它对该管道生命周期内插入的每个 block 都要执行一遍。

做法:先用一次 JSONExtract(...) 把 JSON 解析成类型化 Tuple,再用 getSubcolumn 逐字段读取。

反例(每字段一次解析):

NODE typed_events
SQL >
    SELECT
        at AS timestamp,
        visitParamExtractString(payload, 'field_a') AS field_a,
        visitParamExtractInt(payload, 'field_b') AS field_b,
        visitParamExtractBool(payload, 'field_c') AS field_c,
        simpleJSONExtractString(payload, 'field_d') AS field_d
    FROM raw_events

TYPE MATERIALIZED
DATASOURCE typed_events_ds

正例(总共一次解析):

NODE typed_events
SQL >
    WITH
        JSONExtract(payload, 'Tuple(
            field_a String,
            field_b Int64,
            field_c Bool,
            field_d String
        )') AS payload_json
    SELECT
        at AS timestamp,
        getSubcolumn(payload_json, 'field_a') AS field_a,
        getSubcolumn(payload_json, 'field_b') AS field_b,
        getSubcolumn(payload_json, 'field_c') AS field_c,
        getSubcolumn(payload_json, 'field_d') AS field_d
    FROM raw_events

TYPE MATERIALIZED
DATASOURCE typed_events_ds

补充细节:缺失字段取类型默认值;多个派生表达式依赖同一字段时,用 WITH 别名复用。

物化视图的三个常见陷阱

  1. 物化视图本质是 insert 触发器:对源数据源的 delete/truncate 不会影响物化视图;
  2. 转换与入库针对源数据源插入的每个 block 执行,因此 GROUP BYORDER BYDISTINCTLIMIT 等操作需要 AggregatingMergeTreeSummingMergeTree 这类能处理聚合的引擎配合;
  3. 用 JOIN 生成的物化视图数据源,只在 FROM 侧数据源发生新写入时才会自动更新。

进阶:物化视图中右侧 JOIN 的预过滤

materialized-join-prefilter 规则 解决一个隐蔽的摄取瓶颈:物化视图在源数据源(FROM 最左表)每插入一个 block 时会以该 block 为 FROM 重新执行管道 SQL,而 JOIN / ASOF JOIN 右侧的表如果没有显式限制,会被全量扫描——右表越大,每次插入越贵,最终导致摄取变慢、延迟,甚至 MEMORY_LIMIT_EXCEEDED

适用信号:JOIN 右侧是随时间无界增长的数据源;出现慢插入、摄取滞后、源数据源内存尖峰;JOIN 本身已有高选择性条件(键等值、时间边界)。

模式:把右侧数据源替换成子查询,只保留“可能与当前插入批次匹配”的行。两个过滤叠加:

  1. 键预过滤——右表只保留 join 键出现在左表插入批次中的行;
  2. 时间预过滤——对 left.time >= right.timeASOF join,把 right.time 约束在插入批次的 [min(left.time) - INTERVAL N <unit>, max(left.time)] 区间内。下界是右行与左行之间可接受的最大间隔,应写成显式、可配置的常量。

关键机制:子查询里对左侧源表的引用解析为正在插入的 block 而非全表——这是预过滤廉价的原因。规则给出了完整的改造前后对照(ASOF LEFT JOIN 场景):

改造前,enrichment_table 每次插入被全表扫描:

NODE mv_node
SQL >
    SELECT
        e.tenant_id, e.entity_id, e.event_name, e.event_time,
        x.source_time AS resolved_time
    FROM events_table e
    ASOF LEFT JOIN enrichment_table x
        ON e.tenant_id = x.tenant_id
        AND e.entity_id = x.entity_id
        AND e.ref_id = x.ref_id
        AND e.event_time >= x.source_time
    WHERE e.event_name IN ('event_x', 'event_y')

TYPE materialized
DATASOURCE mv_target

改造后,右表按“批次内出现的键 + 30 天时间窗”双重约束:

NODE mv_node
SQL >
    SELECT
        e.tenant_id, e.entity_id, e.event_name, e.event_time,
        x.source_time AS resolved_time
    FROM events_table e
    ASOF LEFT JOIN (
        SELECT tenant_id, entity_id, ref_id, source_time
        FROM enrichment_table
        WHERE source_time >= (
                SELECT min(event_time)
                FROM events_table
                WHERE event_name IN ('event_x', 'event_y')
            ) - INTERVAL 30 DAY
          AND source_time <= (
                SELECT max(event_time)
                FROM events_table
                WHERE event_name IN ('event_x', 'event_y')
            )
          AND (tenant_id, entity_id, ref_id) IN (
                SELECT tenant_id, entity_id, ref_id
                FROM events_table
                WHERE event_name IN ('event_x', 'event_y')
            )
    ) x
        ON e.tenant_id = x.tenant_id
        AND e.entity_id = x.entity_id
        AND e.ref_id = x.ref_id
        AND e.event_time >= x.source_time
    WHERE e.event_name IN ('event_x', 'event_y')

TYPE materialized
DATASOURCE mv_target

每个右侧 JOIN 都要独立包一层子查询,不要共用。规则附带的检查清单:子查询只选 JOIN 与 SELECT 用到的列;键预过滤的 IN 元组复刻左表源相同的 WHEREINTERVAL N <unit> 是单一显式字面量,便于日后调整;内层子查询对左表源的 WHERE 必须与外层管道一致,保证插入 block 被一致读取。

注意事项(Gotchas)

  • 时间下界是“成本换正确性”的交易——早于 min(left.time) - N 的右行即使本应命中 ASOF 匹配也会被排除,N 要按真实数据间隔取够大,并写进文档;
  • 多个右侧 JOIN 各自语义不同,预过滤不可复用;
  • 内层子查询中的键提取必须与外层完全一致(同样的 cast、同样的 JSONExtract/toInt64OrZero 包装),否则 IN 元组匹配不上;
  • 预过滤不改变窗口内数据的 MV 正确性,但改变窗口外数据的结果——这一点要写进管道级 DESCRIPTION 中显式声明。

去重与 Lambda 架构

deduplication-patterns 规则 给出了四套策略的选择矩阵:

策略 适用场景
查询时去重(argMaxLIMIT BY、子查询) 原型阶段或小数据集
ReplacingMergeTree 大数据集、需要按 key 取最新行
周期快照(Copy Pipes) 新鲜度不关键、需要 rollup 或不同排序键
Lambda 架构 既要新鲜度,又需要 MV 处理不了的复杂转换

维度表/小表通常“周期全量替换”最合适。

查询时去重的三种写法

-- argMax:按 key 取最新值
SELECT post_id, argMax(views, updated_at) as views
FROM posts GROUP BY post_id

-- LIMIT BY
SELECT * FROM posts ORDER BY updated_at DESC LIMIT 1 BY post_id

-- 子查询
SELECT * FROM posts WHERE (post_id, updated_at) IN (
    SELECT post_id, max(updated_at) FROM posts GROUP BY post_id
)

ReplacingMergeTree 的数据源配置与三条硬约束:

ENGINE "ReplacingMergeTree"
ENGINE_SORTING_KEY "unique_id"
ENGINE_VER "updated_at"
ENGINE_IS_DELETED "is_deleted"  -- 可选,UInt8:1=已删除,0=有效
  • 查询必须带 FINAL 或改用其他去重手段;
  • 去重发生在 merge 期间(异步、不可控);
  • 不要在 ReplacingMergeTree 之上建 AggregatingMergeTree 物化视图——MV 只能看到流入的 block,看不到 merge 后的状态,重复行会一直存在。
SELECT * FROM posts FINAL WHERE post_id = {{Int64(post_id)}}

快照式去重(Copy Pipes) 的适用场景:ReplacingMergeTree + FINAL 太慢、需要随更新变化的不同排序键、需要下游 MV 做 rollup。copy_mode 默认 append;表不大且无法控制重复发生时用 COPY_MODE replace 做全量刷新;能控制重复产生并可增量处理时保持 append

NODE generate_snapshot
SQL >
    SELECT post_id, argMax(views, updated_at) as views, max(updated_at) as updated_at
    FROM posts_raw
    GROUP BY post_id

TYPE COPY
TARGET_DATASOURCE posts_snapshot
COPY_SCHEDULE 0 * * * *
COPY_MODE replace

Lambda 架构适合:对 ReplacingMergeTree 做聚合(MV 失效,原因同上)、需要全表扫描的窗口函数、CDC 工作负载、uniqState 性能问题、查询时需要 JOIN 的端点。模式分三层:

  1. 批处理层:Copy Pipe 周期生成去重快照/中间表;
  2. 实时层:查询自上次快照以来的新数据;
  3. 服务层:UNION ALL 合并两者:
SELECT * FROM posts_snapshot
UNION ALL
SELECT post_id, argMax(views, updated_at) as views, max(updated_at) as updated_at
FROM posts_raw
WHERE updated_at > (SELECT max(updated_at) FROM posts_snapshot)
GROUP BY post_id

新鲜度与成本的权衡:Copy Pipe 越频繁快照越新鲜但成本越高;不频繁则批层偏旧,但实时层可以补齐缺口——按查询模式与数据量取平衡。

argMax 的 null 陷阱argMaxMerge 会优先选非空值,即使其时间戳更旧。规则给出的绕法是聚合前把 null 转成 epoch 哨兵值,再在下游查询中处理哨兵:

SELECT post_id,
    argMaxState(CASE WHEN flagged_at IS NULL THEN toDateTime('1970-01-01 00:00:00') ELSE flagged_at END, updated_at) as flagged_at
FROM posts
GROUP BY post_id

Copy Pipe 与 Sink Pipe:两类“非默认”管道

Copy Pipe

copy-files 规则:默认不创建、除非明确要求;放在 /copies;不显式要求时不写 COPY_SCHEDULE;必须 TYPE COPY + TARGET_DATASOURCEcopy_mode 默认 append,但规则建议显式写出;另一个取值是 replace

DESCRIPTION Copy Pipe to export sales hour every hour to the sales_hour_copy Data Source

NODE daily_sales
SQL >
    %
    SELECT toStartOfDay(starting_date) day, country, sum(sales) as total_sales
    FROM teams
    WHERE day BETWEEN toStartOfDay(now()) - interval 1 day AND toStartOfDay(now())
    and country = {{ String(country, 'US')}}
    GROUP BY day, country

TYPE COPY
TARGET_DATASOURCE sales_hour_copy
COPY_SCHEDULE 0 * * * *
COPY_MODE append

Sink Pipe

sink-files 规则:默认不创建、除非明确要求;放在 /sinks;支持的外部系统只有 Kafka、S3、GCS;Sink 依赖 connection,能复用就复用已有 connection;不显式要求时不写 EXPORT_SCHEDULE;必须 TYPE SINK + EXPORT_CONNECTION_NAME

DESCRIPTION Sink Pipe to export sales hour every hour using my_connection

NODE daily_sales
SQL >
    %
    SELECT toStartOfDay(starting_date) day, country, sum(sales) as total_sales
    FROM teams
    WHERE day BETWEEN toStartOfDay(now()) - interval 1 day AND toStartOfDay(now())
    and country = {{ String(country, 'US')}}
    GROUP BY day, country

TYPE sink
EXPORT_CONNECTION_NAME "my_connection"
EXPORT_BUCKET_URI "s3://tinybird-sinks"
EXPORT_FILE_TEMPLATE "daily_prices"
EXPORT_SCHEDULE "*/5 * * * *"
EXPORT_FORMAT "csv"
EXPORT_COMPRESSION "gz"
EXPORT_STRATEGY "truncate"

Connection

connection-files 规则:内容不能为空、名称唯一、属性名不缩进;仅支持 kafka、gcs、s3 三种类型——用户要求其他类型时如实报告、不创建。示例:

# Kafka
TYPE kafka
KAFKA_BOOTSTRAP_SERVERS {{ tb_secret("PRODUCTION_KAFKA_SERVERS", "localhost:9092") }}
KAFKA_SECURITY_PROTOCOL SASL_SSL
KAFKA_SASL_MECHANISM PLAIN
KAFKA_KEY {{ tb_secret("PRODUCTION_KAFKA_USERNAME", "") }}
KAFKA_SECRET {{ tb_secret("PRODUCTION_KAFKA_PASSWORD", "") }}

# S3
TYPE s3
S3_REGION {{ tb_secret("PRODUCTION_S3_REGION", "") }}
S3_ARN {{ tb_secret("PRODUCTION_S3_ARN", "") }}

# GCS(服务账号)
TYPE gcs
GCS_SERVICE_ACCOUNT_CREDENTIALS_JSON {{ tb_secret("PRODUCTION_GCS_SERVICE_ACCOUNT_CREDENTIALS_JSON", "") }}

# GCS(HMAC)
TYPE gcs
GCS_HMAC_ACCESS_ID {{ tb_secret("gcs_hmac_access_id") }}
GCS_HMAC_SECRET {{ tb_secret("gcs_hmac_secret") }}

测试与 Fixture 规范

tests 规则 约定:

  • 测试文件名必须与 pipe 同名;测试文件内场景名唯一;
  • 参数格式 param1=value1&param2=value2;用户给出的参数保持原样大小写与格式;无参数时创建单个空参数测试;
  • 期望结果基于 fixture 数据设计,不要通过调用端点或 SQL 反推数据;创建测试前先分析端点表所用的 fixture 文件;
  • expected_result 恒为空字符串,由工具回填;
  • 只在明确要求时才创建测试(例如“给这个端点建测试”);如果要求只是“测试”或“调用”端点,用 tb endpoint data 而不是建测试。

测试文件格式:

- name: kpis_single_day
  description: Test hourly granularity for a single day
  parameters: date_from=2024-01-01&date_to=2024-01-01
  expected_result: ''

Fixture:放在 /fixtures,文件名与被填充的数据源同名(fixtures/<datasource_name>.ndjson.csv);用 tb datasource append <name> --file fixtures/<name>.ndjson 加载到 local 或 branch 环境;fixture 要小而确定、覆盖测试所需场景(边界、日期范围、不同参数值),并纳入版本控制。Ghost 仓库的 analytics_events.ndjson 就是这样一个 fixture:每行一个 page_hit 事件,payload 里带 site_uuidmember_statuspost_uuidutm_* 等字段,与 analytics_events 数据源的 SCHEMA 一一对应,覆盖了 bot 访问、free/paid/comped 会员、UTM 参数缺失等多种测试场景。

运行测试

tb test run                        # 跑全部测试
tb test run tests/my_endpoint      # 跑指定测试文件
tb test update tests/my_endpoint   # 用当前输出更新期望结果

Ghost 的 tests 目录 与 30 个端点一一对应(api_kpis.yamlapi_top_pages.yamlapi_active_visitors.yaml……),文件名即 pipe 名的约定在真实工程中得到了严格遵循。

总结:一套可检查的 Tinybird 工作流

把技能规范的各规则串起来,完整的日常工作流是:

  1. tb info 确认 CLI 上下文与 .tinyb 凭证位置;
  2. datasources/(Landing)、pipes//materializations/(Cleaned/Aggregation)、endpoints/(API)、copies/sinks/connections/ 的分层布局组织 datafile,属性名不缩进,指令名大写;
  3. 数据源默认 MergeTree、物化目标 AggregatingMergeTree;schema 变更不兼容时加 FORWARD_QUERY,完成后清理;TTL 与分区键对齐,配合 ttl_only_drop_parts=1
  4. SQL 保持 SELECT-only、% 声明参数、类型化参数函数 + 硬编码默认值、参数名不撞列名、参数不加引号;操作顺序遵循“过滤→选列→JOIN→聚合→排序→LIMIT”;
  5. 物化视图中用 State 聚合、目标端用 AggregateFunction 列、读取用 Merge 组合器;同一 JSON 列多字段提取时用 JSONExtract(...Tuple) 单次解析;右表无界增长的 JOIN 用键+时间双预过滤包裹;
  6. 优化端点时先查 tinybird.pipe_stats_rt?explain=true 取证,再按结构化规则直接修、按运行时阈值(p95、rows_read/rows_returned、内存/CPU 占比)决定是否上物化视图、PREWHERE、skip index 或近似函数;
  7. 去重按数据规模在查询时去重、ReplacingMergeTree、Copy Pipe 快照与 Lambda 架构之间选择;记住 MV 只看到 block、看不到 merge 状态这一根本约束;
  8. tb build 指向 dev_mode 配置的 local/branch;生产部署走 tb --cloud deploy(CI 用 --check 预校验);测试用 tb endpoint data 手工验证、tb test run 回归,期望值从 fixture 推导。

这套规范与 Ghost 流量分析工程的实际代码相互印证:从 analytics_events 落地层LowCardinality 类型与 FORWARD_QUERY 迁移,到 mv_session_data_v2 的 State 聚合链,再到 api_top_pages 端点 的参数模板与分页写法,都能在上述规则中找到对应条款。遵循这套 datafile 规范,可以在保证端点低延迟读写的同时,让 Tinybird 项目保持可测试、可构建、可部署的工程形态。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.13 K
2.75 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
857
1.35 K
docsdocs
暂无描述
Markdown
897
5.8 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
529
593
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
915
1.83 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.58 K
1.01 K
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.35 K
1.46 K
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.01 K
515
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
547
388