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 文件
标签:域名枚举, 测试用例