阅山

  • WIN
    • CSharp
    • JAVA
    • OAM
    • DirectX
    • Emgucv
  • UNIX
    • FFmpeg
    • QT
    • Python
    • Opencv
    • Openwrt
    • Twisted
    • Design Patterns
    • Mysql
    • Mycat
    • MariaDB
    • Make
    • OAM
    • Supervisor
    • Nginx
    • KVM
    • Docker
    • OpenStack
  • WEB
    • ASP
    • Node.js
    • PHP
    • Directadmin
    • Openssl
    • Regex
  • APP
    • Android
  • AI
    • Algorithm
    • Deep Learning
    • Machine Learning
  • IOT
    • Device
    • MSP430
  • DIY
    • Algorithm
    • Design Patterns
    • MATH
    • X98 AIR 3G
    • Tucao
    • fun
  • LIFE
    • 美食
    • 关于我
  • LINKS
  • ME
Claves
阅山笑看风云起,意气扬帆向日辉
  1. 首页
  2. Platforms
  3. 正文

dbt在制造业工厂数据治理的配置教程

2026-09-06

dbt中文教程文档见下:

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 软件版本(实测)

组件版本说明
Python3.14.4虚拟环境位于仓库根目录 venv/
dbt-core1.12.3数据转换框架
dbt-duckdb1.11.0DuckDB 适配器
duckdb1.5.5嵌入式分析数据库
dbt-core-experimental-parser2.0.0rc1dbt-core 1.12 依赖的实验性解析器
PyYAML6.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.odsODS 层默认 view 物化,落到 ods schema
models.dbt_factory.dwdDWD 层默认 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 分层规范一致):

层命名物化职责
ODSods_<系统>_<表>view与源表 1:1;增量/实时表追加 _ods_loaded_at 列;声明 source freshness
DWDdwd_<域>_<实体>table / incremental去重(QUALIFY/ROW_NUMBER)、空值(COALESCE)、日期/金额标准化(宏)、字段拆分、条件标签
DWSdws_<域>_<粒度>table日/车间粒度 GROUP BY 汇总、PIVOT 行转列、UNPIVOT 列转行、UNNEST 拆多行
ADSads_<主题>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_datefmt_date(col)三种字符串日期(YYYY-MM-DD/YYYY/MM/DD/YYYYMMDD)统一为 DATE,空值保持 NULL{{ fmt_date('order_date') }} as order_date
fmt_datetimefmt_datetime(col)'YYYY-MM-DD HH:MM:SS' 字符串转 TIMESTAMP{{ fmt_datetime('update_time') }}
fmt_amountfmt_amount(col_fen)分(BIGINT)→ 元,round 2,decimal(18,2){{ fmt_amount('total_fen') }} as total_yuan
cents_to_yuancents_to_yuan(col_fen)fmt_amount 的语义别名(供自助建模按名取用){{ cents_to_yuan('price_fen') }}
bool_labelbool_label(cond, true_label, false_label)布尔条件 → 中文标签{{ bool_label('so2 > 35 or dust > 10 or voc > 100', '超标', '达标') }} as over_label
reach_labelreach_label(value, limit)达标判定(值 ≤ 限值则"达标",NULL 则"无数据"){{ reach_label('voc', 100) }}
dedup_bydedup_by(partition_cols, order_col)生成 QUALIFY row_number() 去重片段{{ dedup_by('unit_no', 'id desc') }}
generate_schema_namedbt 内建覆盖见第 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_testsnot_null / unique / accepted_values / relationships(如订单行项目 → 订单的外键关系)
自定义 generic not_exceed_thresholdtests/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/TMSlog_25车辆/运单/称重/半成品转运
BIP 协同办公bip_12组织/人员/审批流
安防sec_15视频/报警/巡更/消防
门禁acs_10门禁点/刷卡/权限
环保 EHSenv_20CEMS/烟气/废水/固废/危废
设备 EAMeam_15台账/点检/保养/维修
质检 LIMSlims_12原辅料/半成品/成品理化
实时 SCADArt_15工艺参数、环保高频

1.2 同步方式(sync 字段)

sync表数含义典型表
full46主数据/维度类,init 一次性写入,run 模式不追加mes_unit_profile、erp_material、eam_equipment
incremental158业务流水类,run 模式每 15s 批量追加 520 行mes_prod_record、erp_sales_order、log_weigh
realtime20高频监测类,run 模式每 2s 追加 1050 行全部 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 ALLmes_forming_c1_batch / mes_forming_c2_batch 完全同构dwd_mes_forming_batch_all
重复行acs_swipe_record 含约 3% 整行重复dwd_acs_swipe_dedup
空值各表 remark 列约 5% 为 NULLDWD 通用 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-A01split_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生产明细incrementaldelete+insert(unique_key=prod_id,ts 水位)生产记录 JOIN 单元档案/班组/生产计划,计划完成率公式列与单元状态标签
dwd_mes_unit_daily生产单元日报清洗incrementaldelete+insert(unique_key=id)日期统一,产出效率/电耗标签
dwd_acs_swipe_dedup门禁刷卡去重incrementalappend(ts 水位)QUALIFY row_number 去除源表约 3% 整行重复
dwd_env_cems_hour废气CEMS小时均值清洗incrementalappend(ts 水位)超标标签+分污染物异常 flag(SO2>35/粉尘>10/VOCs>100)
dwd_ems_meter_reading计量点读数清洗incrementalappend(ts 水位)JOIN 计量点档案(先去重),供 DWS PIVOT
dwd_log_weighbridge地磅称重清洗incrementalappend(ts 水位)毛重-皮重核算净重公式列,净重差异>0.5t 标异常
dwd_erp_sales_order销售订单incrementalappend(id 水位)金额分→元(fmt_amount),日期统一,状态标签
dwd_ems_unit_consume产品单耗清洗incrementalappend(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 快照snapshottimestamp 策略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 个数据测试,全部通过)

类别内容例子
内建 genericnot_null / unique / accepted_values / relationships主键非空唯一;枚举值域(单元状态、达标标签等);订单行项目 → 订单外键
自定义 genericnot_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 /
2101
2260)。工艺流程覆盖多个车间/产线,主要物料流向如下:

原料仓 ──> 成型/加工车间 ──> 装配/包装车间 ──> 半成品/成品 ──> 质检
  │                                              │
  └─ 余料/不合格品回收处理                      └─> 成品仓 ──> 销售发运
工艺废气/废水/固废 ──> 环保处理系统(净化 + 在线监测)──> 达标排放/合规处置

配套支撑:动力车间(变配电/空压/水泵)、维修车间(设备维修与保养)、物流(原料到货、
半成品转运、成品发运)、环保(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_dayKPI 卡片 + 产量/电耗双轴趋势线 + 达标率仪表盘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_avgdws_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)~520EMS 源表扩展(见路线图)
天然气/蒸汽/水单耗能耗定额表 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 侧只读打开
定时导出 parquetPowerBI、帆软 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_* 表时同样只读打开,写入权只归模拟器/采集程序。
标签: 暂无
最后更新:2026-09-06

阅山

知之为知之 不知为不知

点赞

COPYRIGHT © 2099 登峰造极境. ALL RIGHTS RESERVED.

Theme Kratos Made By Seaton Jiang

蜀ICP备14031139号-5

川公网安备51012202000587号