commit 3f3dc443da248ae6cb2619524a897577b2fbb48f Author: chenwu Date: Thu Aug 20 18:22:16 2026 +0800 refactor CHPA hierarchy SQL diff --git a/README.md b/README.md new file mode 100644 index 0000000..e3208c6 --- /dev/null +++ b/README.md @@ -0,0 +1,54 @@ +# REFACTOR-MA + +Databricks 数仓 SQL 重构仓库。 + +## 首批 review 范围 + +本批次只重构原 `CHPA/01` 中两个层级维表脚本: + +- `sql/chpa/01_dwd/01_dwd_ims_atc_hierarchy.sql` +- `sql/chpa/01_dwd/02_dwd_ims_nfc_hierarchy.sql` + +工作区 `CHPA/` 下的原脚本保持不变,便于逐项对照。 + +## 兼容性决策 + +首批代码继续写入现有物理表: + +- `dwd.dwd_ims_atc_hierarchy` +- `dwd.dwd_ims_nfc_hierarchy` + +这样不会立即影响 pack property 和 DWS 等现有下游任务。推荐的新表名及迁移规则记录在 `docs/sql_refactoring_standard.md`;是否执行物理改名,待本批 review 后决定。 + +## 目录结构 + +```text +RE/ +|-- README.md +|-- docs/ +| `-- sql_refactoring_standard.md +|-- sql/ +| `-- chpa/ +| `-- 01_dwd/ +| |-- 01_dwd_ims_atc_hierarchy.sql +| `-- 02_dwd_ims_nfc_hierarchy.sql +`-- validation/ + `-- chpa/ + `-- 01_dwd/ + `-- validate_hierarchy_refactor.sql +``` + +## 执行顺序 + +1. 在覆盖目标表前运行 `validation/chpa/01_dwd/validate_hierarchy_refactor.sql`。 +2. 核对两个 `DESCRIBE TABLE` 结果;兼容性检查必须全部返回 `passed = true`。 +3. review 重复路径和未匹配层级指标,确认它们符合现有业务口径。 +4. 通过后,按文件名前缀顺序执行 `sql/chpa/01_dwd/` 中的重构脚本。 + +两个重构脚本都要求现有 DWD 源表、目标表已经创建。 + +## 本次 review 重点 + +1. 表命名模型是否适用于整个数仓。 +2. 物理表改名采用兼容视图还是一次性联动下游。 +3. 替换调度任务前,在 Databricks 中对比新旧结果的行数、业务键重复数和空键数。 diff --git a/docs/sql_refactoring_standard.md b/docs/sql_refactoring_standard.md new file mode 100644 index 0000000..c841df9 --- /dev/null +++ b/docs/sql_refactoring_standard.md @@ -0,0 +1,119 @@ +# Databricks SQL 重构规范 + +## 1. 范围和原则 + +本规范适用于 CHPA Databricks 数仓脚本。第一阶段目标是在不改变业务结果的前提下提高可维护性。 + +- 先保持业务口径,再做性能优化。 +- 每个脚本只负责一个主要目标表。 +- 文件头声明源表、目标表、数据粒度、写入方式和依赖。 +- 持久化写入必须显式列出目标列和查询列,禁止使用 `SELECT *`。 +- 逻辑重构与生产表改名分开实施。 +- 为行数、唯一性、非空约束和未匹配记录建立验证查询或任务检查。 + +## 2. 文件和目录命名 + +目录结构统一为: + +```text +sql/<业务域>/<阶段>/<序号>_<目标表>.sql +``` + +规则: + +- 目录和文件名使用小写 snake_case。 +- 路径和文件名不使用空格。 +- 使用两位序号明确 notebook/job 的执行顺序。 +- 一个脚本只写入一个主要持久化目标。像 `01 dwd_update.sql` 这种多目标脚本,应按目标表和职责拆分。 + +示例: + +```text +sql/chpa/01_dwd/01_dwd_ims_atc_hierarchy.sql +``` + +## 3. 物理表命名 + +新表使用以下格式: + +```text +[.]<分层_schema>.<来源系统>_<对象类型>_<业务实体>[_<限定词>] +``` + +未启用 Unity Catalog 或依赖默认 catalog 时可以省略 catalog。分层已经由 schema 表达,表名不再重复 `dwd`、`dws` 或 `dm`。 + +- 来源系统:如 `ims`、`gnd`、`pharbers`。 +- 对象类型:`td` 表示维度/主数据,`tf` 表示事实数据。 +- 业务实体:使用稳定的业务名词,如 `atc_hierarchy`、`pack_property`。 +- 限定词:按需表达粒度或变体,如 `monthly`、`province`。 + +`01` 组表名建议: + +| 现有物理表名 | 推荐规范表名 | 首批处理方式 | +| --- | --- | --- | +| `dwd.dwd_ims_atc_hierarchy` | `dwd.ims_td_atc_hierarchy` | 保留现名 | +| `dwd.dwd_ims_nfc_hierarchy` | `dwd.ims_td_nfc_hierarchy` | 保留现名 | +| `dwd.dwd_ims_td_manufacturer_corp` | `dwd.ims_td_manufacturer_corporation` | 后续 review | +| `dwd.dwd_ims_td_pack_property` | `dwd.ims_td_pack_property` | 后续 review | +| `dwd.dwd_gnd_pharbers_prov_fact` | `dwd.pharbers_tf_province_sales` | 先确认业务粒度 | + +推荐迁移顺序: + +1. 创建规范名称的新表,并用兼容视图暴露旧表名。 +2. 新旧表并行运行,对比数据结果。 +3. 在一次受控发布中修改全部下游引用。 +4. 所有消费者迁移后再移除旧名称兼容视图。 + +## 4. 字段命名 + +- 新字段使用小写 snake_case。 +- 标识使用 `_id`,业务编码使用 `_code`,名称或描述使用 `_name`、`_description`。 +- 日期时间后缀按真实类型使用 `_date`、`_timestamp` 或 `_at`。 +- `atc`、`nfc`、`prod`、`pack` 等已形成业务共识或受输出契约约束的缩写可以保留。 +- 纯逻辑重构中不直接修改持久化输出字段名。 + +## 5. SQL 结构和格式 + +脚本统一按以下顺序组织: + +1. Databricks notebook 标记和脚本契约头。 +2. 带显式目标列的 `INSERT OVERWRITE`。 +3. 源数据标准化 CTE。 +4. 业务规则 CTE。 +5. 按目标表字段顺序显式编写最终 `SELECT`。 + +格式规则: + +- SQL 关键字使用大写。 +- CTE 和别名使用小写 snake_case。 +- 使用四个空格缩进。 +- `SELECT` 每行一个字段。 +- 每个 JOIN 条件单独成行。 +- 使用 `atc_level_1` 等有业务意义的别名,维护代码中不使用 `t1`、`a` 等无语义别名。 +- 同一作用域出现多个关系时,所有字段都带关系限定符。 + +## 6. 写入和数据质量规则 + +- 只有可确定性重跑的全量任务可以使用 `INSERT OVERWRITE`。 +- 不在一个脚本中混合无关的 `UPDATE` 和目标表构建逻辑。 +- 编码补零等标准化逻辑集中到一个明确的处理阶段,避免在多个任务中重复表达式。 +- 每个脚本必须声明预期目标粒度。 +- 最少验证以下指标: + - 新结果与旧结果的总行数; + - 声明业务键的重复数; + - 必填键的空值数; + - 层级关联新增的未匹配记录数。 + +## 7. 首批兼容说明 + +ATC 和 NFC 脚本原样保留以下父级编码规则: + +- ATC 2 级关联 1 级:取前 1 位。 +- ATC 3 级关联 2 级:取前 3 位。 +- ATC 4 级关联 3 级:取前 4 位。 +- NFC 2 级关联 1 级:取前 1 位。 +- NFC 3 级关联 2 级:取前 2 位。 + +任何口径简化都应先用生产样本验证这些规则。 + +首批上线前运行 `validation/chpa/01_dwd/validate_hierarchy_refactor.sql`。新旧双向差集必须为 0,行数必须一致;重复路径和未匹配层级作为业务 review 指标记录,不在未确认阈值前自动判定失败。 diff --git a/sql/chpa/01_dwd/01_dwd_ims_atc_hierarchy.sql b/sql/chpa/01_dwd/01_dwd_ims_atc_hierarchy.sql new file mode 100644 index 0000000..558fe43 --- /dev/null +++ b/sql/chpa/01_dwd/01_dwd_ims_atc_hierarchy.sql @@ -0,0 +1,86 @@ +-- Databricks notebook source +-- ============================================================================= +-- Purpose : Build the flattened IMS ATC level 1-4 hierarchy. +-- Source : dwd.dwd_ims_td_therapeutic_class +-- Target : dwd.dwd_ims_atc_hierarchy +-- Grain : One row per hierarchy path rooted at ATC1_CODE; lower levels may be null. +-- Write mode : Full refresh (INSERT OVERWRITE). +-- Downstream : dwd.dwd_ims_td_pack_property, dws.dws_ims_td_atc_cn +-- Notes : Parent-code matching is preserved from the legacy script. +-- ============================================================================= + +INSERT OVERWRITE TABLE dwd.dwd_ims_atc_hierarchy ( + ATC1_ID, + ATC1_CODE, + ATC1_DES, + ATC2_ID, + ATC2_CODE, + ATC2_DES, + ATC3_ID, + ATC3_CODE, + ATC3_DES, + ATC4_ID, + ATC4_CODE, + ATC4_DES +) +WITH therapeutic_class AS ( + SELECT + therapeutic_id, + therapeutic_code, + therapeutic_name, + therapeutic_level + FROM dwd.dwd_ims_td_therapeutic_class +), +atc_level_1 AS ( + SELECT + therapeutic_id AS atc1_id, + therapeutic_code AS atc1_code, + therapeutic_name AS atc1_des + FROM therapeutic_class + WHERE therapeutic_level = '1' +), +atc_level_2 AS ( + SELECT + therapeutic_id AS atc2_id, + therapeutic_code AS atc2_code, + therapeutic_name AS atc2_des + FROM therapeutic_class + WHERE therapeutic_level = '2' +), +atc_level_3 AS ( + SELECT + therapeutic_id AS atc3_id, + therapeutic_code AS atc3_code, + therapeutic_name AS atc3_des + FROM therapeutic_class + WHERE therapeutic_level = '3' +), +atc_level_4 AS ( + SELECT + therapeutic_id AS atc4_id, + therapeutic_code AS atc4_code, + therapeutic_name AS atc4_des + FROM therapeutic_class + WHERE therapeutic_level = '4' +) +SELECT + atc_level_1.atc1_id, + atc_level_1.atc1_code, + atc_level_1.atc1_des, + atc_level_2.atc2_id, + atc_level_2.atc2_code, + atc_level_2.atc2_des, + atc_level_3.atc3_id, + atc_level_3.atc3_code, + atc_level_3.atc3_des, + atc_level_4.atc4_id, + atc_level_4.atc4_code, + atc_level_4.atc4_des +FROM atc_level_1 +LEFT JOIN atc_level_2 + ON atc_level_1.atc1_code = LEFT(atc_level_2.atc2_code, 1) +LEFT JOIN atc_level_3 + ON atc_level_2.atc2_code = LEFT(atc_level_3.atc3_code, 3) +LEFT JOIN atc_level_4 + ON atc_level_3.atc3_code = LEFT(atc_level_4.atc4_code, 4) +; diff --git a/sql/chpa/01_dwd/02_dwd_ims_nfc_hierarchy.sql b/sql/chpa/01_dwd/02_dwd_ims_nfc_hierarchy.sql new file mode 100644 index 0000000..cbc81ce --- /dev/null +++ b/sql/chpa/01_dwd/02_dwd_ims_nfc_hierarchy.sql @@ -0,0 +1,70 @@ +-- Databricks notebook source +-- ============================================================================= +-- Purpose : Build the flattened IMS NFC level 1-3 hierarchy. +-- Source : dwd.dwd_ims_td_new_form_class +-- Target : dwd.dwd_ims_nfc_hierarchy +-- Grain : One row per hierarchy path rooted at NFC1_CODE; lower levels may be null. +-- Write mode : Full refresh (INSERT OVERWRITE). +-- Downstream : dwd.dwd_ims_td_pack_property, dws.dws_ims_td_nfc_cn +-- Notes : Parent-code matching is preserved from the legacy script. +-- ============================================================================= + +INSERT OVERWRITE TABLE dwd.dwd_ims_nfc_hierarchy ( + NFC1_ID, + NFC1_CODE, + NFC1_DES, + NFC2_ID, + NFC2_CODE, + NFC2_DES, + NFC3_ID, + NFC3_CODE, + NFC3_DES +) +WITH new_form_class AS ( + SELECT + newformclass_id, + newformclass_code, + newformclass_name, + newformclass_level + FROM dwd.dwd_ims_td_new_form_class +), +nfc_level_1 AS ( + SELECT + newformclass_id AS nfc1_id, + newformclass_code AS nfc1_code, + newformclass_name AS nfc1_des + FROM new_form_class + WHERE newformclass_level = '1' +), +nfc_level_2 AS ( + SELECT + newformclass_id AS nfc2_id, + newformclass_code AS nfc2_code, + newformclass_name AS nfc2_des + FROM new_form_class + WHERE newformclass_level = '2' +), +nfc_level_3 AS ( + SELECT + newformclass_id AS nfc3_id, + newformclass_code AS nfc3_code, + newformclass_name AS nfc3_des + FROM new_form_class + WHERE newformclass_level = '3' +) +SELECT + nfc_level_1.nfc1_id, + nfc_level_1.nfc1_code, + nfc_level_1.nfc1_des, + nfc_level_2.nfc2_id, + nfc_level_2.nfc2_code, + nfc_level_2.nfc2_des, + nfc_level_3.nfc3_id, + nfc_level_3.nfc3_code, + nfc_level_3.nfc3_des +FROM nfc_level_1 +LEFT JOIN nfc_level_2 + ON nfc_level_1.nfc1_code = LEFT(nfc_level_2.nfc2_code, 1) +LEFT JOIN nfc_level_3 + ON nfc_level_2.nfc2_code = LEFT(nfc_level_3.nfc3_code, 2) +; diff --git a/validation/chpa/01_dwd/validate_hierarchy_refactor.sql b/validation/chpa/01_dwd/validate_hierarchy_refactor.sql new file mode 100644 index 0000000..6983f97 --- /dev/null +++ b/validation/chpa/01_dwd/validate_hierarchy_refactor.sql @@ -0,0 +1,263 @@ +-- Databricks notebook source +-- MAGIC %md +-- MAGIC # ATC/NFC hierarchy refactor validation +-- MAGIC Run this notebook before executing the refactored overwrite scripts. +-- MAGIC The legacy targets and source tables must represent the same data snapshot. + +-- COMMAND ---------- + +-- Confirm the ATC target contains the expected 12 non-partition output columns. +DESCRIBE TABLE dwd.dwd_ims_atc_hierarchy; + +-- COMMAND ---------- + +-- Confirm the NFC target contains the expected 9 non-partition output columns. +DESCRIBE TABLE dwd.dwd_ims_nfc_hierarchy; + +-- COMMAND ---------- + +CREATE OR REPLACE TEMP VIEW validation_atc_hierarchy_candidate AS +WITH therapeutic_class AS ( + SELECT + therapeutic_id, + therapeutic_code, + therapeutic_name, + therapeutic_level + FROM dwd.dwd_ims_td_therapeutic_class +), +atc_level_1 AS ( + SELECT + therapeutic_id AS atc1_id, + therapeutic_code AS atc1_code, + therapeutic_name AS atc1_des + FROM therapeutic_class + WHERE therapeutic_level = '1' +), +atc_level_2 AS ( + SELECT + therapeutic_id AS atc2_id, + therapeutic_code AS atc2_code, + therapeutic_name AS atc2_des + FROM therapeutic_class + WHERE therapeutic_level = '2' +), +atc_level_3 AS ( + SELECT + therapeutic_id AS atc3_id, + therapeutic_code AS atc3_code, + therapeutic_name AS atc3_des + FROM therapeutic_class + WHERE therapeutic_level = '3' +), +atc_level_4 AS ( + SELECT + therapeutic_id AS atc4_id, + therapeutic_code AS atc4_code, + therapeutic_name AS atc4_des + FROM therapeutic_class + WHERE therapeutic_level = '4' +) +SELECT + atc_level_1.atc1_id AS ATC1_ID, + atc_level_1.atc1_code AS ATC1_CODE, + atc_level_1.atc1_des AS ATC1_DES, + atc_level_2.atc2_id AS ATC2_ID, + atc_level_2.atc2_code AS ATC2_CODE, + atc_level_2.atc2_des AS ATC2_DES, + atc_level_3.atc3_id AS ATC3_ID, + atc_level_3.atc3_code AS ATC3_CODE, + atc_level_3.atc3_des AS ATC3_DES, + atc_level_4.atc4_id AS ATC4_ID, + atc_level_4.atc4_code AS ATC4_CODE, + atc_level_4.atc4_des AS ATC4_DES +FROM atc_level_1 +LEFT JOIN atc_level_2 + ON atc_level_1.atc1_code = LEFT(atc_level_2.atc2_code, 1) +LEFT JOIN atc_level_3 + ON atc_level_2.atc2_code = LEFT(atc_level_3.atc3_code, 3) +LEFT JOIN atc_level_4 + ON atc_level_3.atc3_code = LEFT(atc_level_4.atc4_code, 4) +; + +-- COMMAND ---------- + +CREATE OR REPLACE TEMP VIEW validation_nfc_hierarchy_candidate AS +WITH new_form_class AS ( + SELECT + newformclass_id, + newformclass_code, + newformclass_name, + newformclass_level + FROM dwd.dwd_ims_td_new_form_class +), +nfc_level_1 AS ( + SELECT + newformclass_id AS nfc1_id, + newformclass_code AS nfc1_code, + newformclass_name AS nfc1_des + FROM new_form_class + WHERE newformclass_level = '1' +), +nfc_level_2 AS ( + SELECT + newformclass_id AS nfc2_id, + newformclass_code AS nfc2_code, + newformclass_name AS nfc2_des + FROM new_form_class + WHERE newformclass_level = '2' +), +nfc_level_3 AS ( + SELECT + newformclass_id AS nfc3_id, + newformclass_code AS nfc3_code, + newformclass_name AS nfc3_des + FROM new_form_class + WHERE newformclass_level = '3' +) +SELECT + nfc_level_1.nfc1_id AS NFC1_ID, + nfc_level_1.nfc1_code AS NFC1_CODE, + nfc_level_1.nfc1_des AS NFC1_DES, + nfc_level_2.nfc2_id AS NFC2_ID, + nfc_level_2.nfc2_code AS NFC2_CODE, + nfc_level_2.nfc2_des AS NFC2_DES, + nfc_level_3.nfc3_id AS NFC3_ID, + nfc_level_3.nfc3_code AS NFC3_CODE, + nfc_level_3.nfc3_des AS NFC3_DES +FROM nfc_level_1 +LEFT JOIN nfc_level_2 + ON nfc_level_1.nfc1_code = LEFT(nfc_level_2.nfc2_code, 1) +LEFT JOIN nfc_level_3 + ON nfc_level_2.nfc2_code = LEFT(nfc_level_3.nfc3_code, 2) +; + +-- COMMAND ---------- + +-- Every check must return passed = true before replacing the scheduled jobs. +WITH compatibility_checks AS ( + SELECT + 'atc_row_count' AS check_name, + (SELECT COUNT(*) FROM validation_atc_hierarchy_candidate) AS actual_value, + (SELECT COUNT(*) FROM dwd.dwd_ims_atc_hierarchy) AS expected_value + UNION ALL + SELECT + 'atc_legacy_minus_candidate' AS check_name, + ( + SELECT COUNT(*) + FROM ( + SELECT * FROM dwd.dwd_ims_atc_hierarchy + EXCEPT ALL + SELECT * FROM validation_atc_hierarchy_candidate + ) AS differences + ) AS actual_value, + 0 AS expected_value + UNION ALL + SELECT + 'atc_candidate_minus_legacy' AS check_name, + ( + SELECT COUNT(*) + FROM ( + SELECT * FROM validation_atc_hierarchy_candidate + EXCEPT ALL + SELECT * FROM dwd.dwd_ims_atc_hierarchy + ) AS differences + ) AS actual_value, + 0 AS expected_value + UNION ALL + SELECT + 'nfc_row_count' AS check_name, + (SELECT COUNT(*) FROM validation_nfc_hierarchy_candidate) AS actual_value, + (SELECT COUNT(*) FROM dwd.dwd_ims_nfc_hierarchy) AS expected_value + UNION ALL + SELECT + 'nfc_legacy_minus_candidate' AS check_name, + ( + SELECT COUNT(*) + FROM ( + SELECT * FROM dwd.dwd_ims_nfc_hierarchy + EXCEPT ALL + SELECT * FROM validation_nfc_hierarchy_candidate + ) AS differences + ) AS actual_value, + 0 AS expected_value + UNION ALL + SELECT + 'nfc_candidate_minus_legacy' AS check_name, + ( + SELECT COUNT(*) + FROM ( + SELECT * FROM validation_nfc_hierarchy_candidate + EXCEPT ALL + SELECT * FROM dwd.dwd_ims_nfc_hierarchy + ) AS differences + ) AS actual_value, + 0 AS expected_value +) +SELECT + check_name, + actual_value, + expected_value, + actual_value = expected_value AS passed +FROM compatibility_checks +ORDER BY check_name; + +-- COMMAND ---------- + +-- Profile duplicates and unmatched levels. These are review metrics, not automatic failures. +WITH atc_duplicate_paths AS ( + SELECT COUNT(*) AS path_count + FROM validation_atc_hierarchy_candidate + GROUP BY ATC1_CODE, ATC2_CODE, ATC3_CODE, ATC4_CODE + HAVING COUNT(*) > 1 +), +nfc_duplicate_paths AS ( + SELECT COUNT(*) AS path_count + FROM validation_nfc_hierarchy_candidate + GROUP BY NFC1_CODE, NFC2_CODE, NFC3_CODE + HAVING COUNT(*) > 1 +) +SELECT + 'atc_duplicate_path_rows' AS metric_name, + COALESCE(SUM(path_count - 1), 0) AS metric_value +FROM atc_duplicate_paths +UNION ALL +SELECT + 'atc_null_root_code' AS metric_name, + COUNT_IF(ATC1_CODE IS NULL) AS metric_value +FROM validation_atc_hierarchy_candidate +UNION ALL +SELECT + 'atc_unmatched_level_2' AS metric_name, + COUNT_IF(ATC2_CODE IS NULL) AS metric_value +FROM validation_atc_hierarchy_candidate +UNION ALL +SELECT + 'atc_unmatched_level_3' AS metric_name, + COUNT_IF(ATC3_CODE IS NULL) AS metric_value +FROM validation_atc_hierarchy_candidate +UNION ALL +SELECT + 'atc_unmatched_level_4' AS metric_name, + COUNT_IF(ATC4_CODE IS NULL) AS metric_value +FROM validation_atc_hierarchy_candidate +UNION ALL +SELECT + 'nfc_duplicate_path_rows' AS metric_name, + COALESCE(SUM(path_count - 1), 0) AS metric_value +FROM nfc_duplicate_paths +UNION ALL +SELECT + 'nfc_null_root_code' AS metric_name, + COUNT_IF(NFC1_CODE IS NULL) AS metric_value +FROM validation_nfc_hierarchy_candidate +UNION ALL +SELECT + 'nfc_unmatched_level_2' AS metric_name, + COUNT_IF(NFC2_CODE IS NULL) AS metric_value +FROM validation_nfc_hierarchy_candidate +UNION ALL +SELECT + 'nfc_unmatched_level_3' AS metric_name, + COUNT_IF(NFC3_CODE IS NULL) AS metric_value +FROM validation_nfc_hierarchy_candidate +ORDER BY metric_name;