Metadata-Version: 2.4
Name: tsysmart-tsprep
Version: 1.0.6
Summary: Time series data preprocessing toolkit
Author-email: TSysmart Team <support@tsysmart.com>
License: Proprietary
Project-URL: Homepage, https://github.com/tsysmart/tsdataprep
Project-URL: Documentation, https://github.com/tsysmart/tsdataprep#readme
Project-URL: Repository, https://github.com/tsysmart/tsdataprep.git
Keywords: time-series,data-preprocessing,data-cleaning,anomaly-detection
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
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
Requires-Python: >=3.8
Description-Content-Type: text/markdown
License-File: LICENSE
Requires-Dist: pandas>=1.3.0
Requires-Dist: numpy>=1.20.0
Requires-Dist: pytz>=2021.1
Requires-Dist: PyYAML>=5.4
Requires-Dist: tsysmart_utils>=1.7
Requires-Dist: tsysmart_proj
Requires-Dist: tsysmart_appmng
Requires-Dist: sqlalchemy>=1.4
Requires-Dist: fastapi>=0.115.11
Requires-Dist: uvicorn>=0.34.0
Requires-Dist: pydantic>=2.0
Requires-Dist: prometheus-client
Requires-Dist: tomlkit>=0.14
Requires-Dist: tsysmart_pkg>=1.3
Requires-Dist: openpyxl>=3.0
Provides-Extra: dev
Requires-Dist: pytest>=6.0; extra == "dev"
Requires-Dist: black>=21.0; extra == "dev"
Requires-Dist: flake8>=3.9; extra == "dev"
Dynamic: license-file

# tsysmart-tsprep - 数据预处理

时间序列数据预处理模块，提供**数据清洗**、**差异记录**、**实时处理**和**数据库建表/升级**能力。支持窄表/宽表数据读写，可通过 pip wheel 安装。

---

## 安装

```bash
# 构建 wheel
cd source/base/python/tsdataprep
python -m build --wheel --outdir dist .

# 安装
pip install dist/tsysmart_tsprep-1.0.4-py3-none-any.whl

# 验证
python -c "from tsdataprep import tsclean, tsdataprep; print('OK')"
```

## 目录结构

```
tsdataprep/
├── pyproject.toml              # 打包配置（wheel 构建）
├── __init__.py                 # 包入口
├── __main__.py                 # ts-prep CLI（API 服务入口）
├── meta.json                   # 模块元数据（名称、版本、简介、表清单）
├── install.py                  # 安装 / 卸载 / 升级（安装时自动建表）
├── upgrade.py                  # 数据库升级脚本（增量升级 / 智能重建）
├── tables.py                   # 数据库表 SQLAlchemy ORM 模型
├── requirements.txt            # Python 依赖
├── tsclean/                    # 数据清洗核心
│   ├── pipeline.py             # Pipeline 清洗流程编排
│   ├── config/                 # 清洗配置（YAML）
│   └── steps/                  # 清洗步骤（去重/异常/填充/重采样等）
├── tsdataprep/                 # 数据读写 + 实时处理
│   ├── reader.py               # read_data 数据读取
│   ├── writer.py               # write_data 数据写入
│   └── rt/                     # 实时处理模块
│       ├── cache.py            # HistoryCache 历史缓存
│       ├── queue.py            # DataQueue 数据队列
│       └── processor.py        # RealtimeProcessor 实时编排器
└── src/ts_prep/                # ts-prep API 服务（CLI 依赖）
    └── code_repo/              # 清洗业务代码库
```

---

## 数据库表清单（3 张）

| 表名 | 用途 | 是否本模块管理 |
|------|------|--------------|
| `def_src` | 数据源配置表（读写数据的入口配置）| 是 |
| `preprocess_diff` | 差异记录表（清洗/修复记录）| 是 |
| `preprocess_data` | Pipeline src_ids 模式的默认数据表 | 是 |

> 完整表定义见 [tables.py](tables.py)，三张表均由本模块管理，表结构变更走 upgrade。
> 业务数据表（窄表/宽表）不在此列，由 `def_src` 配置指定表名。

---

## 快速开始

### 1. 安装并建表

