aaronhierrezuelo/fraud-detection-pipeline
GitHub: aaronhierrezuelo/fraud-detection-pipeline
一个端到端的交易欺诈检测流水线,涵盖数据接入、PySpark 特征工程、模型训练、基于经济收益的阈值优化以及 FastAPI 模型服务部署。
Stars: 0 | Forks: 0
# 欺诈检测 Pipeline
这是一个用于交易欺诈检测的**端到端** Pipeline,由 Airflow 进行编排,使用 PySpark 进行 feature engineering,模型通过 FastAPI 提供服务,并全部使用 Docker 进行容器化。
## 架构
```
[ingest] ─► [features (PySpark)] ─► [train] ─► [evaluate]
│ │
└──────────────── orquestado por Airflow ──────────────┘
todo levantado con Docker Compose
modelo entrenado ─► [FastAPI /predict]
```
## 技术栈
| 组件 | 工具 | 文件 |
|---|---|---|
| 编排 | Apache Airflow | `dags/fraud_pipeline_dag.py` |
| 数据接入 | pandas → Parquet | `src/ingest.py` |
| Feature engineering | PySpark | `src/features_spark.py` |
| 训练 | scikit-learn | `src/train.py` |
| 评估 (收益) | numpy | `src/evaluate.py` |
| 服务 | FastAPI | `api/main.py` |
| 容器 | Docker / Compose | `docker/`, `docker-compose.yml` |
## 结构
```
fraud-detection-pipeline/
├── dags/ # DAG de Airflow
├── src/ # config, ingest, features_spark, train, evaluate
├── api/ # FastAPI (model serving)
├── docker/ # Dockerfile.api
├── data/raw/ # fraud_data.xlsx (NO se versiona)
├── data/processed/ # parquets generados
├── models/ # modelo serializado (NO se versiona)
├── notebooks/ # EDA / experimentación
├── docker-compose.yml
└── requirements.txt
```
## 运行方式
### 选项 A — 本地运行(推荐用于逐步学习)
```
python -m venv .venv && source .venv/bin/activate # Windows: .venv\Scripts\activate
pip install -r requirements.txt
python -m src.ingest # xlsx -> parquet
python -m src.features_spark # PySpark feature engineering
python -m src.train # entrena y guarda el modelo
python -m src.evaluate # optimiza threshold por ganancia
uvicorn api.main:app --reload # API en http://localhost:8000
```
测试 API:
```
curl -X POST http://localhost:8000/predict \
-H "Content-Type: application/json" \
-d '{"features": {"a":0,"b":10,"monto":37.5,"j":"UY"}}'
```
### 选项 B — 全部使用 Docker
```
docker compose up --build
# Airflow UI: http://localhost:8080 (user/pass 在容器日志中)
# API: http://localhost:8000/health
```
在 Airflow 的 UI 中,激活并触发 DAG `fraud_detection_pipeline`。
## 结果(本地运行,平衡的 Logit 模型)
| 指标 | 数值 |
|---|---|
| ROC-AUC | 0.788 |
| PR-AUC | 0.623 |
| 最佳阈值 (按收益) | 0.36 |
| 使用模型的收益 | $51.070 |
| 全部批准的收益 (baseline) | -$34.600 |
| **相比 baseline 的提升** | **+$85.670** |
“全部批准”的 baseline 产生了**负**收益:批准欺诈交易的成本超过了合法交易的收益。该模型通过基于经济收益而非统计指标来优化截断点,从而扭转了这一局面。
## 学习路线图(建议顺序)
1. **模型优先**(`train.py` + EDA notebook) — 确保其在无基础设施的情况下正常运行。
2. **PySpark**(`features_spark.py`) — 将 feature engineering 转移到 Spark。
3. **Airflow**(`dags/`) — 编排完整的流程。
4. **FastAPI**(`api/`) — 对模型进行服务。
5. **Docker Compose** — 将所有内容组合运行。
6. **完善** — 架构图、测试(`pytest`)、扩展到大型数据集。
每个文件中都包含 `TODO` 注释,指出了为了学习建议实现的内容。
标签:Apache Airflow, Apex, AV绕过, FastAPI, PySpark, Scikit-learn, 数据工程, 机器学习, 欺诈检测, 请求拦截, 逆向工具