iHazelly/Banking-Transaction-Data-Pipeline-
GitHub: iHazelly/Banking-Transaction-Data-Pipeline-
基于AWS的事件驱动无服务器数据管道,模拟银行欺诈检测场景,实现从数据接收到可视化全流程自动化。
Stars: 0 | Forks: 0
# 🏦 银行交易数据管道 — AWS 生产级





一个**事件驱动、端到端的数据管道**,模拟真实世界的银行欺诈检测系统。
它接收交易数据 (PaySim),执行数据质量检查,使用 **AWS Glue (PySpark)** 转换和分区数据,将其存储为 **Parquet**,通过 **Amazon Athena** 使其可查询,并通过**交互式静态仪表板**可视化关键指标——全部自动化且无服务器。
## 📌 业务问题
银行每天处理数百万笔交易。传统的批处理作业在午夜运行,需要耗费数小时。如果文件损坏或作业失败,直到早上才会有人知道 → 导致欺诈检测延迟和财务损失。
**本项目通过构建一个自动化管道来解决这个问题,该管道能够:**
- 当新 CSV 文件到达 S3 时**立即触发**。
- 在处理之前**验证数据质量**。
- 将数据**转换并分区**为列式 Parquet 格式(查询速度提升 10 倍)。
- **记录数据血缘**并通过电子邮件发送**成功/失败警报**。
- 通过交互式仪表板**提供实时洞察**,该仪表板在每次运行后自动更新。
## 🏗️ 架构概述
*高级架构 – 使用 AWS 服务的事件驱动管道。*
| 层级 | 服务 | 角色 |
|-------|---------|------|
| **接收** | S3 | 原始 CSV 文件存入 `/incoming/` |
| **编排** | EventBridge → Step Functions | 在文件上传时触发管道,管理重试 |
| **模式检测** | Glue Crawler | 自动推断原始数据的模式 |
| **转换** | Glue ETL (PySpark) | 数据质量检查、清洗、丰富、分区、Parquet 转换 |
| **存储** | S3 | 分区的 Parquet 文件 (`year/month/day/type`) |
| **查询** | Amazon Athena | 在 Parquet 上进行无服务器 SQL 分析 |
| **监控** | CloudWatch + SNS | 日志记录、指标、关于成功/失败的电子邮件警报 |
| **可视化** | S3 静态网站 (ApexCharts) | 具有自动刷新功能的交互式仪表板 |
## 🛠️ 使用的技术与服务
| 类别 | 工具 / 服务 |
|----------|------------------|
| **云** | AWS (S3, Glue, Athena, Step Functions, EventBridge, SNS, CloudWatch, IAM) |
| **数据处理** | PySpark (AWS Glue 4.0), Parquet, Snappy 压缩 |
| **查询引擎** | Amazon Athena (Presto) |
| **编排** | Step Functions, EventBridge |
| **监控与警报** | CloudWatch, SNS (电子邮件) |
| **可视化** | HTML, Tailwind CSS, ApexCharts.js |
| **基础设施即代码** | AWS CloudFormation |
| **版本控制** | Git, GitHub |
## 📊 数据集
- **来源**: [PaySim Financial Dataset](https://www.kaggle.com/datasets/ealaxi/paysim1) (Kaggle)
- **行数**: 6.3M(使用的样本:10k 用于测试)
- **关键列**: `step` (时间), `type` (TRANSFER, CASH_OUT 等), `amount`, `nameOrig`, `nameDest`, `oldbalanceOrg`, `newbalanceOrig`, `isFraud`
- **欺诈分布**: 只有 `TRANSFER` 和 `CASH_OUT` 包含欺诈交易。
## 🧠 管道的作用(分步说明)
1. **文件上传** → CSV 文件被放入 `s3://banking-raw-.../incoming/`
2. **EventBridge** → 检测到 S3 `PutObject` 事件并触发 Step Functions。
3. **Step Functions** → 编排:
- 运行 Glue Crawler(模式更新)
- 运行 Glue ETL 作业(转换并加载)
- 在 Athena 中运行 `MSCK REPAIR TABLE`
- 通过 SNS 发送成功/失败电子邮件
4. **Glue ETL (PySpark)** 执行以下操作:
- 数据质量检查(空值、amount>0、有效类型、余额一致性)
- 数据丰富(账户类型、可疑标志、时间特征)
- 将数据写入按 `year/month/day/type` 分区的 **Parquet**
- 将**数据血缘** (JSON) 记录到 S3 血缘 bucket 中
- 将**摘要 JSON** 导出到仪表板 bucket (`data/summary.json`)
5. **Athena** → 在 Parquet 位置上的外部表,准备好进行 SQL 查询。
6. **静态仪表板** (S3 网站) → 获取 `summary.json` 并渲染:
- KPI 卡片(总交易量、欺诈案例、欺诈率、总交易额)
- 条形图(按交易类型划分的欺诈情况)
- 饼图(交易量份额)
- 趋势的 Sparkline
- 最后更新时间戳
## 🚀 部署说明 (CloudFormation)
### 前提条件
- 具有权限的 AWS 账户(或 AWS Learner Lab)。
- 已安装 AWS CLI(可选,您可以使用 Console)。
### 步骤
**1. 克隆存储库**
```
git clone https://github.com/yourusername/banking-pipeline.git
cd banking-pipeline
```
**2. 将 Glue 脚本上传到 S3**
```
aws s3 cp glue/banking_etl_job.py s3://banking-raw-/scripts/
```
**3. 部署 CloudFormation 栈**
- 转到 AWS Console → CloudFormation → Create stack → Upload a template file → 选择 `cloudformation/banking-pipeline-iac.yaml`
- 提供参数:
- `EnvironmentName`: `dev`
- `BucketSuffix`: 唯一的后缀(例如,`yourname-2026`)
- `EmailAddress`: 您的电子邮件(用于接收警报)
- `GlueScriptS3Key`: `scripts/banking_etl_job.py`
- 确认 IAM 功能 → Create stack
- 等待 `CREATE_COMPLETE`(约 2-3 分钟)
**4. 上传测试数据**
将 `dataset/sample_10k.csv` 上传到您的原始 bucket:
```
aws s3 cp dataset/sample_10k.csv s3://banking-raw-/incoming/
```
**5. 监控管道**
- Step Functions 控制台 → 观察执行图表
- Glue 控制台 → Jobs → 查看运行日志
- 检查电子邮件以获取成功/失败通知
**6. 启用仪表板**
- 转到 S3 控制台 → `banking-dashboard-` bucket
- 属性 → 静态网站托管 → 启用 → 索引文档:`index.html`
- 将 `dashboard/index.html` 上传到 bucket 根目录
- 打开网站 URL — 一旦 Glue 作业生成 `data/summary.json`,仪表板就会显示实时数据
## 🔐 安全与成本优化
- **加密**: 在所有 bucket 上启用 SSE-S3
- **公共访问**: 所有 bucket 都已阻止公共访问,仪表板 bucket 除外(静态网站)
- **IAM**: 最小权限角色(此实现中使用 LabRole)
- **分区**: 数据按 `year/month/day/type` 分区 → 显著降低 Athena 扫描成本
- **生命周期策略**: 处理后的数据在 30 天后移动到 `STANDARD_IA`,在 365 天后过期
- **成本**: 完整的开发和测试管道总成本不到 5 美元(AWS Learner Lab 40 美元额度)
## 🧠 展示的技能
| 技能 | 证据 |
|---|---|
| 数据管道 (ETL/ELT) | Glue PySpark 脚本 — 提取、转换、加载 |
| 数据质量 | 空值检查、范围验证、业务规则验证 |
| 数据血缘 | 记录源、转换、目标的 JSON 日志 |
| 分区与列式格式 | Parquet + Snappy,Hive 风格分区 |
| 无服务器编排 | Step Functions + Lambda 事件触发器 |
| 监控与警报 | CloudWatch 日志 + SNS 电子邮件警报 |
| 数据湖与查询 | S3 上的 Athena 外部表 |
| 仪表板与可视化 | HTML + ApexCharts,S3 静态托管 |
| 基础设施即代码 | CloudFormation (YAML) |
| 安全与成本意识 | Bucket 策略、加密、生命周期规则 |
## 🔮 未来改进
- **实时流处理** — 用 Amazon Kinesis + Lambda 替换批处理,实现近乎瞬时的欺诈检测
- **高级数据质量** — 集成 Great Expectations 或 Deequ 以进行自动化 DQ 报告
- **ML 集成** — 使用 SageMaker 预测欺诈概率并在仪表板中展示
- **CI/CD 管道** — 使用 GitHub Actions 自动化部署
- **图分析** — 使用 Neo4j 进行循环转账环检测(洗钱模式)
## 🤖 AI 使用声明
在开发此项目期间,AI 工具被用于:
- 语言翻译和句子润色
- 代码建议、调试和结构指导
- 为 README 和文档提供写作协助
- 头脑风暴和概念支持
但是,所有核心架构决策、数据建模、特征工程、管道配置、结果解释和最终的技术验证均由作者 (Paradorn Khanongsuwan) 完成。所有 AI 生成的输出均已过人工验证和调整。
## 📜 引用
```
@misc{khanongsuwan_2026_banking_pipeline,
title={Banking Transaction Data Pipeline – Production-Grade AWS ETL with Real-time Dashboard},
author={Khanongsuwan, Paradorn},
year={2026},
howpublished={\url{https://github.com/iHazelly/Banking-Transaction-Data-Pipeline-}}
}
```
## 📬 联系方式
- **GitHub**: github.com/iHazelly
欢迎提出 issue 或 pull request 以进行改进!
## ✅ 总结
本项目是一个完整的、生产级的数据管道,展示了在 AWS 上进行现代数据工程的各个阶段——从接收数据到可视化。旨在实现可重现、高性价比且可用于作品集展示。非常适合展示**数据工程师**、**分析工程师**或**数据平台工程师**角色所需的技能,尤其是在银行和金融领域。
标签:AWS, C语言, DPI, PySpark, 数据工程, 数据管道, 网络调试, 自动化, 软件工程, 逆向工具