SharathHL99/Fraud_Detection_Rule_Engine
GitHub: SharathHL99/Fraud_Detection_Rule_Engine
一个基于 Spring Boot 的异步批量数据导入与校验系统,通过流式处理和幂等策略高效地将大型 CSV/Excel 文件持久化到数据库。
Stars: 0 | Forks: 0
## 批量数据导入与校验系统 – 高级版
这是一个 Spring Boot 后端服务,通过使用流式处理、异步处理和批量插入,高效地上传、校验、转换和持久化大型 CSV/Excel 数据集。
## 目录
概述
技术栈
架构
数据库 Schema
项目结构
前置条件
设置说明
配置
运行应用程序
API Endpoints
请求/响应示例
输入文件示例
测试
处理大文件与边缘情况
性能优化
幂等性策略
故障排除
## 概述
该系统允许企业用户上传包含批量数据的大型 CSV/Excel 文件。每条记录都会进行流式读取,使用 Bean Validation 进行校验,并以分批的方式持久化。系统会跟踪导入任务,并提供成功和失败记录的摘要,而无需将整个文件加载到内存中。
核心功能
通过 REST API 上传 CSV/Excel 文件
流式解析并校验每条记录
持久化有效记录;拒绝无效记录并记录错误原因
异步任务处理(非阻塞式上传响应)
数据库批量插入以提升性能
幂等性文件处理(检测并跳过重复上传)
导入摘要报告(成功计数、失败计数、每条记录的错误信息)
## 技术栈
层级 技术
编程语言 Java 17
框架 Spring Boot 3.x
持久层 Spring Data JPA (Hibernate)
数据库 MySQL / PostgreSQL
校验 Jakarta Bean Validation
文件解析 Apache Commons CSV / Apache POI(针对 Excel 使用流式 SXSSF/SAX)
异步处理 Spring `@Async` + `ThreadPoolTaskExecutor`
构建工具 Maven
API 文档 Swagger / springdoc-openapi
测试 JUnit 5, Mockito, Testcontainers
架构
```
Client
│
▼
Controller Layer (ImportController)
│ - Accepts file upload
│ - Creates ImportJob record (status: PENDING)
│ - Returns jobId immediately
▼
Service Layer (ImportService) \[runs @Async]
│ - Streams file (row by row)
│ - Validates each record (Bean Validation)
│ - Buffers valid records into batches
│ - Persists batches (JPA batch insert)
│ - Logs invalid records with error\_message
│ - Updates ImportJob status (PROCESSING → COMPLETED/FAILED/PARTIALLY\_COMPLETED)
▼
Repository Layer (Spring Data JPA)
│
▼
Database (ImportJob, ImportRecord tables)
```
## 分层结构
Controller – REST endpoints,请求/响应 DTO
Service – 业务逻辑,编排解析 + 校验 + 持久化
Parser/Strategy – 可插拔的 CSV/Excel 读取器(基于文件类型的策略模式)
Validator – Bean Validation + 自定义业务规则检查
Repository – Spring Data JPA repositories
Async Executor Config – 专用于导入任务的线程池
## 数据库 Schema
`import\_job`
列名 类型 备注
id BIGINT (PK) 自动生成
file_name VARCHAR 原始文件名
file_hash VARCHAR 文件内容的 SHA-256 哈希值(用于幂等性)
status VARCHAR `PENDING`, `PROCESSING`, `COMPLETED`, `PARTIALLY\_COMPLETED`, `FAILED`
total_records INT 检测到的总行数
success_count INT 成功插入的记录数
failed_count INT 被拒绝的记录数
created_at TIMESTAMP 任务创建时间
updated_at TIMESTAMP 最后状态更新时间
`import\_record`
列名 类型 备注
id BIGINT (PK) 自动生成
job_id BIGINT (FK) 关联 `import\_job.id`
row_number INT 源文件中的行号
data JSON/TEXT 原始或解析后的记录数据
status VARCHAR `SUCCESS`, `FAILED`
error_message TEXT 校验/错误详情(可为空)
created_at TIMESTAMP 插入时间
对 `import\_job.file\_hash` 设置唯一约束可防止重复处理相同的文件。
项目结构
```
bulk-import-system/
├── src/main/java/com/example/bulkimport/
│ ├── controller/
│ │ └── ImportController.java
│ ├── service/
│ │ ├── ImportService.java
│ │ └── FileParserFactory.java
│ ├── parser/
│ │ ├── CsvFileParser.java
│ │ └── ExcelFileParser.java
│ ├── validator/
│ │ └── RecordValidator.java
│ ├── model/
│ │ ├── ImportJob.java
│ │ └── ImportRecord.java
│ ├── repository/
│ │ ├── ImportJobRepository.java
│ │ └── ImportRecordRepository.java
│ ├── dto/
│ │ ├── ImportJobResponse.java
│ │ └── ImportSummaryResponse.java
│ ├── config/
│ │ ├── AsyncConfig.java
│ │ └── SwaggerConfig.java
│ └── exception/
│ └── GlobalExceptionHandler.java
├── src/main/resources/
│ ├── application.yml
│ └── db/migration/ (Flyway scripts, if used)
├── src/test/java/...
├── sample-files/
│ ├── sample\_valid.csv
│ ├── sample\_with\_errors.csv
│ └── sample\_large\_1L\_records.csv
├── postman/
│ └── Bulk\_Import\_API.postman\_collection.json
├── pom.xml
└── README.md
```
## 前置条件
JDK 17+
Maven 3.8+
MySQL 8.x 或 PostgreSQL 14+(运行中的实例)
Postman(可选,用于 API 测试)
建议为大文件测试预留至少 2 GB 的可用堆内存
设置说明
克隆代码库
```
git clone
cd bulk-import-system
```
创建数据库
```
CREATE DATABASE bulk\_import\_db;
```
配置数据库凭证
更新 `src/main/resources/application.yml`(参见配置说明)。
构建项目
```
mvn clean install
```
运行数据库迁移(如果使用 Flyway/Liquibase)
```
mvn flyway:migrate
```
配置
`application.yml`(核心配置项):
```
server:
port: 8080
spring:
datasource:
url: jdbc:mysql://localhost:3306/bulk\_import\_db
username: root
password: root
jpa:
hibernate:
ddl-auto: validate
properties:
hibernate:
jdbc:
batch\_size: 500
order\_inserts: true
order\_updates: true
show-sql: false
servlet:
multipart:
max-file-size: 200MB
max-request-size: 200MB
import:
async:
core-pool-size: 4
max-pool-size: 8
queue-capacity: 50
batch:
size: 500
storage:
temp-dir: ./uploads/temp
```
运行应用程序
```
mvn spring-boot:run
```
## 应用程序启动地址:`http://localhost:8080`
Swagger UI:`http://localhost:8080/swagger-ui/index.html`
## API Endpoints
方法 Endpoint 描述
POST `/api/imports/upload` 上传 CSV/Excel 文件以启动导入任务
GET `/api/imports/{jobId}` 获取任务状态和摘要
GET `/api/imports/{jobId}/records?status=FAILED` 获取分页记录(可按状态过滤)
GET `/api/imports` 列出所有导入任务(分页显示)
GET `/api/imports/{jobId}/errors/export` 将失败记录下载为 CSV
请求/响应示例
上传文件
```
curl -X POST http://localhost:8080/api/imports/upload \\
-F "file=@sample-files/sample\_with\_errors.csv"
```
响应(202 Accepted)
```
{
"jobId": 101,
"fileName": "sample\_with\_errors.csv",
"status": "PROCESSING",
"message": "File accepted. Processing started asynchronously."
}
```
检查任务状态
```
curl http://localhost:8080/api/imports/101
```
响应
```
{
"jobId": 101,
"fileName": "sample\_with\_errors.csv",
"status": "PARTIALLY\_COMPLETED",
"totalRecords": 1000,
"successCount": 950,
"failedCount": 50,
"createdAt": "2026-07-14T10:15:00Z",
"updatedAt": "2026-07-14T10:16:32Z"
}
```
重复文件重新上传响应(409 Conflict)
```
{
"error": "DUPLICATE\_FILE",
"message": "This file has already been processed under jobId 101.",
"existingJobId": 101
}
```
输入文件示例
位于 `/sample-files` 目录下:
文件 用途
`sample\_valid.csv` 所有记录均通过校验
`sample\_with\_errors.csv` 有效和无效行的混合(缺少字段、格式错误)
`sample\_large\_1L\_records.csv` 约 100,000 行,用于负载/性能测试
CSV 格式示例:
```
name,email,phone,age
John Doe,john@example.com,9876543210,29
Jane Smith,invalid-email,9876543211,31
,missing-name@example.com,9876543212,25
```
测试
运行所有单元测试和集成测试:
```
mvn test
```
## 测试覆盖范围包括:
Bean Validation 规则(单元测试)
CSV/Excel 流式解析器(使用示例文件的单元测试)
Service 层批量插入逻辑(使用 Testcontainers 的集成测试)
幂等性检查(拒绝重复文件哈希)
异步任务完成轮询测试
大文件模拟测试(10 万+ 行)用于内存/性能基准测试
## 处理大文件与边缘情况
场景 处理策略
大文件(10 万+ 记录) 流式解析器(Commons CSV iterator / POI SAX API);绝不将整个文件加载到内存中
无效的文件格式 在创建任务之前,于上传时返回 `400 Bad Request` 拒绝处理
文件内的重复记录 标记为 `FAILED` 并附带错误信息;任务继续执行
重复上传文件 通过 SHA-256 文件哈希检测;返回现有任务信息,而不是重新处理
处理过程中系统崩溃 任务状态保持为 `PROCESSING`;启动时的恢复/核对任务会标记过期任务,并允许安全恢复或重新触发
部分行处理失败 行级 try/catch;失败的行连同原因被记录到 `import\_record`,有效的行继续处理
## 性能优化
流式 I/O:通过缓冲 iterator 读取 CSV;通过 POI 的流式(SAX/event)API 读取 Excel —— 无论文件大小如何,内存消耗保持恒定
批量插入:启用 Hibernate `hibernate.jdbc.batch\_size`、`order\_inserts` 和 `order\_updates`;在每个批处理周期中将记录 flush 并从持久化上下文中 clear
异步处理:Upload endpoint 立即返回 `jobId`;实际处理在专用的 `@Async` 线程池上运行,保持 API 的响应速度
分块事务:每个批次独立提交,以避免单个大型长时间运行的事务,并限制回滚的影响范围
连接池:调整 HikariCP 池大小以匹配批处理并发度
## 幂等性策略
上传时,计算文件内容的 SHA-256 哈希值。
检查 `import\_job` 表中是否存在具有相同 `file\_hash` 的行。
如果找到:返回现有任务的详细信息,而不是创建新任务或重新处理。
如果未找到:创建新的 `ImportJob` 并继续处理。
这保证了同一个文件永远不会被处理两次,即使在并发上传尝试下也是如此(通过对 `file\_hash` 的数据库唯一约束来强制执行)。
## 故障排除
问题 可能的修复方法
`MaxUploadSizeExceededException` 增加 `spring.servlet.multipart.max-file-size` / `max-request-size`
任务卡在 `PROCESSING` 检查异步线程池日志;确认启动恢复任务已运行
大文件插入缓慢 确认 Hibernate 批处理属性已激活;检查 `import\_record` 上的数据库索引
未检测到重复文件 确保 `file\_hash` 列具有唯一约束,并且哈希处理覆盖完整的文件字节
交付物检查清单
[ ] Spring Boot 项目源代码
[ ] API 文档(Swagger UI / Postman collection)
[ ] 输入文件示例(`/sample-files`)
[ ] 数据库 schema(DDL / 迁移脚本)
[ ] 包含设置和执行步骤的本 README 文件
标签:域名枚举, 测试用例