Metadata-Version: 2.3
Name: wedata-dq-expectations
Version: 2.1.0
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
pip install wedata-dq-expectations
```

```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

# 初始化运行时（DLC 环境下自动读取任务标识，用于质量结果留痕）
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` 仅首次创建，表已存在则报错 |
| `safe_check` | 覆盖保护开关，默认 `True`（见下） |

### 覆盖保护

`overwrite` 模式下默认开启覆盖保护：只有由本 SDK 创建/管理的表才允许被覆盖。若目标表是已存在的业务表（非本 SDK 管理），写入会被拦截并报错，避免误覆盖生产数据。

如确需覆盖非本 SDK 管理的表，可显式关闭：`@table(..., safe_check=False)`（请谨慎使用），或改用 `mode="append"`。

## 查看校验结果

每次运行，各规则的校验指标会写入质量结果表，可直接用 SQL 查询：

```sql
SELECT table_name, rule_name, severity, action, passed, failed, pass_rate
FROM system_catalog.wedata.data_quality_expectations;
```

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

### 结果表位置（可配置）

结果表默认写入 `system_catalog.wedata.data_quality_expectations`。若你的引擎无法访问该 catalog（例如报 `only support namespace with 1 level`、`REQUIRES_SINGLE_PART_NAMESPACE` 等），可按以下任一方式改到你能访问的库表，优先级从高到低：

```python
# 方式一：调用时显式指定（最高优先级）
init_runtime(result_table="DataLakeCatalog.my_db.data_quality_expectations")
```

```bash
# 方式二：环境变量（引擎级统一配置，无需改脚本）
export WEDATA_DQ_RESULT_TABLE="DataLakeCatalog.my_db.data_quality_expectations"
```

不设置时使用默认值。查询时对应改成你配置的表名即可。

## 任务标识

`init_runtime()` 用于设置任务标识，让质量结果能关联到具体任务：

```python
# 方式一：DLC 调度环境下自动读取平台注入的系统变量
init_runtime()

# 方式二：显式指定
init_runtime(task_id="my_task_001", task_name="用户画像清洗")

# 方式三：指定读取系统变量的 widget key
init_runtime(task_id_widget="wedata_dq_task_id")

# 可与结果表配置组合使用
init_runtime(task_id="my_task_001", result_table="DataLakeCatalog.my_db.data_quality_expectations")
```

在非 DLC 环境（如本地调试）且未显式传参时，任务标识降级为 `unknown`，不影响主流程运行。

## 环境要求

- Python 3.8+
- 运行环境需已具备 PySpark（DLC 引擎 / Notebook 内核默认自带），本包不额外安装 PySpark。