```bash
# 安装（注册模块 + 自动建表），project_code 可选：
#   指定 project_code → 建到该项目 schema
#   不指定           → 建到租户默认库（public）
python -m tsdataprep.install --tenant_code tsysmart_simu --project_code ____forecast --action install

# 查询模块信息
python -m tsdataprep.install --tenant_code tsysmart_simu --action query
```

### 2. 数据清洗（Pipeline）

```python
from tsysmart_utils import DBHelper
from tsdataprep.tsclean import Pipeline

db = DBHelper(tenant_code="tsysmart_simu", project_code="____forecast")

# 创建清洗流程（支持内置配置名 / YAML 文件 / 配置字典）
pipeline = Pipeline('basic_cleaning')

# 方式1：传入 DataFrame 直接清洗（diff 在返回结果的 attrs 中）
df = read_data(src_id='test_narrow_data', db=db)
result = pipeline.run(data=df, db=db, column_source_map={'value': 'test_narrow_cleaned'})
diff_records = result.attrs['diff_records']  # 差异记录

# 方式2：传入 src_ids 列表，自动从数据库读取
result = pipeline.run(src_ids=['src_001', 'src_002'], db=db)
```

### 3. 数据读写（窄表/宽表）

```python
from tsdataprep.tsdataprep import read_data, write_data

# 读取
df = read_data(src_id='test_narrow_data', db=db)

# 写入（自动识别窄表/宽表）
count = write_data(df, src_id='test_narrow_cleaned', db=db)
```

### 4. 实时处理（rt 模块）

```python
from tsdataprep.tsdataprep import RealtimeProcessor

def on_anomaly(src_id, row, tag):
    print(f"[告警] {src_id}: {tag}")

processor = RealtimeProcessor(
    db,
    pipeline_config='basic_cleaning',
    window_size='7D',       # 历史窗口：只保留最近 7 天
    freq='15min',           # 数据频率
    on_anomaly=on_anomaly,  # 异常告警回调
)
processor.start()  # 启动后台消费线程

# 传感器上报数据（立即返回，不阻塞）
processor.submit('src_001', 123.4, timestamp)

processor.stop()  # 优雅停止
```

> 核心机制：**HistoryCache** 首次从数据库加载历史到内存，之后直接追加，避免重复读库；**DataQueue** 异步消费，不阻塞上报。

---

## 安装 / 卸载 / 升级

通过 `tsysmart_appmng` 框架管理：

```bash
# 安装（注册模块 + 自动建表）
python -m tsdataprep.install --tenant_code tsysmart_simu --project_code ____forecast --action install

# 卸载（注销模块，不删表）
python -m tsdataprep.install --tenant_code tsysmart_simu --action uninstall

# 升级（更新注册 + 执行 upgrade.py 数据迁移）
python -m tsdataprep.install --tenant_code tsysmart_simu --project_code ____forecast --action upgrade

# 查询模块信息
python -m tsdataprep.install --tenant_code tsysmart_simu --action query
```

| 参数 | 必填 | 说明 |
|------|------|------|
| `--action` / `-a` | 是 | `install` / `uninstall` / `upgrade` / `query` |
| `--tenant_code` / `-t` | 是 | 租户代码 |
| `--project_code` | 否 | 单个项目/schema（不填则用租户默认库 public）|
| `--all_schemas` | 否 | 指定 schema 列表（空格分隔）|
| `--rebuild` | 否 | 升级时使用重建模式 |
| `--force` | 否 | 配合 rebuild 强制全部重建 |

---

## 数据库升级（upgrade.py）

独立于 `install.py` 的升级脚本，支持两种模式：

### 增量升级（默认）

对比 `tables.py` 模型与数据库中实际表结构，仅对**新增列**执行 `ALTER TABLE ADD COLUMN`。不删数据，不删列，类型变更仅报告。

```bash
# 升级指定项目
python -m tsdataprep.upgrade --tenant_code tsysmart_simu --project_code ____forecast

# 升级多个 schema
python -m tsdataprep.upgrade --tenant_code tsysmart_simu --all_schemas schema1 schema2
```

### 智能重建（--rebuild）

逐个表对比结构，**仅重建有变化的表**。流程：

