Metadata-Version: 2.3
Name: wedata-dq-expectations
Version: 2.7.0rc6
Summary: WeData 数据质量 Expectations SDK — PySpark 声明式数据质量校验装饰器，支持 ETL 过程中的质量卡点、目标表持久化与覆盖保护。
Author: WeData Team
Author-email: WeData Team <wedata@tencent.com>
License: MIT
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.8
Classifier: Programming Language :: Python :: 3.9
Classifier: Programming Language :: Python :: 3.10
Classifier: Programming Language :: Python :: 3.11
Classifier: Topic :: Software Development :: Quality Assurance
Requires-Python: >=3.8
Project-URL: Homepage, https://wedata.tencent.com
Project-URL: Documentation, https://wedata.tencent.com/docs
Description-Content-Type: text/markdown

# wedata-dq-expectations

WeData 数据质量 Expectations SDK。在 PySpark / Notebook 脚本中，用声明式的方式在 **ETL 加工过程中**对数据做质量校验，把脏数据在落库前拦下来。

```bash
# 建议 pin 到「当前 minor 版本区间」：只收该 minor 的补丁，不跨 minor 升到含破坏性变更的新版本。
# 写法 >=X.Y,<X.(Y+1)，其中 X.Y 替换为你使用时的最新稳定版（下方 2.6 仅为示例）：
pip install "wedata-dq-expectations>=2.6,<2.7"
```

```python
from wedata.dq import table, expect, expect_or_drop, expect_or_fail, init_runtime
```

## 适用场景

传统质量规则只能挂在 ETL 节点**之后**做事后检查——脏数据已经落库才被发现。本 SDK 让质量校验前移到加工过程里，在数据写入下游之前完成校验，并按策略自动**留痕 / 隔离 / 阻断**。

典型场景：

- **入口校验**：读完源表、进入加工前，先校验源数据是否符合预期。
- **落库前门禁**：写目标表之前校验，不合格的数据不落库。
- **关键字段兜底**：主键非空、金额为正、状态枚举合法等硬约束，一旦违反立即中断任务，避免污染下游。

## 快速上手

给你的表转换函数加上装饰器，声明「这张表应该满足哪些质量规则」，SDK 会在函数返回 DataFrame 后自动校验、按策略处理，并把结果写入下游表：

```python
from wedata.dq import table, expect, expect_or_drop, expect_or_fail, init_runtime

# 初始化运行时（自动采集工作流溯源参数；可传 result_table 指定结果表位置）
init_runtime()

@table(name="my_catalog.my_db.customers_clean", persist=True, mode="overwrite")
@expect("valid_age", "age BETWEEN 18 AND 100", severity="warn")        # 只留痕，不动数据
@expect_or_drop("valid_email", "email IS NOT NULL AND email != ''")    # 丢弃违规行
@expect_or_fail("has_id", "customer_id IS NOT NULL")                   # 违规则中断任务
def transform_customers():
    return spark.table("raw_customers")

transform_customers()
```

运行后，`customers_clean` 表里只保留通过校验（且未被丢弃）的数据，每条规则的通过率、违规行数等指标自动落到质量结果表，可直接查询。

## 规则声明与处理策略

用三个装饰器声明规则，对应三种处理策略。规则表达式是标准 Spark SQL 布尔表达式。

| 装饰器 | 策略 | 违规时的行为 | 适用 |
|--------|------|-------------|------|
| `@expect(name, expr)` | 留痕 | 记录通过率与违规数，**数据照常写入** | 探索性分析、质量观测 |
| `@expect_or_drop(name, expr)` | 丢弃 | 违规行被过滤，不写入目标表 | 生产清洗、严格质量要求 |
| `@expect_or_fail(name, expr)` | 阻断 | 发现违规立即中断任务，已写数据回滚、不产生下游表 | 关键业务数据、零容忍 |

参数说明：

- `name`：规则名称，用于在结果中标识该规则。
- `expr`：Spark SQL 布尔表达式，返回真表示该行通过校验，如 `"amount > 0"`、`"status IN ('pending','completed')"`。
- `severity`：严重级别（`warn` / `error` / `critical`），仅用于结果标记，不影响处理逻辑。

多个装饰器可叠加在同一个函数上，规则会依次执行。

## 写入目标表

`@table` 负责把校验后的 DataFrame 写入目标表：

```python
@table(name="catalog.schema.table", persist=True, mode="overwrite")
```

| 参数 | 说明 |
|------|------|
| `name` | 目标表名，支持 `catalog.schema.table` 三段式；缺省时取函数名（去掉 `transform_` 前缀） |
| `persist` | `True` 持久化写表；`False`（默认）仅注册临时视图，供同一脚本内后续步骤引用 |
| `mode` | `overwrite` 全量覆盖（默认）/ `append` 追加 / `safe` 仅首次创建，表已存在则报错 |

