dbt中文教程文档见下:
本教程面向制造业工厂数据治理场景,基于 dbt + DuckDB 构建从源系统到数仓分层(ODS/DWD/DWS/ADS)的配置与加工实践。
文档目录
- 01 dbt 工程配置说明
- 02 数据处理说明
- 03 BI 分析规划
01 dbt 工程配置说明
工程:
dbt_factory/(通用制造工厂数据中台)。数仓库data/dwh.duckdb,源库data/raw.duckdb(只读挂载)。
实测:dbt build全量 402/402 全部通过(详见docs/run_report.md)。
1. 环境准备
1.1 软件版本(实测)
| 组件 | 版本 | 说明 |
|---|---|---|
| Python | 3.14.4 | 虚拟环境位于仓库根目录 venv/ |
| dbt-core | 1.12.3 | 数据转换框架 |
| dbt-duckdb | 1.11.0 | DuckDB 适配器 |
| duckdb | 1.5.5 | 嵌入式分析数据库 |
| dbt-core-experimental-parser | 2.0.0rc1 | dbt-core 1.12 依赖的实验性解析器 |
| PyYAML | 6.0.3 | 自助建模生成器依赖 |
1.2 安装步骤
python3 -m venv venv
./venv/bin/pip install dbt-core==1.12.3 dbt-duckdb==1.11.0 duckdb==1.5.5 pyyaml
1.3 网络坑:dbt-core-experimental-parser 手动安装
dbt-core-experimental-parser 的 wheel 托管在 GitHub Release 上,pip 直连下载经常超时失败。
解决办法:经 gh-proxy 镜像手动下载 wheel 后本地安装:
# 1) 通过 gh-proxy 镜像下载 wheel(本工程实际使用 https://gh-proxy.com/ 前缀代理;
# 注意 wheel 文件名与 release tag 如下,平台为 manylinux_2_28_x86_64)
wget "https://gh-proxy.com/https://github.com/dbt-labs/dbt-core/releases/download/v2.0.0-rc.1/dbt_core_experimental_parser-2.0.0rc1-py3-none-manylinux_2_28_x86_64.whl"
# 2) 校验 sha256 后本地安装(期望值见 sdist 内 _dbt_sa_build/assets.json,
# 本工程实测匹配)
./venv/bin/pip install ./dbt_core_experimental_parser-2.0.0rc1-py3-none-manylinux_2_28_x86_64.whl
# 3) 再装 dbt-core / dbt-duckdb(此时依赖已满足,不再回源下载)
./venv/bin/pip install dbt-core==1.12.3 dbt-duckdb==1.11.0
验证:./venv/bin/pip list | grep dbt 应能看到上表中的全部组件。
1.4 初始化与自检
# 生成源数据(224 张源表,每表 1000 行初始数据)
./venv/bin/python simulator/simulator.py init
# dbt 连接自检
cd dbt_factory && ../venv/bin/dbt debug --profiles-dir .
# 期望输出: Connection test: [OK connection ok] All checks passed!
2. profiles.yml 逐项解释
文件:dbt_factory/profiles.yml
dbt_factory:
target: dev
outputs:
dev:
type: duckdb
path: ../data/dwh.duckdb
threads: 8
attach:
- path: ../data/raw.duckdb
alias: raw
read_only: true
| 配置项 | 值 | 说明 |
|---|---|---|
dbt_factory(顶层) | profile 名 | 与 dbt_project.yml 的 profile: 'dbt_factory' 对应 |
target: dev | 默认 target | 多环境(dev/prod)切换入口,本工程只有 dev |
type: duckdb | 适配器类型 | 使用 dbt-duckdb 适配器 |
path: ../data/dwh.duckdb | 数仓库文件 | 所有模型产物(ods/dwd/dws/ads 各 schema)写入此文件 |
threads: 8 | 并发线程数 | dbt DAG 中无依赖关系的模型最多 8 路并行构建 |
attach.path: ../data/raw.duckdb | 源库文件 | 模拟器写入的 224 张源表所在库 |
attach.alias: raw | 挂载别名 | 源表以 raw.main.<表名> 寻址,sources.yml 中 database: raw 与此对应 |
attach.read_only: true | 只读挂载 | dbt 侧只读源库,写权限只留给模拟器(DuckDB 单文件单写者约束,避免写冲突) |
配套宏 macros/generate_schema_name.sql 的作用:默认情况下 dbt 会把 +schema 拼成<target.schema>_<custom_schema>(如 main_ods),本工程覆盖该宏,直接使用 +schema
配置的名字(ods/dwd/dws/ads),使 dwh.duckdb 中出现干净的四个 schema:
{% macro generate_schema_name(custom_schema_name, node) -%}
{%- if custom_schema_name is none -%}
{{ target.schema }}
{%- else -%}
{{ custom_schema_name | trim }}
{%- endif -%}
{%- endmacro %}
3. dbt_project.yml 逐项解释
文件:dbt_factory/dbt_project.yml
name: 'dbt_factory'
version: '1.0.0'
profile: 'dbt_factory'
model-paths: ["models"]
macro-paths: ["macros"]
test-paths: ["tests"]
snapshot-paths: ["snapshots"]
analysis-paths: ["analyses"]
seed-paths: ["seeds"]
target-path: "target"
clean-targets:
- "target"
- "dbt_packages"
models:
dbt_factory:
ods:
+materialized: view
+schema: ods
dwd:
+materialized: table
+schema: dwd
dws:
+materialized: table
+schema: dws
ads:
+materialized: table
+schema: ads
snapshots:
dbt_factory:
+schema: dwd
| 配置段 | 说明 |
|---|---|
name / version / profile | 工程名;profile 名指向 profiles.yml 中的 dbt_factory |
*-paths | 各类资源目录:模型/宏/测试/快照/分析/种子数据 |
target-path / clean-targets | 编译与产物输出目录;dbt clean 清空的目录 |
models.dbt_factory.ods | ODS 层默认 view 物化,落到 ods schema |
models.dbt_factory.dwd | DWD 层默认 table 物化(单个模型可在 SQL 头部用 {{ config(materialized='incremental', ...) }} 覆盖为增量) |
models.dbt_factory.dws / ads | 汇总层/应用层默认 table 物化 |
snapshots.dbt_factory.+schema: dwd | 快照(SCD2)产物也放 dwd schema |
要点:分层物化策略是"层默认值 + 模型级覆盖"。8 个增量模型就是在 DWD 层 table 默认值上
单独声明 materialized='incremental' 实现的。
4. 目录结构与各层职责
dbt_factory/
├── dbt_project.yml # 工程配置(分层物化 + schema)
├── profiles.yml # 连接配置(dwh 库 + raw 只读挂载)
├── models/
│ ├── ods/ # sources.yml + ods.yml + 224 个视图(tools/gen_ods.py 生成,勿手改)
│ ├── dwd/ # 13 个明细模型 + dwd.yml(去重/清洗/标准化,含 8 个 incremental)
│ ├── dws/ # 8 个日汇总模型 + dws.yml(GROUP BY / PIVOT / UNPIVOT)
│ └── ads/ # 9 个应用模型 + ads.yml(BI 主题宽表,含自助建模产物 ads_my_energy)
├── snapshots/
│ └── eam_equipment_snapshot.yml # 设备台账 SCD2 快照(timestamp 策略)
├── macros/ # fmt_date / fmt_datetime / fmt_amount / cents_to_yuan /
│ # bool_label / reach_label / dedup_by / generate_schema_name
├── tests/
│ ├── generic/not_exceed_threshold.sql # 自定义 generic 测试
│ └── test_*.sql # 5 个 singular 测试
├── analyses/ seeds/ # 分析语句 / 种子数据(预留)
└── target/ # 编译产物、manifest.json、catalog.json、运行日志
各层职责(与 DESIGN.md 分层规范一致):
| 层 | 命名 | 物化 | 职责 |
|---|---|---|---|
| ODS | ods_<系统>_<表> | view | 与源表 1:1;增量/实时表追加 _ods_loaded_at 列;声明 source freshness |
| DWD | dwd_<域>_<实体> | table / incremental | 去重(QUALIFY/ROW_NUMBER)、空值(COALESCE)、日期/金额标准化(宏)、字段拆分、条件标签 |
| DWS | dws_<域>_<粒度> | table | 日/车间粒度 GROUP BY 汇总、PIVOT 行转列、UNPIVOT 列转行、UNNEST 拆多行 |
| ADS | ads_<主题> | table | 面向 BI 角色的主题宽表与指标表 |
5. 常用命令速查
所有命令在 dbt_factory/ 目录下执行(--profiles-dir . 指向工程内 profiles.yml):
# 连接与配置自检
../venv/bin/dbt debug --profiles-dir .
# 全量构建(run + test + snapshot,推荐日常入口)
../venv/bin/dbt build --profiles-dir .
# 只跑模型不跑测试 / 只跑测试
../venv/bin/dbt run --profiles-dir .
../venv/bin/dbt test --profiles-dir .
# 按层/按模型选择执行
../venv/bin/dbt run --select dwd --profiles-dir . # 只跑 DWD 层
../venv/bin/dbt run --select dwd_mes_prod_detail --profiles-dir .
../venv/bin/dbt build --select +ads_factory_cockpit_day --profiles-dir . # 含上游依赖
# 强制全量重建增量模型(忽略水位线)
../venv/bin/dbt run --select dwd_mes_prod_detail --full-refresh --profiles-dir .
# 源数据新鲜度检查(178 张配置 freshness 的源表)
../venv/bin/dbt source freshness --profiles-dir .
# 文档与血缘
../venv/bin/dbt docs generate --profiles-dir .
../venv/bin/dbt docs serve --profiles-dir . # 浏览器查看 DAG 血缘图
实测耗时参考(docs/run_report.md):首次全量 dbt build 65s(PASS=401),模拟器追加 60s
后第二次 61s(PASS=401),含自助建模产物的终验 build 402/402;source freshness 9.1s。
6. 增量任务类型说明
本工程覆盖 dbt 的四种典型任务形态,落点如下:
| 类型 | 说明 | 本工程例子 | 适用场景 |
|---|---|---|---|
view | 不落数据,查询时实时计算 | ODS 层全部 224 个模型 | 源表直通、轻量封装,要求"所见即源" |
table | 每次全量重建 | dwd_mes_forming_batch_all、dwd_eam_device、dwd_erp_sales_order_item、dwd_bip_approval_flow、dwd_sec_alarm;DWS/ADS 全层 | 数据量小、无稳定水位键、或逻辑含 UNNEST/多路 JOIN 不便增量的模型 |
incremental + append | 只追加水位线之后的新行 | dwd_acs_swipe_dedup、dwd_env_cems_hour、dwd_ems_meter_reading、dwd_log_weighbridge(ts 水位);dwd_erp_sales_order、dwd_ems_unit_consume(id 水位) | 纯流水表,历史行不会回改 |
incremental + delete+insert | 按 unique_key 先删后插,可覆盖回改 | dwd_mes_prod_detail(ts 水位)、dwd_mes_unit_daily(id 水位) | 业务键可能重复进入新批数据、需要幂等覆盖 |
snapshot(SCD2) | 记录维度历史版本(dbt_valid_from/to) | snapshots/eam_equipment_snapshot(timestamp 策略,updated_at: install_ts,unique_key=id) | 主数据/维度变更留痕,如设备台账、组织人员 |
水位线 > vs >= 的坑(实测踩过)
append策略必须用>(严格大于)。初版用ts >= max(ts)时,边界行会被重复 append,
导致 4 个模型的 unique 测试失败;改为ts >后正常(前提是模拟器保证 ts 严格单调递增)。delete+insert策略保留>=无碍——命中 unique_key 的历史行会先被删除再重新插入,
边界重复不会留痕。
对照示例:
-- append(dwd_acs_swipe_dedup.sql):严格大于
{% if is_incremental() %}
where ts > (select coalesce(max(swipe_time), timestamp '1900-01-01') from {{ this }})
{% endif %}
-- delete+insert(dwd_mes_unit_daily.sql):大于等于,靠先删后插保证幂等
{% if is_incremental() %}
where id >= (select coalesce(max(id), 0) from {{ this }})
{% endif %}
增量执行入口:dbt run 按水位线增量;加 --full-refresh 则忽略水位全量重建。
7. 宏使用说明
宏文件位于 dbt_factory/macros/,在模型中以 {{ 宏名(...) }} 调用:
| 宏 | 签名 | 作用 | 示例 |
|---|---|---|---|
fmt_date | fmt_date(col) | 三种字符串日期(YYYY-MM-DD/YYYY/MM/DD/YYYYMMDD)统一为 DATE,空值保持 NULL | {{ fmt_date('order_date') }} as order_date |
fmt_datetime | fmt_datetime(col) | 'YYYY-MM-DD HH:MM:SS' 字符串转 TIMESTAMP | {{ fmt_datetime('update_time') }} |
fmt_amount | fmt_amount(col_fen) | 分(BIGINT)→ 元,round 2,decimal(18,2) | {{ fmt_amount('total_fen') }} as total_yuan |
cents_to_yuan | cents_to_yuan(col_fen) | fmt_amount 的语义别名(供自助建模按名取用) | {{ cents_to_yuan('price_fen') }} |
bool_label | bool_label(cond, true_label, false_label) | 布尔条件 → 中文标签 | {{ bool_label('so2 > 35 or dust > 10 or voc > 100', '超标', '达标') }} as over_label |
reach_label | reach_label(value, limit) | 达标判定(值 ≤ 限值则"达标",NULL 则"无数据") | {{ reach_label('voc', 100) }} |
dedup_by | dedup_by(partition_cols, order_col) | 生成 QUALIFY row_number() 去重片段 | {{ dedup_by('unit_no', 'id desc') }} |
generate_schema_name | dbt 内建覆盖 | 见第 2 节,控制 schema 命名 | — |
fmt_date 的实现(正则识别三种格式,兜底 try_cast):
{% macro fmt_date(col) -%}
case
when {{ col }} is null then null
when regexp_matches({{ col }}, '^\d{4}-\d{2}-\d{2}$') then cast({{ col }} as date)
when regexp_matches({{ col }}, '^\d{4}/\d{2}/\d{2}$') then cast(strptime({{ col }}, '%Y/%m/%d') as date)
when regexp_matches({{ col }}, '^\d{8}$') then cast(strptime({{ col }}, '%Y%m%d') as date)
else try_cast({{ col }} as date)
end
{%- endmacro %}
8. 测试体系与 freshness
8.1 测试分类(合计 147 个数据测试,全部通过)
| 类别 | 位置 | 数量/例子 |
|---|---|---|
| dbt 内建 generic | 各层 yml 的 data_tests | not_null / unique / accepted_values / relationships(如订单行项目 → 订单的外键关系) |
自定义 generic not_exceed_threshold | tests/generic/not_exceed_threshold.sql | 如 dws_production_day.power_per_unit ≤ 1000、ads_factory_cockpit_day.emission_compliance_rate ≤ 1 |
| singular 测试 | tests/test_*.sql | 共 5 个:刷卡去重校验、金额非负、驾驶舱日期连续、单位产品电耗区间、排放超标率 <5% |
自定义 generic 测试的实现(查出超限行即失败):
{% test not_exceed_threshold(model, column_name, threshold) %}
select {{ column_name }}
from {{ model }}
where {{ column_name }} > {{ threshold }}
{% endtest %}
yml 中的写法(dbt 1.12 要求参数放 arguments: 下,否则有弃用警告):
data_tests:
- not_exceed_threshold:
arguments:
threshold: 1000
8.2 source freshness 配置
在 models/ods/sources.yml 中按同步方式分级配置,共 178 张源表:
- incremental 表:
warn_after: {count: 24, period: hour} - realtime 表:
warn_after: {count: 1, period: hour} loaded_at_field支持 SQL 表达式原样插入:"ts"、"cast(rpt_date as timestamp)"、"strptime(rpt_date, '%Y%m%d')"、"strptime(ym, '%Y%m')"等。
执行:../venv/bin/dbt source freshness --profiles-dir .,结果写入 target/sources.json。
实测:174 pass + 4 warn + 0 error(4 张 warn 均为月份粒度表,strptime(ym,'%Y%m') 落到月初
导致账龄 >24h,属预期告警,演示了 freshness 告警机制)。
9. 数据血缘
../venv/bin/dbt docs generate --profiles-dir .
../venv/bin/dbt docs serve --profiles-dir .
target/manifest.json— 全部节点(source/模型/测试/快照/宏)及依赖关系的完整清单;target/catalog.json— 各节点在数仓中的实际列与类型。
dbt docs serve 启动本地站点后可交互查看 source → ods → dwd → dws → ads 的完整 DAG 血缘图,
任意节点可下钻列级说明(来自各层 yml 的 description)。
10. 自助建模工具 self_service
业务同学用 YAML 描述模型,self_service/model_generator.py 渲染成标准 dbt SQL 落到dbt_factory/models/<layer>/ 下,可直接参与 dbt build,也可 --run 立即单跑。
# 只生成 SQL
./venv/bin/python self_service/model_generator.py self_service/examples/ads_my_energy.yaml
# 生成并立即运行(实测 PASS=1)
./venv/bin/python self_service/model_generator.py self_service/examples/ads_my_energy.yaml --run
YAML 字段说明:
| 字段 | 必填 | 说明 |
|---|---|---|
name | 是 | 模型名(小写蛇形),同时是输出文件名 |
layer | 是 | ods/dwd/dws/ads,决定输出目录;物化默认继承 dbt_project.yml 分层配置 |
description | 否 | 中文描述,写入生成 SQL 头部注释 |
materialized | 否 | 覆盖层默认物化(view/table) |
from | 是 | {ref: 模型名} 或 {source: [源名, 表名]},可加 alias |
columns | 否 | 选择列,支持 表达式 as 别名 |
joins | 否 | [{to: {ref: ...}, type: left, on: "...", alias: ...}] |
where | 否 | 过滤条件 |
group_by | 否 | 分组列 |
aggregates | 否 | 聚合表达式,如 sum(elec_kwh) as total_elec_kwh |
having / order_by | 否 | 聚合过滤 / 排序 |
示例(self_service/examples/ads_my_energy.yaml):
name: ads_my_energy
layer: ads
description: "自助建模示例: 各车间能耗合计(电/天然气/水), 按电量降序"
from:
ref: dws_energy_workshop_day
columns:
- workshop
group_by:
- workshop
aggregates:
- round(sum(elec_kwh), 2) as total_elec_kwh
- round(sum(gas_m3), 2) as total_gas_m3
- round(sum(water_t), 2) as total_water_t
- count(distinct stat_date) as stat_days
order_by: total_elec_kwh desc
生成器只依赖标准库 + pyyaml,不引入任何 dbt 包;生成文件头部带"请勿手改"注释,
终验 dbt build 402/402 已包含该自助生成模型。
02 数据处理说明
面向数据开发者:从 224 张源表到 4 层数仓的完整加工逻辑、16 类数据处理任务落点、
增量/实时同步链路与数据质量保障体系。实测数据见docs/run_report.md。
1. 源系统与数据模拟说明
1.1 源系统总览(11 个系统,224 张表)
由 simulator/(catalog.py 定义 + simulator.py 生成)写入 data/raw.duckdb(schema=main),
每表初始 1000 行。表清单与当前行数见 docs/source_catalog.md。
| 系统 | 前缀 | 表数 | 内容 |
|---|---|---|---|
| MES 制造执行 | mes_ | 45 | 原料/成型/加工/装配/包装 等生产数据 |
| ERP 企业资源 | erp_ | 30 | 采购/销售/库存/财务/主数据 |
| EMS 能源管理 | ems_ | 25 | 电/天然气/蒸汽/水/压缩空气 |
| 物流 WMS/TMS | log_ | 25 | 车辆/运单/称重/半成品转运 |
| BIP 协同办公 | bip_ | 12 | 组织/人员/审批流 |
| 安防 | sec_ | 15 | 视频/报警/巡更/消防 |
| 门禁 | acs_ | 10 | 门禁点/刷卡/权限 |
| 环保 EHS | env_ | 20 | CEMS/烟气/废水/固废/危废 |
| 设备 EAM | eam_ | 15 | 台账/点检/保养/维修 |
| 质检 LIMS | lims_ | 12 | 原辅料/半成品/成品理化 |
| 实时 SCADA | rt_ | 15 | 工艺参数、环保高频 |
1.2 同步方式(sync 字段)
| sync | 表数 | 含义 | 典型表 |
|---|---|---|---|
full | 46 | 主数据/维度类,init 一次性写入,run 模式不追加 | mes_unit_profile、erp_material、eam_equipment |
incremental | 158 | 业务流水类,run 模式每 | mes_prod_record、erp_sales_order、log_weigh |
realtime | 20 | 高频监测类,run 模式每 | 全部 rt_*、env_cems_*、acs_swipe_record、acs_face_record、log_veh_gps |
所有 incremental/realtime 表都有 ts(TIMESTAMP)事件时间列,追加时单调递增,供 dbt
incremental 模型做水位过滤;无 ts 的流水表(如 erp_sales_order、mes_unit_daily)用自增id 做水位。DuckDB 单文件单写者:模拟器是唯一写者,每批追加一个事务;dbt 侧经
profiles.yml read_only: true 只读挂载。
1.3 模拟器用法
# 初始化:建全部 224 张表,每表 1000 行(幂等,重复执行先 DROP 再建)
# 初始数据日期分布:full/incremental 表在最近 90 天内,realtime 表在最近 24 小时内
./venv/bin/python simulator/simulator.py init
# 持续追加:--duration 秒数后自动退出,0 = 无限运行
./venv/bin/python simulator/simulator.py run --duration 600
# 重新生成表目录文档 docs/source_catalog.md
./venv/bin/python simulator/simulator.py doc
1.4 故意预置的数据特征(供加工演示)
| 特征 | 预置方式 | 加工落点 |
|---|---|---|
| JOIN 关联 | mes_prod_record.unit_no ↔ mes_unit_profile.unit_no(有限基数编码池) | dwd_mes_prod_detail |
| UNION ALL | mes_forming_c1_batch / mes_forming_c2_batch 完全同构 | dwd_mes_forming_batch_all |
| 重复行 | acs_swipe_record 含约 3% 整行重复 | dwd_acs_swipe_dedup |
| 空值 | 各表 remark 列约 5% 为 NULL | DWD 通用 COALESCE |
| 日期格式混乱 | YYYY-MM-DD / YYYY/MM/DD / YYYYMMDD 三种字符串散落不同表 | fmt_date 宏 |
| 分单位金额 | erp_* 等表金额字段以"分"存 BIGINT(*_fen) | fmt_amount 宏 |
| 异常值 | 约 1% 行超标(env_cems_* 排放、ems_* 电耗、mes_voltage_log 工艺电压等) | SQL flag + dbt tests |
| 字段拆分 | 设备编码 EQ-YL-BS-001、车牌 蒙A12345、批次号 F1-20260906-A01 | split_part |
| 嵌套/复合字段 | erp_sales_order.item_list(JSON 数组)、bip_approval.approver_list(逗号分隔) | UNNEST |
| 宽表/长表 | ems_meter_reading 含 energy_type 列;env_cems_hour 含 so2/nox/dust/voc 并列列 | PIVOT / UNPIVOT |
2. 四层加工逻辑
整体链路:raw.duckdb(224 源表)→ ODS(224 view)→ DWD(13 模型 + 1 snapshot)→
DWS(8 table)→ ADS(9 table)。
2.1 ODS 层(224 个 view,与源表 1:1)
- 由
tools/gen_ods.py以simulator/catalog.py为唯一契约批量生成,勿手改; - 每个视图直通源表,增量/实时表追加
current_timestamp as _ods_loaded_at装载时间列:
-- models/ods/ods_acs_swipe_record.sql
select
*,
current_timestamp as _ods_loaded_at
from {{ source('acs', 'acs_swipe_record') }}
sources.yml声明 178 张增量/实时表的 freshness(见第 5 节)。
2.2 DWD 层(13 模型 + 1 snapshot)
| 模型 | 中文名 | 物化 | 增量策略 | 说明 |
|---|---|---|---|---|
| dwd_mes_prod_detail | 生产明细 | incremental | delete+insert(unique_key=prod_id,ts 水位) | 生产记录 JOIN 单元档案/班组/生产计划,计划完成率公式列与单元状态标签 |
| dwd_mes_unit_daily | 生产单元日报清洗 | incremental | delete+insert(unique_key=id) | 日期统一,产出效率/电耗标签 |
| dwd_acs_swipe_dedup | 门禁刷卡去重 | incremental | append(ts 水位) | QUALIFY row_number 去除源表约 3% 整行重复 |
| dwd_env_cems_hour | 废气CEMS小时均值清洗 | incremental | append(ts 水位) | 超标标签+分污染物异常 flag(SO2>35/粉尘>10/VOCs>100) |
| dwd_ems_meter_reading | 计量点读数清洗 | incremental | append(ts 水位) | JOIN 计量点档案(先去重),供 DWS PIVOT |
| dwd_log_weighbridge | 地磅称重清洗 | incremental | append(ts 水位) | 毛重-皮重核算净重公式列,净重差异>0.5t 标异常 |
| dwd_erp_sales_order | 销售订单 | incremental | append(id 水位) | 金额分→元(fmt_amount),日期统一,状态标签 |
| dwd_ems_unit_consume | 产品单耗清洗 | incremental | append(id 水位) | 单位产品电耗异常值 flag(>550 或 <400) |
| dwd_mes_forming_batch_all | 成型批次合并 | table | 全量 | 成型一+成型二 UNION ALL,批次号拆出车间码/日期/序号 |
| dwd_eam_device | 设备台账 | table | 全量 | 设备编码拆分 系统/工序/序号,COALESCE 空值处理 |
| dwd_erp_sales_order_item | 销售订单行项目 | table | 全量 | item_list JSON 数组 UNNEST 拆多行,单价分→元 |
| dwd_bip_approval_flow | 审批流明细 | table | 全量 | approver_list 逗号分隔拆多行,一行一审批人 |
| dwd_sec_alarm | 安防报警清洗 | table | 全量 | 报警级别颜色标签,处置闭环标签 |
| eam_equipment_snapshot | 设备台账 SCD2 快照 | snapshot | timestamp 策略 | unique_key=id,updated_at=install_ts;变更产生 dbt_valid_from/to 版本行 |
2.3 DWS 层(8 个 table,日/车间粒度汇总)
| 模型 | 中文名 | 说明 |
|---|---|---|
| dws_production_day | 生产日汇总(日×产线) | 产量/产出效率/单位产品电耗(公式列:电量÷产量,产量加权) |
| dws_energy_workshop_day | 各车间能耗日汇总 | energy_type 行转列(PIVOT),电/天然气/蒸汽/水/压缩空气 |
| dws_env_pollutant_day | 排放污染物日汇总 | CEMS 4 污染物列 UNPIVOT 列转行后按 日×排口×污染物 汇总 |
| dws_device_health_day | 设备健康日汇总 | 故障次数/停机时长/点检异常数(三流 FULL OUTER JOIN) |
| dws_logistics_day | 物流日汇总 | 地磅净重/运单量/半成品转运量(三流 FULL OUTER JOIN) |
| dws_access_security_day | 门禁安防日汇总 | 刷卡次数/人数/报警数/重大报警数 |
| dws_forming_day | 成型生产日汇总(日×车间) | 批次数/产量/平均产品纯度 |
| dws_sales_day | 销售日汇总 | 订单数/订单额(元)/发货数/发货量(t),订单与发货 FULL OUTER JOIN |
2.4 ADS 层(9 个 table,BI 主题)
| 模型 | 中文名 | 说明 |
|---|---|---|
| ads_factory_cockpit_day | 厂长驾驶舱日报(大宽表) | 产量/单位产品电耗/产出效率/排放达标率/安全事件数/订单额 |
| ads_production_unit_ranking | 单元级排名与异常单元清单 | 产量/电耗 rank,高耗或低产出效率标"异常关注" |
| ads_energy_cost_day | 能源成本日报(日×车间) | 能耗 × 能源均价估算成本(元) |
| ads_env_compliance_day | 环保达标日报(日×排口) | 超标次数/达标率/达标标签 |
| ads_equipment_overview | 设备总览(车间×系统) | 设备数/运行数/故障数/停机时长/运行占比 |
| ads_logistics_efficiency | 物流效率日报 | 称重/运单/转运/进出厂车辆,平均单车净重与单运单吨位 |
| ads_safety_access_day | 安全与门禁日报 | 刷卡/报警/事故/隐患全量指标 |
| ads_sales_delivery | 销售发运明细(订单粒度) | 订单额/订购量/发运量/发运进度 |
| ads_my_energy | 各车间能耗合计(自助建模产物) | self_service 生成器产出,证明自助建模链路可用 |
3. 16 类数据处理任务逐一说明
3.1 左右合并 JOIN
- 实现方式:
left join关联维表;维表编码有限基数池存在重复档案行,JOIN 前先用dedup_by宏去重取最新。 - 落点:
dwd_mes_prod_detail(生产记录 JOIN 单元档案/班组/生产计划)。
-- dwd_mes_prod_detail.sql
unit as (
select unit_no, line_no, workshop, status
from {{ ref('ods_mes_unit_profile') }}
{{ dedup_by('unit_no', 'id desc') }}
)
...
from prod t
left join unit c on t.unit_no = c.unit_no
left join team tm on tm.workshop = c.workshop
left join plan p on p.unit_no = t.unit_no and p.plan_date_d = cast(t.ts as date) and p.shift = t.shift
3.2 上下合并 UNION ALL
- 实现方式:两张同构表(成型一/成型二车间批次)UNION ALL 纵向合并,各自补车间标签列。
- 落点:
dwd_mes_forming_batch_all。
-- dwd_mes_forming_batch_all.sql
select
'成型一车间' as workshop,
id, ts as prod_time, batch_no, ...
from {{ ref('ods_mes_forming_c1_batch') }}
union all
select
'成型二车间' as workshop,
id, ts as prod_time, batch_no, ...
from {{ ref('ods_mes_forming_c2_batch') }}
3.3 去重 QUALIFY / ROW_NUMBER
- 实现方式:
qualify row_number() over (partition by 业务键 order by ...) = 1,封装为dedup_by宏。 - 落点:
dwd_acs_swipe_dedup(源表约 3% 整行重复);维表去重还见于 dwd_mes_prod_detail、dwd_ems_meter_reading。
-- dwd_acs_swipe_dedup.sql
select
id, ts as swipe_time, cast(ts as date) as swipe_date,
card_no, emp_no, person, door_code, direction,
coalesce(remark, '无') as remark
from {{ ref('ods_acs_swipe_record') }}
{% if is_incremental() %}
where ts > (select coalesce(max(swipe_time), timestamp '1900-01-01') from {{ this }})
{% endif %}
{{ dedup_by('id, card_no, emp_no, door_code, direction, ts', 'ts') }}
3.4 空值处理 COALESCE
- 实现方式:备注、车间、计量点等可空字段统一兜底;
nullif防除零。 - 落点:DWD 全层通用,如 dwd_mes_prod_detail:
-- dwd_mes_prod_detail.sql
coalesce(c.workshop, '未知车间') as workshop,
coalesce(t.remark, '无') as remark,
round(t.output_qty / nullif(p.plan_qty, 0) * 100, 1) as plan_fulfill_pct,
3.5 日期格式统一(宏)
- 实现方式:
fmt_date宏正则识别YYYY-MM-DD/YYYY/MM/DD/YYYYMMDD三种字符串并转 DATE,空值保持 NULL。 - 落点:dwd_erp_sales_order、dwd_mes_unit_daily、dwd_eam_device、dwd_mes_forming_batch_all(批次号内嵌 YYYYMMDD)等。
-- dwd_erp_sales_order.sql
{{ fmt_date('order_date') }} as order_date,
-- dwd_mes_forming_batch_all.sql:批次号第 2 段是 YYYYMMDD 日期
{{ fmt_date("split_part(batch_no, '-', 2)") }} as batch_date,
3.6 金额格式统一(宏)
- 实现方式:源表金额以"分"存 BIGINT(
*_fen),fmt_amount宏转元:cast(round(x/100.0, 2) as decimal(18,2))。 - 落点:dwd_erp_sales_order、dwd_erp_sales_order_item、ads_energy_cost_day。
-- dwd_erp_sales_order.sql
{{ fmt_amount('total_fen') }} as total_yuan,
-- dwd_erp_sales_order_item.sql
cast(round(j.qty * j.price_fen / 100.0, 2) as decimal(18,2)) as item_amount_yuan
3.7 异常值识别(Tests + SQL 双轨)
- SQL 轨:在 DWD 打异常 flag 列。落点 dwd_ems_unit_consume(电耗)、dwd_env_cems_hour(排放)、dwd_log_weighbridge(净重):
-- dwd_ems_unit_consume.sql(正常区间 450~550 kWh/t,约 1% 异常高耗 600~800)
{{ bool_label('unit_consume > 550 or unit_consume < 400', '异常', '正常') }} as power_flag,
- Tests 轨:singular 测试做区间/比率校验,落点
tests/test_unit_power_range.sql(单位产品电耗 400~600,
加unit_cnt >= 5样本量门槛,排除小样本边缘日的统计噪声)、tests/test_env_exceed_rate.sql(排放超标率 <5%)。
3.8 行转列 PIVOT
- 实现方式:DuckDB 原生
pivot ... on ... using sum(...)。 - 落点:
dws_energy_workshop_day(energy_type 行转列为五种能源列)。
-- dws_energy_workshop_day.sql
piv as (
pivot src
on energy_type
using sum(reading)
group by stat_date, workshop
)
select
stat_date, workshop,
round("电", 2) as elec_kwh, round("天然气", 2) as gas_m3, ...
from piv
3.9 列转行 UNPIVOT
- 实现方式:DuckDB 原生
unpivot ... on ... into name ... value ...,4 个并列污染物列转成长表再汇总。 - 落点:
dws_env_pollutant_day。
-- dws_env_pollutant_day.sql
with long as (
unpivot (
select stat_date, outlet_code, so2, nox, dust, voc
from {{ ref('dwd_env_cems_hour') }}
)
on so2, nox, dust, voc
into name pollutant value conc
)
select stat_date, outlet_code, pollutant,
count(*) as sample_cnt, round(avg(conc), 2) as conc_avg, ...
from long group by 1, 2, 3
3.10 字段拆分(split_part)
- 实现方式:
split_part(col, '-', n)拆结构化编码 + CASE 映射中文名。 - 落点:
dwd_eam_device(设备编码EQ-YL-BS-001拆出系统/工序/序号)、dwd_mes_forming_batch_all(批次号拆车间码/日期/序号)。
-- dwd_eam_device.sql
split_part(equip_code, '-', 2) as sys_code,
case split_part(equip_code, '-', 2)
when 'YL' then '原料系统' when 'SC' then '生产系统' when 'CX' then '成型系统'
when 'HB' then '环保系统' when 'DL' then '动力系统' when 'WX' then '维修系统'
else '其他'
end as sys_name,
3.11 拆分成多行 UNNEST
- 实现方式:JSON 数组用
unnest(from_json(...));逗号分隔串用unnest(string_split(...))。 - 落点:
dwd_erp_sales_order_item(订单行项目)、dwd_bip_approval_flow(审批人列表)。
-- dwd_erp_sales_order_item.sql
from {{ ref('ods_erp_sales_order') }} o,
unnest(from_json(o.item_list,
'[{"material":"VARCHAR","qty":"DOUBLE","price_fen":"BIGINT"}]')) as u(j)
-- dwd_bip_approval_flow.sql
from {{ ref('ods_bip_approval') }} a,
unnest(string_split(a.approver_list, ',')) as u(approver)
3.12 分组汇总 GROUP BY
- 实现方式:DWS 全层按 日/产线/车间 分组聚合;多路指标流用 FULL OUTER JOIN 按日合并 + COALESCE 补零。
- 落点:
dws_production_day、dws_logistics_day等全部 8 个 DWS 模型。
-- dws_production_day.sql
select
rpt_date as stat_date,
cast(substr(unit_no, 1, 1) as integer) as line_no,
count(distinct unit_no) as unit_cnt,
round(sum(output_qty) / 1000.0, 2) as output_t,
...
from {{ ref('dwd_mes_unit_daily') }}
group by 1, 2
3.13 公式列
- 实现方式:汇总层按业务公式派生指标,
nullif防除零。 - 落点:
dws_production_day(单位产品电耗)、dwd_log_weighbridge(净重核算)、ads_logistics_efficiency(平均单车净重)。
-- dws_production_day.sql:单位产品电耗 = Σ(单耗×产量)/Σ产量(产量加权)
round(sum(power_consume * output_qty) / nullif(sum(output_qty), 0), 1) as power_per_unit,
-- dwd_log_weighbridge.sql:毛重-皮重核算净重,与仪表净重比对
round(gross_t - tare_t, 2) as calc_net_t,
round(net_t - (gross_t - tare_t), 2) as net_diff_t,
3.14 条件标签列 CASE WHEN
- 实现方式:多分支 CASE 或
bool_label/reach_label宏生成中文标签。 - 落点:
dwd_mes_prod_detail(单元状态)、dwd_env_cems_hour(达标/超标)、dwd_sec_alarm(级别颜色)、ads_production_unit_ranking(异常关注)。
-- dwd_mes_prod_detail.sql
case c.status
when '运行' then '正常生产'
when '停机' then '停产'
when '检修' then '检修中'
else '启动调试'
end as unit_status_label,
-- dwd_env_cems_hour.sql
{{ bool_label('so2 > 35 or dust > 10 or voc > 100', '超标', '达标') }} as over_label,
3.15 自助建模(自研 YAML 生成器)
- 实现方式:业务同学写 YAML(from/columns/group_by/aggregates/order_by),
self_service/model_generator.py渲染成标准 dbt SQL 落到对应层目录。 - 落点:
ads_my_energy(实测生成并dbt run成功,PASS=1)。字段说明见docs/01_dbt配置说明.md第 10 节。
3.16 数据质量检查与血缘
- 质量检查:dbt tests 体系,合计 147 个数据测试(内建 generic + 自定义 generic
not_exceed_threshold+ 5 个 singular),
随dbt build同步执行,实测全部通过。 - 血缘:
dbt docs generate产出target/manifest.json(节点与依赖)与target/catalog.json(列级目录),dbt docs serve交互查看 source → ods → dwd → dws → ads 全链路 DAG。
4. 增量与实时同步链路
模拟器 run(--duration) ──追加──> raw.duckdb ──只读挂载──> dbt incremental(水位线) ──> DWD ──全量重建──> DWS/ADS
incremental 表 ~15s 一批 5~20 行 ts 严格单调递增
realtime 表 ~2s 一批 10~50 行 append 用 > 水位, delete+insert 用 unique_key 幂等
实测验证(docs/run_report.md 第 3 节):模拟器追加 60s 后再次 dbt build(PASS=401,61s),
增量模型前后行数对比:
| 模型(策略) | 源表 前→后 | DWD 前→后 | 说明 |
|---|---|---|---|
| dwd_mes_prod_detail(delete+insert) | 1041 → 1081(+40) | 1041 → 1081(+40) | 正确追加,无重复 |
| dwd_acs_swipe_dedup(append) | 1732 → 2354(+622) | 1674 → 2275(+601) | 新增批次内 21 行重复被 QUALIFY 去除 |
| dwd_mes_unit_daily(delete+insert) | 1039 → 1068(+29) | 1039 → 1068(+29) | 正确追加 |
水位线边界教训:append 模型初版用 ts >= max(ts) 导致边界行重复 append、4 个模型 unique 测试失败,
改为 ts > max(ts) 后正常;delete+insert 模型因先删后插,保留 >= 无碍。DWS/ADS 为全量 table,
每次 build 基于最新 DWD 全量重建,天然拿到最新汇总。
5. 数据质量保障体系
5.1 测试分类统计(147 个数据测试,全部通过)
| 类别 | 内容 | 例子 |
|---|---|---|
| 内建 generic | not_null / unique / accepted_values / relationships | 主键非空唯一;枚举值域(单元状态、达标标签等);订单行项目 → 订单外键 |
| 自定义 generic | not_exceed_threshold(列值不得超阈值) | dws_production_day.power_per_unit ≤ 1000;ads_factory_cockpit_day.emission_compliance_rate ≤ 1 |
| singular(5 个) | 自定义 SQL,返回非空即失败 | 刷卡去重校验、金额非负、驾驶舱日期连续(generate_series 补全日历比对)、单位产品电耗区间(400~600,unit_cnt≥5)、排放超标率 <5% |
5.2 freshness 监控
- 178 张源表配置 freshness:incremental 表
warn_after 24h,realtime 表warn_after 1h; loaded_at_field支持表达式(cast(... as timestamp)/strptime(...)),适配无 ts 列的表;- 实测 9.1s:174 pass + 4 warn + 0 error。4 张 warn(ems_energy_plan / env_emission_fee /
env_env_report / mes_target_plan)为月份粒度表,strptime(ym,'%Y%m')落月初账龄 >24h,属预期。
5.3 异常值阈值(DESIGN.md 达标线)
| 指标 | 阈值 | 处理 |
|---|---|---|
| VOCs | > 100 mg/m³ 超标 | dwd_env_cems_hour flag + dws_env_pollutant_day 超标计数 |
| 粉尘 | > 10 mg/m³ 超标 | 同上 |
| SO2 | > 35 mg/m³ 超标 | 同上(NOx 参考 >100) |
| 单位产品电耗 | <400 或 >550 kWh/t 异常 | dwd_ems_unit_consume power_flag;日汇总区间测试 400~600 |
| 产出效率 | < 94% 偏低 | dwd_mes_unit_daily eff_label |
| 地磅净重 | 毛重-皮重与仪表净重差 >0.5t 异常 | dwd_log_weighbridge weigh_flag |
6. 已知数据特征与注意事项
沿用 docs/run_report.md 的"已知数据特征"(非缺陷):
mes_unit_profile.unit_no、ems_meter.meter_code等编码列为有限基数随机池,源表内不唯一,
DWD JOIN 前均先 QUALIFY 去重取最新;ODS 层相应不配 unique 测试。acs_swipe_record源表故意含 3% 整行重复,ODS 不配 id unique,由dwd_acs_swipe_dedup去重后测试。- 模拟器中发货单
order_no独立生成,ads_sales_delivery仅少量巧合匹配,delivery_pct 多数为 0/NULL。 - 排放达标率仅近 1~2 天有值:
env_cems_hour为 realtime 表,初始数据仅最近 24 小时。 - 数仓核验参考:dwh.duckdb 中 ods=224 view、dwd=14 表(13 模型+1 snapshot)、dws=8 表、ads=9 表;
eam_equipment_snapshot 1000 行/1000 唯一 id/0 条已闭版本(源为静态主数据,SCD2 结构就绪);
ads_factory_cockpit_day 91 行(2026-06-08 ~ 2026-09-06,日期连续 singular test 通过)。
03 BI 分析规划
面向 BI 团队与管理层:基于 dbt_factory 数仓 ADS 层 9 张主题表,规划全厂角色驾驶舱/看板体系。
数仓库:data/dwh.duckdb(ads schema)。所有 ADS 模型均已在dbt build中验证通过。
1. 工厂工艺背景
本厂为通用制造工厂,2 条产线、每条产线 160 个生产单元(单元号 11011160 /2260)。工艺流程覆盖多个车间/产线,主要物料流向如下:
2101
原料仓 ──> 成型/加工车间 ──> 装配/包装车间 ──> 半成品/成品 ──> 质检
│ │
└─ 余料/不合格品回收处理 └─> 成品仓 ──> 销售发运
工艺废气/废水/固废 ──> 环保处理系统(净化 + 在线监测)──> 达标排放/合规处置
配套支撑:动力车间(变配电/空压/水泵)、维修车间(设备维修与保养)、物流(原料到货、
半成品转运、成品发运)、环保(CEMS/废水/固废/危废)、安防门禁与质检(LIMS)。
行业标杆参考值(DESIGN.md):单位产品电耗 ~480 kWh/t、综合单位电耗 ~520 kWh/t、
原料/辅料单耗按产品定额管理、产出效率 ~94%;
排放达标线:VOCs < 100 mg/m³、粉尘 < 10 mg/m³、SO2 < 35 mg/m³。
2. 全厂角色 BI 体系规划
刷新频率约定:实时链路(rt_*、env_cems_* 等 realtime 源表)看板按 1 分钟刷新;
日级 ADS 表按 T+1(随每日 dbt build 更新)。
| 角色 | 驾驶舱/看板 | 核心指标 | 对应 ADS 模型 | 建议图表 | 刷新 |
|---|---|---|---|---|---|
| 厂长 | 全厂驾驶舱日报 | 成品产量、单位产品电耗、产出效率、排放达标率、安全事件数、未处置报警、订单数/订单额 | ads_factory_cockpit_day | KPI 卡片 + 产量/电耗双轴趋势线 + 达标率仪表盘 | T+1 |
| 生产副厂长 | 生产运行看板 | 日×产线产量/产出效率/高耗单元数、单元级产量与电耗排名、异常关注单元清单 | ads_factory_cockpit_day、ads_production_unit_ranking | 产线对比柱状图 + 单元号排名条形图 + 异常单元明细表 | T+1 |
| 生产车间主任 | 单元状况监控看板 | 工艺电压/工艺温度/产线电流实时值、单单元日电耗与产出效率、异常关注单元 | ads_production_unit_ranking(T+1);rt_unit_voltage/rt_unit_temp/rt_line_current(实时) | 实时折线 + 单元热力图 + 异常单元列表 | 实时 1 分钟 + T+1 |
| 成型车间主任 | 成型生产看板 | 日×车间批次数/产量/平均产品纯度、成型一 vs 二对比 | dws_forming_day(经 ADS 或直接消费 DWS) | 车间对比柱状图 + 产量趋势线 | T+1 |
| 设备部 | 设备健康看板 | 车间×系统设备数/运行占比/故障数/停机时长、日故障与点检异常趋势 | ads_equipment_overview、dws_device_health_day | 运行占比环形图 + 停机时长趋势 + 故障 TOP 系统 | T+1 |
| 能源部 | 能耗成本看板 | 日×车间电/气/汽/水消耗与成本、单位产品电耗及标杆对比 | ads_energy_cost_day、ads_factory_cockpit_day | 能耗堆叠柱状图 + 成本占比饼图 + 电耗 vs 标杆线 | T+1 |
| 安环部 | 环保达标看板 | 日×排口超标次数/达标率、分污染物浓度均值与峰值、安全事件/事故/隐患未整改数 | ads_env_compliance_day、ads_safety_access_day;env_cems_min(实时) | 排口达标率热力表 + 污染物浓度趋势 + 实时 CEMS 曲线 | 实时 1 分钟 + T+1 |
| 物流部 | 物流效率看板 | 地磅净重/称重异常数、运单量、半成品转运量、进出厂车辆数、平均单车净重 | ads_logistics_efficiency | 进出厂流量对比 + 转运量趋势 + 异常称重明细 | T+1 |
| 销售经营 | 销售发运看板 | 订单数/订单额、发运量、订单粒度发运进度 | ads_sales_delivery、dws_sales_day | 订单额趋势 + 客户订单明细表 + 发运进度条 | T+1 |
| 财务 | 经营成本看板 | 订单额、各车间能源成本(元)、单位产品电耗 | ads_factory_cockpit_day、ads_energy_cost_day | 金额 KPI 卡片 + 成本趋势 | T+1 |
3. 指标体系
3.1 产量类
| 指标 | 口径 | 来源模型 |
|---|---|---|
| 成品日产量(t) | dws_production_day.output_t 全厂合计 | ads_factory_cockpit_day |
| 成品日产量 | dws_forming_day.output_t(日×车间) | dws_forming_day |
| 周/月产量 | 日表按周/月再聚合(BI 侧完成) | 上述日表 |
3.2 质量类
| 指标 | 口径 | 来源模型 |
|---|---|---|
| 产品平均纯度(目标 优级品) | dws_forming_day.purity_avg | dws_forming_day |
| 产出效率(标杆 ~94%) | 日×产线均值、全厂均值 | dws_production_day、ads_factory_cockpit_day |
| 当前实测参考 | 2026-09-06:产出效率 93.88% | ads_factory_cockpit_day |
3.3 能耗类(含行业标杆对比)
| 指标 | 标杆值 | 来源模型 |
|---|---|---|
| 单位产品电耗(kWh/t) | ~480(实测样例 486.6) | dws_production_day、ads_factory_cockpit_day |
| 综合单位电耗(kWh/t) | ~520 | EMS 源表扩展(见路线图) |
| 天然气/蒸汽/水单耗 | 能耗定额表 mes_energy_quota 对比 | dws_energy_workshop_day |
| 车间能耗成本(元) | — | ads_energy_cost_day |
3.4 设备类
| 指标 | 口径 | 来源模型 |
|---|---|---|
| 设备运行占比 | 运行台数/总台数(车间×系统) | ads_equipment_overview.running_ratio |
| 故障次数/停机时长 | 按日统计 | dws_device_health_day |
| 点检异常数 | 点检记录 result='异常' 计数 | dws_device_health_day |
| 维修及时率、OEE | 需扩展 eam_repair_order/eam_downtime 建模(见路线图) | 待建 |
3.5 环保类
| 指标 | 达标线 | 来源模型 |
|---|---|---|
| VOCs/粉尘/SO2 超标次数与达标率 | 100 / 10 / 35 mg/m³ | dws_env_pollutant_day、ads_env_compliance_day |
| 全厂排放达标率 | 实测样例 99.6% | ads_factory_cockpit_day |
| 危废合规 | env_hazard_waste 台账(待建模入仓) | 见路线图 |
3.6 物流类
| 指标 | 来源模型 |
|---|---|
| 半成品转运量/次数(生产单元→成型车间) | dws_logistics_day、ads_logistics_efficiency |
| 平均单车净重、单运单吨位、称重异常数 | ads_logistics_efficiency |
| 进出厂车辆数(车辆周转) | ads_logistics_efficiency(gate_in/out_cnt) |
3.7 安全类
| 指标 | 来源模型 |
|---|---|
| 刷卡次数/人数、门禁异常 | dws_access_security_day、ads_safety_access_day |
| 报警数/重大报警数/未处置报警数(报警处置率) | ads_safety_access_day |
| 事故事件数/伤亡数、隐患数/未整改数 | ads_safety_access_day |
3.8 经营类
| 指标 | 来源模型 |
|---|---|
| 订单数/订单额(实测样例 74 笔 / 1.83 亿元) | ads_factory_cockpit_day、dws_sales_day |
| 订单发运进度(订单交付率) | ads_sales_delivery.delivery_pct |
| 库存周转 | erp_inventory/erp_stock_in/out 待建模(见路线图) |
4. ADS 层模型与 BI 主题映射表
| ADS 模型 | 中文名 | 粒度 | BI 主题 | 主要消费角色 |
|---|---|---|---|---|
| ads_factory_cockpit_day | 厂长驾驶舱日报 | 日(91 行,2026-06-08~2026-09-06 连续) | 全厂综合 | 厂长、生产副厂长、财务 |
| ads_production_unit_ranking | 单元级排名与异常单元清单 | 生产单元(218 个在产单元) | 生产运行 | 生产副厂长、生产车间主任 |
| ads_energy_cost_day | 能源成本日报 | 日×车间 | 能耗成本 | 能源部、财务 |
| ads_env_compliance_day | 环保达标日报 | 日×排口 | 环保达标 | 安环部 |
| ads_equipment_overview | 设备总览 | 车间×系统 | 设备健康 | 设备部 |
| ads_logistics_efficiency | 物流效率日报 | 日 | 物流效率 | 物流部 |
| ads_safety_access_day | 安全与门禁日报 | 日 | 安全门禁 | 安环部、综合部 |
| ads_sales_delivery | 销售发运明细 | 订单 | 销售经营 | 销售经营、财务 |
| ads_my_energy | 各车间能耗合计(自助建模示例) | 车间 | 自助分析示例 | 业务科室 |
DWS 层 8 张日汇总表(dws_production_day、dws_forming_day、dws_energy_workshop_day、
dws_env_pollutant_day、dws_device_health_day、dws_logistics_day、dws_access_security_day、
dws_sales_day)可作为二级明细看板的直接数据源。
5. 落地路线图
阶段一:现有 ADS 对接 BI 工具(0~1 个月)
- BI 工具直连 dwh.duckdb 的 ads/dws schema,落地第 2 节 10 个角色看板;
- 每日定时
dbt build更新 T+1 数据;dbt test+source freshness作为数据质量门禁,
测试失败时看板挂"数据待核"标记; - 交付物:厂长驾驶舱 + 各部门日报看板。
阶段二:实时大屏接 rt_/env_ 链路(1~3 个月)
- 20 张 realtime 源表(rt_* 15 张 + env_cems_*、acs_swipe_record、acs_face_record、log_veh_gps)
经 ODS 视图直通 BI,按 1 分钟刷新; - 建设生产车间实时大屏(工艺电压/工艺温度/产线电流)、环保车间 CEMS 实时排放大屏、厂区门禁热力;
- 按需新增分钟级增量聚合模型(incremental append,复用现有水位线模式)。
阶段三:自助建模推广到业务科室(3 个月起)
- 以 self_service 生成器 + ads_my_energy 示例为模板,培训各科室用 YAML 自助建模型;
- 建立"YAML 评审 → 生成 → 纳入每日 build"的轻量流程,生成的模型自动继承分层物化与测试体系;
- 补齐经营域建模:库存周转(erp_inventory/stock_in/out)、OEE 与维修及时率(eam_repair_order/
eam_downtime)、危废合规台账(env_hazard_waste)等当前未入 ADS 的主题。
6. BI 工具对接建议
DuckDB 为单文件嵌入式库,三种对接方式按工具能力选择:
| 方式 | 适用工具 | 说明 |
|---|---|---|
| duckdb 驱动直连 | Superset(duckdb-engine)、Python/R 自定义报表 | 直接打开 data/dwh.duckdb,查 ads/dws schema;注意 DuckDB 单写者约束,BI 侧只读打开 |
| 定时导出 parquet | PowerBI、帆软 FineReport/FineBI 等无 duckdb 驱动的工具 | 每日 build 后用 COPY (SELECT * FROM ads.xxx) TO 'xxx.parquet' 导出,BI 消费 parquet/再入库 |
| 导出 CSV/Excel | 临时取数、科室自查 | 小表直接导出,或经自助建模生成器定制模型后导出 |
注意事项:
dbt build执行期间 dwh.duckdb 被 dbt 独写,BI 刷新窗口应避开 build 窗口(实测全量 build 约 61~65s,凌晨窗口充足);- 血缘与口径文档直接复用
dbt docs serve站点(manifest.json/catalog.json 已含全部中文列描述),
作为 BI 指标口径的唯一权威来源; - 实时链路消费 raw.duckdb 的 rt_*/env_cems_* 表时同样只读打开,写入权只归模拟器/采集程序。