```
有变化的表：
  ├── 有数据 → CREATE TABLE _upgrade_bak_{表名} AS SELECT（备份到同一 schema）
  ├── 有数据 → 导出 bakdata/{时间戳}_{tenant}_{schema}/{表名}.csv（文件备份）
  ├── DROP TABLE {表名} CASCADE
  ├── CREATE TABLE {表名}（按 tables.py 新结构）
  ├── INSERT INTO {表名} SELECT FROM _upgrade_bak_{表名}（恢复数据）
  └── DROP TABLE _upgrade_bak_{表名}（清理备份）
无变化的表 → 跳过
```

```bash
# 智能重建（自动备份 + 恢复）
python -m tsdataprep.upgrade --tenant_code tsysmart_simu --project_code ____forecast --rebuild
```

### 强制重建（--rebuild --force）

跳过结构对比，所有表全部 `DROP + CREATE`。**不备份**。

```bash
python -m tsdataprep.upgrade --tenant_code tsysmart_simu --project_code ____forecast --rebuild --force
```

### 参数一览

| 参数 | 必填 | 说明 |
|------|------|------|
| `--tenant_code` | 是 | 租户代码 |
| `--project_code` | 否 | 单个项目/schema |
| `--all_schemas` | 否 | 指定 schema 列表 |
| `--rebuild` | 否 | 启用智能重建模式 |
| `--force` | 否 | 配合 `--rebuild`，跳过智能检测，全部重建不备份 |

### 备份位置

智能重建会生成**双重备份**：

1. **PG 备份表**：`_upgrade_bak_{表名}`，与源表位于同一 schema，恢复完成后自动删除
2. **文件备份**：`bakdata/{YYYYMMDD_HHMMSS}_{tenant}_{schema}/{表名}.csv`，位于 `tsdataprep/bakdata/` 目录下，**不会自动删除**，方便追溯每次升级

---

## 清洗配置（YAML）

内置配置（`tsclean/config/defaults/`）：

| 配置名 | 说明 |
|--------|------|
| `basic_cleaning` | 基础清洗：校验 + 去重 + 异常 + 填充 |
| `full_cleaning` | 完整清洗 |
| `complete_cleaning` | 最全清洗 |
| `minimal` | 最小清洗 |
| `anomaly` | 仅异常检测 |

自定义配置：

```python
config = {
    'steps': [
        {'name': 'dedup', 'enabled': True, 'config': {'strategy': 'last'}},
        {'name': 'anomaly', 'enabled': True, 'config': {
            'method': 'zscore', 'z_threshold': 3.0, 'action': 'interpolate'
        }},
        {'name': 'fill_missing', 'enabled': True, 'config': {'strategy': 'linear'}},
    ]
}
pipeline = Pipeline(config)
```

---

## 测试

```bash
# 数据清洗测试（测试2：真实数据）
python 测试2/test_clean_data.py

# 实时处理测试（测试7）
python 测试7_实时处理/test_realtime.py

# 运行实时处理服务示例（信号处理 + 心跳 + 模拟上报）
python 测试7_实时处理/run_realtime.py
```

---

## 开发说明

### 修改表结构

1. 修改 `tables.py` 中的 ORM 模型（如新增列）
2. 更新 `meta.json` 版本号
3. 运行 `upgrade.py`（增量模式）自动 ADD COLUMN，不丢数据
4. 或运行 `upgrade.py --rebuild` 重建有变化的表（有数据自动备份恢复）

### 新增清洗步骤

1. 在 `tsclean/steps/` 下新建 `xxx.py`，实现 `def xxx_step(df, **config) -> df`
2. 在 `tsclean/steps/__init__.py` 的 `STEP_REGISTRY` 中注册
3. 在配置文件中引用该步骤名

### 依赖

| 库 | 用途 |
|----|------|
| `pandas` / `numpy` | 数据处理 |
| `sqlalchemy` | ORM / 数据库抽象 |
| `tsysmart_utils` | 数据库连接辅助（DBHelper）|
| `tsysmart_proj` | 数据库连接管理（平台依赖）|
| `tsysmart_appmng` | 模块安装管理（平台依赖）|
| `fastapi` / `uvicorn` | API 服务（src/ts_prep）|
