# 金融交易网络风险图谱分析系统 **Repository Path**: APIJSON/BankTransferNetwork ## Basic Information - **Project Name**: 金融交易网络风险图谱分析系统 - **Description**: 基于 APIJSON + Spark + Kafka 构建的金融交易网络风险图谱分析系统,针对反洗钱(AML)及团伙欺诈场景,以客户为节点、账户间转账为边,利用 Spark Core RDD 算子进行图结构计算,结合自定义累加器 AccumulatorV2 实现风险传导链路深度分析。本项目为副本,请给原仓库右上角点亮 ⭐️ Star - **Primary Language**: Unknown - **License**: MulanPSL-2.0 - **Default Branch**: master - **Homepage**: https://gitee.com/guxurui2024/BankTransferNetwork - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 1 - **Created**: 2026-08-09 - **Last Updated**: 2026-08-09 ## Categories & Tags **Categories**: Uncategorized **Tags**: APIJSON, Spark, Kafka ## README # 金融交易网络风险图谱分析系统 基于 APIJSON + Spark + Kafka 构建的金融交易网络风险图谱分析系统,针对反洗钱(AML)及团伙欺诈场景,以客户为节点、账户间转账为边,利用 Spark Core RDD 算子进行图结构计算,结合自定义累加器 AccumulatorV2 实现风险传导链路深度分析。 ## 系统架构 ``` ┌─────────────────────────────────────────────────────────────────┐ │ 模块一:数据采集 │ │ Java 数据模拟器 → Kafka (bank_transfer_network) │ │ 生成 12000+ 条转账记录,模拟 500 个账户,含三类风险模式 │ └──────────────────────┬──────────────────────────────────────────┘ │ Kafka Topic ▼ ┌─────────────────────────────────────────────────────────────────┐ │ 模块二:数据处理 (Spark Core + Spark SQL) │ │ Kafka 消费 → RDD 算子链 → 自定义累加器 → 风险检测 → MySQL │ │ map/flatMap/filter/groupByKey/join/cogroup/reduceByKey │ │ AccumulatorV2 / 广播变量 / MEMORY_AND_DISK 持久化 │ └──────────────────────┬──────────────────────────────────────────┘ │ MySQL (risk_graph) ▼ ┌─────────────────────────────────────────────────────────────────┐ │ 模块三:数据可视化 (ECharts + APIJSON) │ │ APIJSON 后端 (:8080/get) ← Axios ← ECharts 前端 │ │ 力导向图 / 仪表盘+柱状图 / 桑基图 / 可排序表格 │ └─────────────────────────────────────────────────────────────────┘ ``` ## 技术栈 | 组件 | 版本 | 用途 | |------|------|------| | APIJSON | 6.1.0 | 腾讯开源 JSON ORM 框架,后端通用接口 | | Spring Boot | 2.6.4 | Web 框架 | | Spark | 3.3.2 (Scala 2.12) | 大数据处理引擎 | | Kafka | 3.3.1 | 消息队列 | | ECharts | 5.4.3 | 前端数据可视化 | | MySQL | 5.7+ | 数据存储 | | FastJSON | 2.0.43 | JSON 解析 | | Java | 1.8 | 运行环境 | ## 项目结构 ``` APIJSONDemo/ ├── pom.xml # Maven 构建配置(含 Spark/Kafka/Scala 依赖) │ ├── src/main/java/apijson/demo/ │ ├── DemoApplication.java # SpringBoot 启动类(端口 8080) │ ├── DemoController.java # APIJSON 通用接口控制器 │ ├── DemoSQLConfig.java # 数据库配置(risk_graph) │ └── simulator/ │ └── TransferDataSimulator.java # 【模块一】数据采集模拟器(Java→Kafka) │ ├── src/main/scala/apijson/demo/spark/ │ ├── TransferRecord.scala # 转账记录样例类 │ ├── TransferStat.scala # 账户交易统计样例类 │ ├── Edge.scala # 转账关系边样例类 │ ├── RiskIndicatorAccumulator.scala # 【模块二】自定义 AccumulatorV2 累加器 │ └── RiskGraphAnalysis.scala # 【模块二】Spark 核心分析作业 │ ├── src/main/resources/ │ └── application.properties # SpringBoot 配置 │ ├── sql/ │ ├── risk_graph_tables.sql # 风险图谱建表 + 示例数据 │ ├── register_risk_tables.sql # APIJSON access 表注册 │ └── sys.sql # APIJSON 系统表 │ ├── frontend/static/ │ ├── index.html # 【模块三】大屏入口(侧边导航) │ └── visualization/ │ ├── network_overview.html # 网络总览(力导向图) │ ├── risk_monitor.html # 风险指标监控(仪表盘+柱状图) │ ├── fund_tracing.html # 资金链路追踪(桑基图) │ └── high_risk_accounts.html # 高风险账户列表(可排序表格) │ ├── libs/ # APIJSON 本地 JAR 依赖 ``` ## 模块详解 ### 模块一:数据采集(Java → Kafka) **文件**:`src/main/java/apijson/demo/simulator/TransferDataSimulator.java` **功能**: - 生成 12000 条账户转账记录,模拟 500 个账户之间的转账关系网络 - 通过 Kafka Producer 发送到 Topic:`bank_transfer_network` - 使用 `from_account` 作为 key,保证同账户交易有序 **模拟三类风险模式**: | 模式 | 占比 | 说明 | |------|------|------| | 资金闭环 | ~5% | A→B→C→D→A,资金在多个账户间循环流转(3~4阶闭环) | | 资金归集 | ~8% | 多个分散账户在短时间内向同一目标账户转入资金 | | 大额拆分 | ~5% | 一笔大额资金入境后拆分为多笔小额转出(1小时内≥5笔,单笔<2万) | | 正常转账 | ~82% | 日常消费、工资代发、还贷等正当用途 | **运行方式**: ```bash # 确保 Kafka 已启动(localhost:9092) # 在 IDE 中直接运行 TransferDataSimulator.main() ``` ### 模块二:数据处理(Spark Core + Spark SQL) **文件**:`src/main/scala/apijson/demo/spark/RiskGraphAnalysis.scala` **核心技术点**: 1. **RDD 转换算子链**: - `map`:TransferRecord → (account, TransferStat) 键值对 - `groupByKey`:按 account 分组所有转入/转出记录 - `flatMap` + `filter`:统计独特交易对手数、交易频次、金额均值/标准差 - `reduceByKey`:构建 (from, to) 聚合边 - `join`:构建二阶关联路径 - `cogroup`:资金拆分检测 2. **自定义 AccumulatorV2**: - `RiskIndicatorAccumulator` 继承 `AccumulatorV2[String, Map[String, Long]]` - 统计风险指标:大额转账、高频转账、资金闭环、资金归集、大额拆分 3. **广播变量**: - 广播风险检测阈值,避免每个任务传递配置 4. **RDD 持久化**: - `persist(StorageLevel.MEMORY_AND_DISK)` 减少重复计算 5. **风险检测算法**: - **资金闭环检测**:join 构建 2~4 阶路径,检测起点=终点的闭环 - **归集账户识别**:统计入度 > 20 的高归集账户 - **资金拆分检测**:单笔转入 > 10万 且 1小时内拆分 ≥ 5笔小额转出(< 2万) 6. **Spark SQL 写入 MySQL**: - 分析结果写入 `risk_edge`、`risk_account`、`risk_indicator`、`fund_flow_path` 表 **运行方式**: ```bash spark-submit --class apijson.demo.spark.RiskGraphAnalysis \ --jars mysql-connector-java-8.0.29.jar \ target/apijson-demo-6.1.0.jar ``` ### 模块三:数据可视化(ECharts + APIJSON) | 页面 | 图表类型 | 功能说明 | |------|----------|----------| | 网络总览 | 力导向图 | 全量账户转账关系网络,节点按金额加权,颜色按风险编码,支持点击查看详情 | | 风险指标监控 | 仪表盘 + 柱状图 | 四类风险指标实时计数仪表盘 + 各类型命中账户数对比 | | 资金链路追踪 | 桑基图 | 资金流转路径,展示源头到终点流向与流量,支持按类型筛选 | | 高风险账户列表 | 可排序表格 | 账户明细表(入度/出度/金额/标签),支持排序筛选搜索 | **前端访问**:浏览器打开 `frontend/static/index.html` ## 快速开始 ### 1. 环境准备 - JDK 1.8 - Maven 3.x - MySQL 5.7+ - Kafka 3.x(可选,仅模块一/二需要) - Spark 3.3.2(可选,仅模块二需要) - IntelliJ IDEA(推荐) ### 2. 数据库初始化 ```sql -- 1. 创建风险图谱数据库和业务表 SOURCE sql/risk_graph_tables.sql; -- 2. 导入 APIJSON 系统表(首次需要) SOURCE sql/sys.sql; -- 3. 注册业务表到 APIJSON access 表 SOURCE sql/register_risk_tables.sql; ``` ### 3. 启动后端服务 ```bash # 修改 DemoSQLConfig.java 中的数据库连接信息(如需要) # 启动 SpringBoot 应用 mvn spring-boot:run ``` ### 4. 访问可视化大屏 浏览器打开 `frontend/static/index.html`,通过侧边导航切换四个可视化页面。 ### 5.(可选)运行数据模拟和 Spark 分析 ```bash # 模块一:启动 Kafka,运行数据模拟器 # 在 IDE 中运行 TransferDataSimulator.main() # 模块二:运行 Spark 分析作业 spark-submit --class apijson.demo.spark.RiskGraphAnalysis \ --jars mysql-connector-java-8.0.29.jar \ target/apijson-demo-6.1.0.jar ``` ## 数据库表说明 | 表名 | 说明 | 数据来源 | |------|------|----------| | `transfer_record` | 转账记录表 | Kafka 原始数据 | | `risk_edge` | 风险关系图谱边表 | Spark 分析结果 | | `risk_account` | 账户风险信息表 | Spark 分析结果 | | `risk_indicator` | 风险指标统计表 | Spark 累加器结果 | | `fund_flow_path` | 资金流转路径表 | Spark 分析结果 | ## License Apache License 2.0