> ⚠️ `mode="overwrite"` 会**全量覆盖**目标表的已有数据。目标表由你在 `name` 中显式指定，请自行确认该表可被覆盖；若不希望覆盖已有数据，改用 `mode="append"`（追加）或 `mode="safe"`（仅允许首次创建，表已存在则报错）。

## 查看校验结果

各规则的校验指标会写入你指定的质量结果表（见下方「结果表位置」），可直接用 SQL 查询：

```sql
SELECT table_name, rule_name, severity, action, passed, failed, pass_rate,
       from_unixtime(check_time / 1000) AS check_time_readable
FROM <你指定的结果表>;
```

字段含义：目标表名、规则名、严重级别、处理策略、通过行数、违规行数、通过率、校验时刻等。

### 结果表位置（必须由你指定）

**结果表位置需要你传入一个自己有写入权限的库表**，SDK 不使用写死的默认位置。通过 `init_runtime` 的 `result_table` 参数指定：

```python
init_runtime(result_table="catalog.schema.wedata_dq_expectations")
```

**不指定时**：SDK 不落物理结果表，仅在日志打印各规则指标，并注册一个当次会话可查的临时视图 `data_quality_expectations`（`SELECT * FROM data_quality_expectations`），主流程不受影响。

> 为什么不给默认位置：写死一个公共 catalog（如 system catalog）在实际环境中常遇到「该 catalog 不存在」或「执行账号无写入权限」的问题，因此改为由你传入自己有权限的位置。

## 运行时初始化与工作流溯源

`init_runtime()` 用于初始化运行时：指定结果表位置，并自动采集工作流溯源信息，让质量结果能关联到具体工作流/执行：

```python
# 工作流调度环境下：自动读取平台注入的系统参数，结果表可选
init_runtime(result_table="catalog.schema.wedata_dq_expectations")

# 也可不指定结果表（不落物理表，仅打印指标 + 临时视图）
init_runtime()
```

### 工作流溯源字段（自动采集）

在工作流调度环境下，SDK 会自动读取工作流内置系统参数并写入结果表，用于把质量结果关联到具体工作流/执行：

| 结果表字段 | 来源系统参数 | 含义 |
|-----------|-------------|------|
| `workspace_id` | `workspaceId` | 工作空间 |
| `workflow_id` | `workflowId` | 工作流 |
| `workflow_execution_id` | `workflowExecutionId` | 工作流执行 |
| `task_execution_id` | `taskExecutionId` | 任务执行（唯一标识本次执行） |

在非工作流环境（如本地调试、直连试运行）取不到这些参数时，对应字段为空，不影响主流程运行。

### 重复执行与幂等写入

任务节点可能多次执行（同一次执行内失败重试、cell 重复运行、工作流多次调度）。结果表写入以 `task_execution_id + table_name` 为幂等键：

- **同一次执行**（`task_execution_id` 相同）重复写入：先删除该执行**该表**的旧记录再写入，**不产生重复**；同一次执行内多个 transform（多张目标表）的结果互不影响。
- **不同调度执行**（`task_execution_id` 不同）：各自保留一份，**溯源历史完整**。
- 结果表需支持行级 `DELETE`（如 Iceberg）；若表格式不支持，会自动降级为追加写入并打印告警。
- 取不到 `task_execution_id`（本地/直连试运行）时无法幂等，按追加写入。

## 版本与升级

本包遵循语义化版本（`MAJOR.MINOR.PATCH`）：

- `PATCH`（如 `x.y.0 → x.y.1`）：修 bug、向后兼容。
- `MINOR`（如 `x.y.* → x.(y+1).0`）：可能含破坏性变更（如结果表结构、API 调整），升级前请看变更说明。
- 预发布版（如 `x.y.0rc1`）：仅供内部测试，`pip install` **默认不会安装**，无需担心被自动拉取。

**强烈建议在脚本里 pin 到当前的 minor 版本区间**，避免我们发布新版本时影响你正在运行的作业。把下方 `X.Y` 替换为你使用时的最新稳定版即可（如 `2.6` → `>=2.6,<2.7`）：

```python
pip install "wedata-dq-expectations>=X.Y,<X.(Y+1)"   # 收补丁，不跨 minor
```

需要升级到新 minor 时，先在测试环境验证，再把区间上调一个 minor。

## 环境要求

- Python 3.8+
- 运行环境需已具备 PySpark（DLC 引擎 / Notebook 内核默认自带），本包不额外安装 PySpark。
- 质量结果表建议使用支持行级 `DELETE` 的表格式（如 Iceberg），以启用按执行的幂等去重。
