一个完善的数据仓库会面临一个反复出现的问题:接入新渠道。
不管是新入驻了拼多多、上线了私域小程序,还是收购了一个新品牌——每次接入的套路都一样:确认数据源 → 建表 → 配同步 → 清洗 → 汇总。但每次都要数仓开发花 3-5 天手动做一遍。
这篇文章讲我们是怎么做的——给 AI 一份标准格式的接入文档,AI 自动完成剩余的全部开发工作。
但前提是:数仓先有 Harness。
一、一个问题
一个典型的新渠道接入流程,长这样:
确认数据源信息 ↓ 走查源表结构 ↓ 建 ODS 层表(DDL) ↓ 配置同步任务(Flink CDC / DataX) ↓ 写 DWD 清洗 SQL ↓ 写 DWS 汇总 SQL(含 JOIN 维表) ↓ 配置调度依赖 ↓ 写质检 SQL ↓ 联调验证上线传统做法,这一步大概 5-6 天。其中大部分工作是重复性的——表结构换了个前缀、字段名稍有差异、但套路完全一致。
问题不是"能不能让 AI 做",而是怎么让 AI 做得放心。
二、前提:数仓的 Harness 体系
AI 能自动接入,前提是数仓已经有完备的 Harness 文档体系。否则 AI 不知道表怎么命名、分区怎么设、维度如何对齐。
Harness 不是为 AI 建的。在 AI 进来之前,数仓团队就已经需要这些规范。AI 是后来加入的受益者。
一个完整的数仓 Harness 包含 10 份文档,分为四类:
核心三件套:
领域文档:
工程规范:
安全底线:
这些文档在 AI 接入之前就已经存在。AI 介入后,自动读取全部文档,在约束范围内工作。
核心红线——写在 ARCHITECTURE.md 里的铁律:
┌─────────────────────────────────────────────┐│ 架构铁律(不可违反) ││ ││ 数据安全 ││ · PII 字段必须脱敏 ││ · 禁止将生产数据导出到本地 ││ ││ 操作红线 ││ · 禁止 DROP TABLE ││ · 禁止 TRUNCATE ││ · 禁止 DELETE 不带 WHERE ││ · 禁止裸写 INSERT 不带 PARTITION ││ ││ 分层约束 ││ · 禁止 ODS → ADS 跨层依赖 ││ · 禁止 ADS 层表被其他 ETL 作为上游 │└─────────────────────────────────────────────┘
三、核心契约:数据源接入设计文档
准备好了上述规范,架构师需要做的事情只有一件:按标准格式写一份接入设计文档。
以接入"拼多多旗舰店"为例:
# 数据源接入文档:拼多多旗舰店data_source: name: pdd_official_store display_name: 拼多多旗舰店 type: MySQL(Binlog) db: pdd_order_db tables: - name: order cn_name: 订单表 description: 拼多多订单主表,每笔订单一条记录 schema: - source_field: id target_field: order_id type: String comment: 订单ID - source_field: order_sn target_field: out_order_no type: String comment: 拼多多平台订单号 - source_field: gmv target_field: pay_amount type: Decimal(18,2) comment: 实付金额 - source_field: sku_id target_field: sku_id type: String comment: 商品SKU ID - source_field: receiver_province target_field: province type: String comment: 收货省 - source_field: receiver_city target_field: city type: String comment: 收货市 - source_field: order_time target_field: order_time type: DateTime comment: 下单时间 - source_field: update_time target_field: update_time type: DateTime comment: 更新时间 - source_field: status target_field: order_status type: Int comment: '订单状态:1待支付 2已支付 3已发货 4已完成 5已取消' - name: refund cn_name: 退款表 description: 退款记录表 schema: - source_field: refund_id target_field: refund_id type: String comment: 退款ID - source_field: order_id target_field: order_id type: String comment: 关联订单ID - source_field: refund_amount target_field: refund_amount type: Decimal(18,2) comment: 退款金额 - source_field: refund_time target_field: refund_time type: DateTime comment: 退款时间 dwd_mapping: # 字段统一映射 field_mapping: - channel_field: order_sn unified_field: out_order_no note: 各渠道保留原始单号 - channel_field: gmv unified_field: pay_amount note: 统一口径 - channel_field: status unified_field: order_status enum_mapping: 1: "pending_pay" 2: "paid" 3: "shipped" 4: "completed" 5: "cancelled" dws_requirements: - metric: gmv dimension: [brand_id, category_id, dt] - metric: order_count dimension: [channel, dt] - metric: 客单价 dimension: [dt]这份文档大约 60-80 行,架构师写一次大概 半天。
核心原则:架构师只写"做什么"——表和字段。AI 负责"怎么做"——表名怎么命、分区怎么设、SQL 怎么写、调度怎么配。
四、AI 接单后的完整工作流
AI 读取接入文档后,按顺序完成 6 个步骤。每一步都有 Harness 约束兜底。
Step 1:生成 ODS 建表 DDL
AI 根据 tables.schema 里的字段列表,结合命名规范和分层设计规范,生成建表语句:
CREATE TABLE ods_trade_order_pdd_di ( order_id String COMMENT '订单ID', out_order_no String COMMENT '拼多多平台订单号', pay_amount Decimal(18,2) COMMENT '实付金额', sku_id String COMMENT '商品SKU ID', province String COMMENT '收货省', city String COMMENT '收货市', order_time DateTime COMMENT '下单时间', update_time DateTime COMMENT '更新时间', order_status Int COMMENT '订单状态', dt Date COMMENT '分区日期') ENGINE = ReplacingMergeTree(update_time)PARTITION BY toYYYYMM(dt)ORDER BY (order_id, dt);Harness 校验:
• 表名 ods_trade_order_pdd_di→ 符合{分层}_{域}_{实体}_{粒度}规范 ✅• 分区策略 toYYYYMM(dt)→ 符合 ODS 层日分区约定 ✅• 字段注释覆盖率 100% → 符合 ETL 规范 ✅ • ENGINE = ReplacingMergeTree → 符合 Binlog 同步场景 ✅
每个校验点都有对应的规范文档支撑。AI 不是凭"经验"写的,是按文档约束生成的。
Step 2:生成同步任务配置
接入文档中 type: MySQL(Binlog) 告诉 AI 需要配置实时同步。AI 生成 Flink CDC 任务:
{ "job_name": "sync_pdd_order_to_ods", "source": { "connector": "mysql-cdc", "hostname": "${pdd_mysql_host}", "port": 3306, "database-name": "pdd_order_db", "table-name": "order", "server-id": 5401 }, "sink": { "connector": "clickhouse", "hosts": ["${clickhouse_host}"], "table": "ods_trade_order_pdd_di", "username": "${clickhouse_user}", "password": "${clickhouse_password}", "partition-strategy": "day" }, "transforms": [ { "type": "route", "source-table": "pdd_order_db.order", "sink-table": "ods_trade_order_pdd_di" } ]}Harness 校验:
• 字段映射完整性 → source 24 个字段 vs sink 24 个字段 ✅ • 密码引用 ${clickhouse_password}而非明文 → 符合凭证铁律 ✅• server-id 不与已有任务冲突 → 符合 Flink 配置规范 ✅
⚠️ 这里 Hook 发挥了作用。 如果 AI 不小心把数据库密码直接写进了配置文件,PostToolUse 的 Output Guard 会立即捕获并告警:
[Harness Output Guard] 检测到敏感信息: 数据库连接串泄露[Harness Output Guard] 已自动清除并记录审计这不是 AI 不可靠——这是 Harness 在兜底。AI 一个人写 6 个步骤,中间出点小错很正常。关键是被 Hook 截住,不让它流到生产环境。
Step 3:生成 DWD 清洗 SQL
利用 dwd_mapping 里的字段映射规则,把 ODS 层数据清洗到 DWD 宽表:
INSERT INTO dwd_trade_order_detail ( order_id, out_order_no, channel, pay_amount, sku_id, province, city, order_time, order_status, dt)SELECT order_id, out_order_no, 'pdd' AS channel, pay_amount, sku_id, province, city, order_time, CASE order_status WHEN 1 THEN 'pending_pay' WHEN 2 THEN 'paid' WHEN 3 THEN 'shipped' WHEN 4 THEN 'completed' WHEN 5 THEN 'cancelled' END AS order_status, dtFROM ods_trade_order_pdd_diWHERE dt = '${biz_date}';Harness 校验:
• 渠道值 'pdd'→ 与接入文档一致,并已加入域划分文档的渠道枚举 ✅• 状态枚举 5 种全部覆盖 → 与接入文档的 enum_mapping 一致 ✅ • INSERT 带 PARTITION( dt = '${biz_date}')→ 符合红线规则 ✅• 字段全部显式列出,没有 SELECT *→ 符合 ETL 规范 ✅
这里 Hook 也在工作: 假设 AI 忘了加 WHERE dt = '${biz_date}',PreToolUse 中的检测规则会将无分区条件的 INSERT 标记为高危操作:
[Harness Guard] WARNING: INSERT INTO without explicit PARTITION condition[Harness Guard] 已自动补全 dt = '${biz_date}'Step 4:生成 DWS 汇总 SQL
根据 dws_requirements 自动生成多维度汇总 SQL:
INSERT INTO dws_trade_order_daily ( channel, brand_id, category_id, gmv, order_count, order_member_count, avg_order_amt, dt)SELECT 'pdd' AS channel, sku.brand_id, sku.category_id, SUM(pay_amount) AS gmv, COUNT(DISTINCT order_id) AS order_count, COUNT(DISTINCT member_id) AS order_member_count, SUM(pay_amount) / COUNT(DISTINCT order_id) AS avg_order_amt, dtFROM dwd_trade_order_detail ordLEFT JOIN dim_sku sku ON ord.sku_id = sku.sku_idWHERE ord.channel = 'pdd' AND dt = '${biz_date}'GROUP BY channel, brand_id, category_id, dt;Harness 校验:
• gmv = SUM(pay_amount)→ 与指标体系文档中 gmv 口径一致 ✅• LEFT JOIN dim_sku→ JOIN 路径在维度建模规范中 ✅• 分组维度 brand_id, category_id→ 与接入文档的 dimension 一致 ✅• 新增字段 avg_order_amt→ 注意:客单价是衍生指标,不在原始数据中 ✅
这里有一个关键设计:AI 没有直接写 pay_amount / COUNT(DISTINCT order_id),因为它知道命名规范要求衍生指标必须用可读别名,且与指标体系文档的口径定义一致。
Step 5:生成调度配置 + 质检脚本
# 调度依赖pipeline: sync_pdd_to_ods: type: Flink cron: continuous # 实时同步 ods_to_dwd_pdd: type: ClickHouse cron: "30 1 * * *" # 每天 1:30 depends_on: - table: ods_trade_order_pdd_di check: "SELECT count(*) > 0 FROM ods_trade_order_pdd_di WHERE dt = '昨日'" dwd_to_dws_pdd: type: ClickHouse cron: "0 3 * * *" # 每天 3:00 depends_on: - table: dwd_trade_order_detail source_channel: pdd check: "SELECT count(*) > 0 FROM dwd_trade_order_detail WHERE channel='pdd' AND dt = '昨日'"# 质检规则quality_checks: - name: pdd_order_null_check description: 订单ID为空检查 sql: > SELECT count(*) AS null_count FROM ods_trade_order_pdd_di WHERE order_id IS NULL AND dt = '昨日' threshold: 0 severity: error - name: pdd_ods_dwd_row_check description: ODS 到 DWD 行数一致性 sql: > SELECT abs(ods_cnt - dwd_cnt) AS diff FROM ( SELECT count(*) AS ods_cnt FROM ods_trade_order_pdd_di WHERE dt = '昨日' ) a, ( SELECT count(*) AS dwd_cnt FROM dwd_trade_order_detail WHERE channel = 'pdd' AND dt = '昨日' ) b threshold: 0.05 severity: warn - name: pdd_gmv_wow_check description: GMV 同比波动检测 sql: > SELECT abs(gmv_today - gmv_yesterday) / gmv_yesterday AS wow_change FROM ( SELECT SUM(pay_amount) AS gmv_today FROM dwd_trade_order_detail WHERE channel = 'pdd' AND dt = '昨日' ) today, ( SELECT SUM(pay_amount) AS gmv_yesterday FROM dwd_trade_order_detail WHERE channel = 'pdd' AND dt = date_sub(昨日, 1) ) yesterday threshold: 1.0 severity: infoHarness 校验:
• 调度时间 1:30→ 符合 ODS→DWD 必须在 2:00 前完成的约定 ✅• 质检规则覆盖 null 检查、行数对比、波动检测三种类型 → 符合 ETL 开发规范中的"每个 ETL 任务必须有质检"要求 ✅ • 没有跨层依赖(ODS→DWD→DWS,逐层)→ 符合分层约束 ✅
Step 6:输出接入完成报告
AI 最后汇总一份完整的接入完成报告,供架构师审批:
📋 拼多多旗舰店数仓接入完成报告═══════════════════════════════════════ODS 层: ✅ 建表 DDL → ods_trade_order_pdd_di ✅ Flink CDC 同步任务配置 ✅ 表名规范校验通过 ✅ 字段注释覆盖率 100% ✅ 分区策略 = toYYYYMM(dt)DWD 层: ✅ 清洗 SQL → ods → dwd_trade_order_detail ✅ 状态枚举映射完整(5 种状态全部覆盖) ✅ 渠道值 'pdd' 与域划分文档一致 ✅ 无 SELECT * 等违规写法DWS 层: ✅ 汇总 SQL → 按 brand_id/category_id/dt 3 个维度 ✅ gmv 口径与指标体系文档一致 ✅ JOIN dim_sku 路径校验通过调度与质检: ✅ 调度依赖配置(共 3 个任务) ✅ 质检规则(共 3 条) ✅ 无跨层依赖违规红线检查: ✅ 无 DROP TABLE ✅ 无 TRUNCATE ✅ 所有 INSERT 带分区条件 ✅ 无明文凭证(全部引用 ${变量})关联文档更新: 📝 域划分文档 → 渠道枚举新增 'pdd' 📝 DWD 状态枚举 → 已有值覆盖全部 5 种状态═══════════════════════════════════════状态:等待架构师审批。架构师需要做的,就是看一遍这份报告,确认签名。
一个完整的新渠道接入流程,从"人写 5-6 天"变成了"人写文档半天 + AI 自动完成 + 人看报告 1 小时"。
五、Harness 的兜底机制全览
整个过程中,Harness 在 4 个层面同时运作:
1. 约束层(CLAUDE.md + ARCHITECTURE.md)
AI 每次启动时自动读取。告诉 AI:这是什么项目、属于哪个域、必须遵守什么规则、不允许碰什么操作。
2. 工具层(settings.json 权限)
控制 AI 能调用什么工具。数仓场景下,Read 放开,Write/Edit/Bash 全开但受 Hook 约束。
3. 中间件层(Hook 拦截器)
三层拦截器:
PreToolUse: 拦截 DROP TABLE → exit(1) 阻断 拦截 TRUNCATE → exit(1) 阻断 检测全表 DELETE → exit(1) 阻断 检测 INSERT 无 PARTITION → 警告 + 自动补全PostToolUse: 扫描输出中的明文凭证 → 清除 + 审计日志 扫描 EXPLAIN 结果 → 超阈值标记 检测 SELECT * → 标记 + 建议修正Session Summary(可手动触发): 记录本次变更的文件列表 记录 Hook 拦截日志 留痕供事后审计4. 编排层(工作流程)
AI 按 6 步顺序执行。每一步的输出都是下一步的输入。如果某一步校验不通过,AI 修正后重试,不跳到下一步。
接入文档 → Step1(DDL) → Step2(同步) → Step3(DWD) → Step4(DWS) → Step5(调度质检) → Step6(报告) 校验 ✅ 校验 ✅ 校验 ✅ 校验 ✅ 校验 ✅ 审批签名这种"生成一小步 + 校验一步"的方式,把出大错的风险分散到了每一步。比一次全部生成再整体校验,更容易定位问题、也更安全。
六、接完不是结束:文档自动更新
接入完成后,AI 还会自动更新受影响的 Harness 文档:
• 域划分文档 → 渠道枚举新增 pdd,新增表列到订单域下• DWS 口径表 → 新渠道的口径与已有口径做一致性校验("拼多多的 gmv 和天猫的 gmv 口径一致吗?") • 质检规则库 → 新增的 3 条质检规则自动归档
这意味着第二次接入新渠道时,Harness 体系更完整了。
七、总结
回到开头的问题:数仓接入的效率瓶颈不是写代码的速度,是让机器放心干活的能力。
这套方案的核心不是 AI 有多聪明,是 Harness 体系让 AI 的聪明变得可靠:
• 接入文档定义"做什么"——架构师写一次 • AI 自动完成"怎么做"——从 DDL 到质检报告 • Harness 兜底"别做错"——红线 + 校验 + 拦截 • 文档自动更新——体系越用越完整
结果:
传统方式: 需求 → 开发写 5-6 天 → 人工验证 → 上线我们的方式: 架构师写文档 0.5 天 → AI 自动完成 → 审报告 1 小时 → 上线数仓接入,架构师写文档,AI 写代码。不是 AI 取代了架构师,是架构师从写重复代码升级到了写契约文档。
至于这套 Harness 体系是怎么从零建起来的、路上踩过哪些坑、每个阶段的效果如何——那是下一篇的话题。
夜雨聆风