yassine140903/Real-Time-User-Behavior-Anomaly-Detection-Platform

GitHub: yassine140903/Real-Time-User-Behavior-Anomaly-Detection-Platform

面向银行柜台业务的实时欺诈与反洗钱异常检测平台,融合双模型ML评分与监管规则引擎,提供可解释的告警队列和完整的MLOps流水线。

Stars: 0 | Forks: 0

# 实时用户行为异常检测平台 这是一个为 Amen Bank 分行(柜台)业务构建的实时欺诈与行为异常检测平台,涵盖四种核心柜员交易类型:retrait(取款)、versement(存款)、virement(转账)和 remise de chèques(支票存款)。该平台专为分行主管和 AML(反洗钱)/合规职能部门设计:它对每一笔发生的交易进行评分,将机器学习异常信号与严格的监管规则相融合,并呈现一个经过排序、可解释的告警队列,而不是满屏的原始交易日志。它针对三个相互重叠的风险类别——经典的交易欺诈、突尼斯银行业监管规定(BCT、CTAF、第41-2024号法律)定义的 AML 拆分/洗钱模式,以及客户或员工日常操作方式中较轻的行为偏移——并将客户和员工视为两个独立的参与者,对其行为进行独立画像和评分。 ## 架构概述 交易流经一个基于 Kafka 兼容 topic 构建的四阶段流式 pipeline:`raw-events` → `enriched-events` → `scored-events` → `decisions`。模拟器(在生产环境中为真实的 core-banking 数据源)将原始柜员事件发布到 `raw-events`。enrichment service 是唯一允许读写行为状态的组件——它从 Redis 中提取客户和员工的“热”画像,计算上下文特征(针对原型基线的 z-score、滚动计数、累计金额、临近阈值标志、重复检测),并重新发布一个丰富后的事件。scoring service 消费丰富后的事件,并融合两个模型:一个是对单事件重建误差进行评分的自编码器(AE),另一个是对偏离客户近期交易历史的行为进行评分的 LSTM 序列模型。由于全新的客户没有可用于对比的序列,这两个分数通过与账户年龄挂钩的 sigmoid 加权平均值进行合并——新账户几乎完全依赖于 AE,成熟账户则越来越倾向于 LSTM——而不是使用固定的权重。每个分数也会转换为相对于参考分布的百分位数,这样原本处于不同原始尺度的 AE 和 LSTM 在融合前就变得具有可比性。在计算融合分数的同时还会计算 SHAP 值,因此每个告警都会附带一份促发该告警的特征排序列表。decision service 是最后一个无状态阶段:它将融合后的分数映射到一个告警层级,然后在此基础上叠加确定性的 AML/监管规则和风险历史调整,并且只会提升层级——规则可以将层级上调,绝不会人为压制真实的 ML 驱动信号。最终会输出四个层级:INFO(无需操作)、REVIEW(排队等待主管处理)、ALERT(高风险,在队列中优先处理)和 BLOCK(强制阻断,立即上报)。两个 sink consumer 独立且并行地持久化 pipeline 的输出——将原始事件存入 `transactions`,将决策存入 `alerts`——同时画像状态在 Redis 中保持“热”状态以实现毫秒级延迟读取,而持久化的“冷”副本则存放在 PostgreSQL 中,由批处理作业每晚刷新。 ## 技术栈 | 类别 | 技术 | |---|---| | 语言 | Python 3.11 | | 流处理 | Redpanda (Kafka API), kafka-python | | ML | PyTorch (Autoencoder + LSTM), scikit-learn, SHAP | | 特征存储 | Redis (热画像、缓冲区、原型基线) | | 数据库 | PostgreSQL 17 (冷画像、交易、告警、评估标签) | | 编排 | Apache Airflow (LocalExecutor, DockerOperator) | | MLOps | MLflow (tracking, model registry) | | 监控 | Prometheus, Alertmanager, kafka-exporter | | 仪表板 | Streamlit, Plotly | | API | FastAPI, Uvicorn | | 容器化 | Docker Compose (20 个服务) | ## 项目结构 ``` . ├── docker-compose.yml # 20-service orchestration (infra, pipeline, MLOps, dashboard) ├── Dockerfile.api # streaming services, API, batch jobs, simulator ├── Dockerfile.airflow # Airflow webserver/scheduler ├── Dockerfile.dashboard # Streamlit supervisor dashboard ├── Dockerfile.webhook # Prometheus→Airflow drift-retrain relay ├── config/ │ ├── archetype.yaml # 5 client archetypes (population mix, op mix, amount/frequency priors) │ ├── init.sql # PostgreSQL schema (clients, transactions, alerts, snapshots) │ ├── prometheus.yml / alert_rules.yml # scrape config + drift/health alert rules │ └── alertmanager.yml # routes drift alerts to the webhook relay ├── dags/ │ ├── nightly_profiles.py # batch profile/baseline refresh → hydrate Redis │ ├── weekly_retrain.py # scheduled retrain (also triggered by drift/feedback) │ └── feedback_monitor.py # daily supervisor-rejection-rate check ├── src/ │ ├── generator/ # synthetic client population + anomaly injection │ ├── streaming/ # simulator, enrichment/scoring/decision services, sinks, DLQ │ ├── enrichment/ # EnrichmentCore — feature computation (pure, no I/O side effects) │ ├── scoring/ # ScoreFusion (AE + LSTM), SHAP explanation core │ ├── decision/ # DecisionService — tiering, regulatory rules, explainer │ ├── training/ # model definitions + train_all entrypoint (MLflow-tracked) │ ├── batch/ # nightly profile/baseline/employee-profile jobs, hydrate_redis │ ├── dashboard/ # Streamlit app (Alerts, KPIs, Simulation, Health pages) │ ├── api/ # FastAPI simulation-control API │ ├── webhook/ # Alertmanager → Airflow retrain relay │ └── monitoring/ # Prometheus metrics definitions ├── models/ # trained artifacts (autoencoder.pt, lstm.pt, scaler, reference scores) ├── scripts/ # one-off/maintenance scripts (seeding, calibration, sanity checks) └── tests/ # decision, fusion, and SHAP end-to-end tests ``` ## 服务 | 服务 | 角色 | 端口 | |---|---|---| | **基础设施** | | | | redpanda | Kafka 兼容的事件代理 | 9092, 8082, 9644 | | redis | 热 profile/buffer/baseline 存储 | 6379 | | postgres | 冷存储:客户、交易、告警、快照 | 5432 | | **流式 Pipeline** | | | | simulator | 将合成的柜员事件发布到 `raw-events` | — | | enrichment-service | 计算上下文特征、热画像查找 | 9100 | | scoring-service | AE + LSTM 融合、SHAP 解释 | 9101 | | decision-service | 层级划分、监管规则、告警组装 | 9102 | | **Sink Consumer** | | | | raw-events-sink | 将原始事件持久化到 `transactions` | 9103 | | decisions-sink | 将决策持久化到 `alerts`,更新 Redis 风险历史 | 9104 | | **批处理 / 编排** | | | | hydrate-redis | 将冷画像/基线从 Postgres 加载到 Redis | — | | db-seed | 将合成的 CSV 数据加载到 Postgres | — | | airflow-init | Airflow 数据库迁移 + 管理员用户引导 | — | | airflow-webserver | Airflow UI | 8080 | | airflow-scheduler | 运行 `nightly_profiles`, `weekly_retrain`, `feedback_monitor` DAG | — | | **MLOps** | | | | training | 训练 AE + LSTM,将运行/制品记录到 MLflow | — | | mlflow | 实验跟踪 + 模型注册中心 UI | 5001 | | **监控** | | | | prometheus | 指标抓取(pipeline 延迟、吞吐量、错误) | 9090 | | alertmanager | 路由漂移/健康告警 | 9093 | | kafka-exporter | 用于 Prometheus 的 Redpanda/Kafka 指标 | 9308 | | webhook-relay | 将漂移告警转发以触发 Airflow 重训练 | — | | **面向用户** | | | | dashboard | Streamlit 主管仪表板 | 8501 | | simulation-api | 用于模拟器的 FastAPI 控制平面 | 8000 | ## 入门指南 **前置条件:** Docker 和 Docker Compose。 ``` # 1. Clone 仓库 git clone cd Real-Time-User-Behavior-Anomaly-Detection-Platform # 2. Build 所有镜像 docker compose build # 3. 从已提交的 CSVs 填充 PostgreSQL docker compose --profile seed up db-seed # 4. 启动整个平台(hydrate-redis 会在 pipeline 服务启动前自动运行) docker compose up -d # 5. 运行模拟器以通过 pipeline 流式传输事件 docker compose --profile simulate run --rm simulator ``` 仪表板可通过 `http://localhost:8501` 访问,MLflow 在 `http://localhost:5001`,Airflow 在 `http://localhost:8080`(admin/admin),Prometheus 在 `http://localhost:9090`。 **启动时发生的事情:** Docker Compose 启动基础设施(RedPanda、Redis、PostgreSQL),然后 `hydrate-redis` 将客户/员工画像从 PostgreSQL 加载到 Redis(不预加载序列缓冲区——`--skip-buffers` 标志已内置在 compose 命令中)。hydration 完成后,流式 pipeline 服务开始消费。模拟器会重放过去 10 天的合成交易,以便缓冲区自然建立,这代表了正在查看仪表板的主管实际会看到的情况。 ## 关键设计决策 **每个阶段使用单一 Kafka topic,按 client_id 分区。** 所有四种操作类型在每个 pipeline 阶段都流经同一个 topic,而不是每种操作类型一个 topic。这保留了每个分区内跨操作的时间顺序——这对于 LSTM 序列模型至关重要,因为它需要按正确的顺序查看客户的 retrait 以及随后的 versement。client_id 分区键保证了按客户排序,同时仍然允许消费者进行水平扩展。 **EnrichmentCore 没有 I/O 副作用。** 所有 Redis/Postgres 访问都发生在包装服务中;`EnrichmentCore.enrich_event()` 接收已经获取的 profile 和 buffer,并返回一个纯粹的特征字典。这使得特征计算可以独立测试,并可在流式服务和用于训练数据的批量回填路径之间重用。 **规则只会升级,永远不会降级。** decision service 的四个层(分数评级 → 监管规则 → 风险历史调整 → 解释)一旦分配,就只能提升层级。真正的 ML 驱动信号永远不会被规则悄无声息地降级。 **`alerts` 到 `transactions` 没有外键。** 这是有意为之,不是疏忽:raw-events sink 和 decisions sink 并行地从不同的 topic 消费,因此一个决策完全有可能在其对应的原始事件之前到达 Postgres。引用完整性由 pipeline 拓扑保证,而不是数据库约束。 **三层冷启动评分。** 序列缓冲区中事件少于 5 个的客户将获得 0.5 的中性融合分数,并带有 cold_start=True 标志——模型完全不会被调用。decision service 仍然会运行所有监管规则(即使是新客户,拆分模式也是可疑的),但将最终层级最高限制为 REVIEW。在 5 到 9 个事件之间,只有 AE 评分(LSTM 权重强制为 0)。在达到 10 个或以上事件时,完整的 AE+LSTM 融合将根据账户成熟度使用 sigmoid 加权。 **带有失败分类的 DLQ。** 每个流式消费者都将 `process_fn` 包装在一个弹性处理器中,该处理器区分瞬时基础设施错误(Redis/Postgres 连接问题——通过指数退避进行重试)和毒丸事件(格式错误的 payload——直接发送到 `dlq` topic 并附带完整的堆栈跟踪和原始事件,不会浪费任何重试)。 **三个独立的重训练触发器。** 模型在每周的 cron 计划上、Prometheus/Alertmanager 检测到的统计漂移上(通过 webhook 中继路由到 Airflow DAG 运行),以及主管反馈出现分歧时(7 天内 ALERT/BLOCK 的拒绝率超过 50%,最少 10 次审查)进行重新训练——这样模型就可以对缓慢的漂移和与人类审查者的突然分歧做出反应,而不仅仅是根据日历。 ## MLOps Pipeline 重训练通过三种方式触发:每周日的 cron(`weekly_retrain` DAG)、通过 webhook 中继转发到 Airflow 的 Prometheus 检测到的漂移告警,以及当主管拒绝了超过一半的 ALERT/BLOCK 告警时触发的每日反馈分歧检查。每个触发器运行的都是同一个 DAG——批量重新富化、特征完整性/NaN 率验证门控,然后是一个容器化的训练运行(`train_all`),将指标、参数和 artifact 记录到 MLflow。提升是自动化的:训练脚本将新的 AUC 与 models/production_metrics.json 进行比较,仅当新的 AUC 达到或超过当前值时,才会覆盖生产模型文件。MLflow 的作用是事后剖析——当提升门控在凌晨 4 点拒绝了一个模型时,工程师可以在周一检查完整的运行情况。 ## 仪表板 位于 `http://localhost:8501` 的 Streamlit 主管仪表板每 10 秒自动刷新一次,并分为四个页面组织:**告警**(带有 SHAP 解释和监管标志的实时、分层队列)、**KPI**(分支和员工级别的业务量、告警率和层级分布视图)、**模拟**(事件模拟器的控制界面)和**健康**(从 Prometheus 提取的 pipeline 延迟、吞吐量和 DLQ 状态)。主管直接在告警页面对告警采取行动——确认或拒绝告警会写回到 `alerts.supervisor_decision`,这会反馈给反馈分歧重训练触发器,从而闭合人工审查和模型重训练之间的循环。 ## 异常场景 合成数据生成器注入了跨越三个难度级别的十三个带标签的异常场景,既用于验证检测,也用于驱动实时模拟演示。 | 场景 | 难度 | 参与者 | 描述 | |---|---|---|---| | amount_spike | 简单 | 客户 | 远高于客户通常金额的单笔交易 | | round_amount | 简单 | 客户 | 可疑的整数交易金额(5000、10000、15000、20000 DT) | | duplicate_virement | 简单 | 客户 | 在短时间内重复的几乎相同的转账 | | volume_spike | 简单 | 员工 | 员工处理的交易量突然激增 | | frequency_burst | 中等 | 客户 | 交易次数相对于客户预期频率的突然激增 | | cumulative_threshold | 中等 | 客户 | 累计接近监管阈值的多笔交易 | | timing_anomaly | 中等 | 客户 | 在正常营业时间之外进行的交易 | | client_concentration | 中等 | 员工 | 异常集中的客户-员工配对 | | smurfing | 困难 | 客户 | 持续将存款拆分至刚好低于 10,000 DT 的报告阈值 | | cheque_structuring | 困难 | 客户 | 支票金额保持在刚好低于 30,000 DT 的法定上限(第41-2024号法律) | | money_mule | 困难 | 客户 | 被用作快速进出转账的过渡账户 | | behavioral_drift | 困难 | 客户 | 几周内交易模式的逐渐转变 | | cross_actor | 困难 | 交叉 | 客户与员工之间协调一致的异常行为 | ## 模型性能 | 模型 | AUC | | |---|---|---| | Autoencoder | 0.9714 | 单事件重建误差 | | LSTM | 0.9355 | 序列偏差分数,仅限成熟客户 | 这些数据来自于对基于原型的模拟器生成的合成数据进行训练和验证,而非生产交易历史。它们证明了融合架构和特征设计能够将注入的异常与正常行为分离开来——这并不是对真实世界检测性能的声明,后者需要针对带标签的生产数据来确立。 ## 生产路线图 要将其从 Docker Compose 演示转移到生产环境,需要:使用 Kubernetes 进行编排和自动缩放以取代 Compose;一个 CI/CD pipeline 用于自动化测试、镜像构建和门控模型提升;使用真正的密钥管理器(Vault 或云 KMS)来代替 `.env` 文件凭证;进行负载测试以验证在现实的分支网络交易量下的吞吐量;为交易、告警和画像快照制定明确的数据保留/归档策略;以及使用集中式日志聚合(ELK 或 Loki)来代替每个容器的日志。 作为在 Amen Bank, DCSI 实习的一部分而构建。
标签:Autoencoder, Kubernetes, LSTM, MLOps, PMD, 凭据扫描, 反洗钱, 实时流处理, 异常检测, 搜索引擎查询, 测试用例, 版权保护, 自定义请求头, 逆向工具, 金融风控