# Server — Code Wiki > **项目名称**: qingjian-server(清鉴) > **定位**: 智能反贪数据研判平台后端服务 > **技术栈**: Solon 3.10.3 / Java 25 / MyBatis-Plus 3.5.12 / DuckDB + SQLite + RocksDB / Elasticsearch (spring-data-elasticsearch 5.5) / JGraphT 1.5.2 --- ## 目录 1. [项目整体架构](#1-项目整体架构) 2. [技术栈与依赖](#2-技术栈与依赖) 3. [目录结构](#3-目录结构) 4. [启动与运行](#4-启动与运行) 5. [多数据源机制](#5-多数据源机制) 6. [公共层详解 (common)](#6-公共层详解-common) 7. [业务模块详解 (module)](#7-业务模块详解-module) 8. [数据流与模块依赖](#8-数据流与模块依赖) 9. [AI Agent 架构](#9-ai-agent-架构) 10. [ETL 清洗管线](#10-etl-清洗管线) 11. [缓存体系](#11-缓存体系) 12. [SSE 实时推送](#12-sse-实时推送) 13. [枚举与常量](#13-枚举与常量) 14. [API 端点索引](#14-api-端点索引) --- ## 1. 项目整体架构 qingjian Server 是一个面向纪检监察和经侦部门的**智能数据研判平台** ,核心能力是处理和分析银行账单、第三方支付流水、通讯话单(CDR)、基站轨迹等多源异构数据,通过图算法、时序分析和异常检测模型挖掘隐蔽的贿赂线索、利益输送链条和攻守同盟关系。 ### 架构分层 ``` ┌─────────────────────────────────────────────────────────┐ │ Controller 层 │ │ (Solon MVC, 路由注解 @Controller @Mapping) │ ├─────────────────────────────────────────────────────────┤ │ Service 层 │ │ (业务逻辑, @Inject 注入, MyBatis-Plus ServiceImpl) │ ├─────────────────────────────────────────────────────────┤ │ Mapper 层 │ │ (MyBatis-Plus BaseMapper, XML Mapper) │ ├─────────────────────────────────────────────────────────┤ │ 数据访问层 │ │ DynamicDataSource → SQLite(plat) / DuckDB(case) │ │ RocksDB → 本地 KV (运营商/银行卡/缓存) │ │ Elasticsearch → 全文检索索引 │ └─────────────────────────────────────────────────────────┘ ``` ### 核心数据流 ``` 文件导入 → DM(预处理+清洗) → Govern(治理/去重/关系计算) → 查询分析 ├── Graph (关系图谱) ├── Trans (交易分析) ├── Call (通话分析) ├── Track (轨迹追踪) ├── OTG (对象轨迹) └── Person(人员/亲密度) ``` --- ## 2. 技术栈与依赖 ### 核心框架 | 依赖 | 版本 | 用途 | |---|---|---| | Solon (solon-parent) | 3.10.3 | 应用框架(非 Spring Boot),IoC/AOP/MVC | | solon-web | 3.10.3 | Web 核心(路由、过滤器、静态资源) | | solon-web-sse | 3.10.3 | Server-Sent Events 支持 | | solon-serialization-jackson | 3.10.3 | JSON 序列化(替代默认 Snack4) | | solon-view-thymeleaf | 3.10.3 | 模板引擎(AI 聊天页面) | | solon-data-dynamicds | 3.10.3 | 动态多数据源 | | solon-ai / solon-ai-agent | 3.10.3 | AI Agent 框架(ReAct Agent) | | solon-ai-skill-web/cli/toolgateway | 3.10.3 | AI Skill 扩展 | ### 数据层 | 依赖 | 版本 | 用途 | |---|---|---| | mybatis-plus-solon-plugin | 3.5.12 | ORM 框架 + Solon 插件 | | mybatis-plus-jsqlparser-4.9 | 3.5.12 | SQL 解析(分页插件) | | sqlite-jdbc (willena) | 3.51.1.0 | SQLite JDBC 驱动(plat 库) | | duckdb_jdbc | 1.5.0.0 | DuckDB JDBC 驱动(case 库) | | HikariCP | 7.0.2 | 连接池 | | rocksdbjni | 10.2.1 | 本地 KV 存储 | ### 搜索与图计算 | 依赖 | 版本 | 用途 | |---|---|---| | spring-boot-starter-data-elasticsearch | Boot 3.5 托管 | 全文检索引擎(Elasticsearch 8.18 服务端) | | ik-analyzer | 9.0.0 | 中文分词器 | | jgrapht-core / ext / io | 1.5.2 | 图计算框架 | ### 工具库 | 依赖 | 版本 | 用途 | |---|---|---| | hutool-all | 5.8.42 | Java 工具集 | | guava | 33.5.0-jre | Google 核心库(Multimap 等) | | mapstruct | 1.6.3 | DTO-Entity 转换器 | | lombok | — | 注解简化(@Data/@Getter 等) | | fesod-sheet | 2.0.1-incubating | Excel/CSV 流式读取 | | oshi-core | 6.9.2 | 系统信息获取(CPU/内存) | | roguemap-core | 1.1.3 | Map 增强工具 | --- ## 3. 目录结构 ``` com.qingjian/ ├── App.java # 应用入口:Solon.start(),端口 8980,contextPath "!/js/a/" ├── common/ # 公共层 │ ├── base/ # 基础类 │ │ ├── Query.java # 请求参数基类(分页/排序/时间范围) │ │ ├── ErrorEnum.java # 错误码枚举 │ │ ├── BasicColumn.java # 基础列定义 │ │ ├── Cleaned.java # 清洗结果标记接口 │ │ └── TreeNode.java # 树形节点 │ ├── cache/ # 缓存 │ │ ├── GlobalCache.java # 全局元数据缓存(表/字段/正则/运营商) │ │ ├── CaseDataCache.java # 案件数据缓存(人员-卡号-手机号映射) │ │ └── RocksDBCache.java # Solon CacheService 的 RocksDB 实现 │ ├── config/ # 配置 │ │ ├── WebConfig.java # 数据源/Jackson/缓存 Bean 配置 │ │ ├── MybatisPlusConfig.java # MyBatis-Plus 分页插件 │ │ ├── DuckdbUnpooledDataSource.java # DuckDB 数据源(自动配置内存/线程) │ │ ├── R.java # 统一响应封装 │ │ ├── GlobalResultInterceptor.java # 全局响应拦截器 │ │ ├── BigDecimalSerializer/Deserializer.java # 金额序列化 │ │ └── LicenseFilter.java # 授权过滤器 │ ├── constants/ # 常量 │ │ ├── StrConsts.java # 字符串常量(数据源 key、缓存 key) │ │ ├── PathConst.java # 文件路径常量(数据库/工作空间/RocksDB) │ │ ├── RegConst.java # 正则表达式常量 │ │ ├── CmdConstants.java # 命令常量 │ │ └── FileTypeConstants.java # 文件类型常量 │ ├── enums/ # 枚举 │ │ ├── FieldTypeEnum.java # 字段类型(TEXT/MONEY/PHONE/DATE_TIME 等) │ │ ├── MappingEnum.java # 映射类型 │ │ ├── DelimiterEnum.java # 分隔符类型 │ │ ├── PersonTagEnum.java # 人员标签 │ │ ├── TimeSeriesTypeEnum.java # 时序类型 │ │ └── ... # 其他枚举 │ ├── event/ # 事件监听 │ │ ├── AppLoadEndEventListener.java # 启动完成(初始化缓存/RocksDB/AI) │ │ └── AppStopEndEventListener.java # 关闭(释放 RocksDB) │ ├── exception/ # 异常 │ │ └── ServerException.java # 服务端异常(code + msg) │ ├── model/ # DTO/Entity/Query(按模块组织) │ │ ├── ai/ # AI 模型(AgentProperties/ChatModeType/ModelType 等) │ │ ├── call/ # 通话(entity/dto/query) │ │ ├── dm/ # 数据管理(entity/dto/query/mapstruct) │ │ ├── govern/ # 数据治理 │ │ ├── graph/ # 图谱 │ │ ├── otg/ # 对象轨迹 │ │ ├── person/ # 人员 │ │ ├── plat/ # 平台 │ │ ├── track/ # 轨迹 │ │ └── trans/ # 交易 │ ├── rocksdb/ # RocksDB 封装 │ │ ├── AbstractRocksDBStorage.java # 抽象基类(生命周期/CRUD/列簇管理) │ │ ├── DefaultRocksDBStorage.java # 默认实现(3 列簇:default/ispData/cacheData) │ │ └── RocksDBHelper.java # 全局单例持有 Storage │ └── utils/ # 工具类 │ ├── GlobalPool.java # 全局线程池(虚拟线程 + 自定义 ThreadFactory) │ ├── Json.java # Jackson ObjectMapper 全局封装 │ ├── DataUtil.java # 数据工具(银行卡前缀/手机号前缀/金额格式化) │ ├── DateUtil.java # 日期工具 │ ├── SlowDownProgress.java # 渐进式进度条(三阶段策略) │ ├── StateManager.java # 案件状态单例(当前案件 ID/工作空间) │ ├── SystemConfManager.java # 系统配置管理(持久化到文件) │ ├── CoordinateTransform.java # 坐标转换(WGS84/GCJ02/BD09) │ ├── CellUtils/CellDll.java # 基站经纬度查询(JNI 调用) │ ├── TogetherUtil.java # 同行分析工具 │ ├── TreeUtils.java # 树形结构工具 │ ├── LicenseUtil.java # 授权验证 │ ├── MachineCodeUtils.java # 机器码生成 │ ├── JsonToMarkdownUtil.java # JSON 转 Markdown │ └── ThreadFactoryImpl.java # 自定义线程工厂 └── module/ # 业务模块 ├── ai/ # AI 智能模块 ├── call/ # 通话记录模块 ├── dm/ # 数据管理模块 ├── govern/ # 数据治理模块 ├── graph/ # 图谱模块 ├── otg/ # 对象轨迹生成模块 ├── person/ # 人员信息模块 ├── plat/ # 平台模块 ├── track/ # 轨迹追踪模块 └── trans/ # 交易记录模块 ``` --- ## 4. 启动与运行 ### 环境要求 - **JDK**: 25+ - **Maven**: 3.x - **操作系统**: Windows / Linux ### 构建命令 ```bash mvn compile -f pom.xml ``` ### 启动流程 1. **`App.main()`** — 入口方法 - 调用 `init()` 检查工作空间目录(如不存在则从原始路径复制) - `Solon.start()` 启动 Solon 框架,注册 CORS 过滤器(`CrossFilter`,优先级 -1) 2. **`AppLoadEndEventListener.onEvent()`** — 应用加载完成回调 - `initDir()` — 创建 tmp/RocksDB/license/agent/skills 等目录 - `initRocks()` — 创建 `DefaultRocksDBStorage` 并启动,注册到 `RocksDBHelper`,初始化银行卡/运营商数据 - `GlobalCache.initCache()` — 初始化表/字段元数据缓存 + 正则缓存 - `CaseDataCache.initCache()` — 从 case 库加载人员库编号映射 - `AiService` 初始化 — AI 服务就绪 3. **服务就绪** — 监听端口 `8980`,contextPath `!/js/a/` ### 关闭流程 - **`AppStopEndEventListener.onEvent()`** — 关闭 RocksDB ### 关键配置 (app.yml) ```yaml server: port: 8980 contextPath: "!/js/a/" solon.app: name: 'QingJian' group: 'VectorExplore' qingjian.db1: class: "org.noear.solon.data.dynamicds.DynamicDataSource" strict: false default: case plat: type: "org.apache.ibatis.datasource.unpooled.UnpooledDataSource" driver: org.sqlite.JDBC case: type: "com.qingjian.common.config.DuckdbUnpooledDataSource" driver: org.duckdb.DuckDBDriver memoryLimitInMiB: 4096 threads: 16 ``` --- ## 5. 多数据源机制 ### 数据源配置 | 数据源 Key | 数据库 | 用途 | 切换方式 | |---|---|---|---| | `plat` | SQLite | 平台配置、表字段元数据、运营商/银行字典 | `DynamicDsKey.with("plat", () -> ...)` | | `case` | DuckDB | 案件数据(通话/交易记录、关系边、去重中间表) | 默认数据源,或 `DynamicDsKey.with("case", () -> ...)` | ### DuckDB 数据源特性 [DuckdbUnpooledDataSource](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/config/DuckdbUnpooledDataSource.java) 是自定义的 DuckDB 数据源实现: - **单进程共享**:所有连接通过 `duckDBConnection.duplicate()` 复制,共享同一 DuckDB 实例 - **自动配置**:根据 CPU 核数和物理内存自动设置 `memory_limit` 和 `threads` - **嵌入式**:DuckDB 以嵌入式模式运行,数据文件位于案件工作空间 ### 数据源切换方式 ```java // 编程式切换 DynamicDsKey.with(StrConsts.DS_KEY_PLAT, () -> { // 此代码块内使用 SQLite 数据源 return tableInfoMapper.selectList(null); }); // 注解式切换(Service 类级别) @DynamicDs(StrConsts.DS_KEY_PLAT) public class KnowledgeInfoService extends ServiceImpl { ... } ``` --- ## 6. 公共层详解 (common) ### 6.1 统一响应封装 [R.java](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/config/R.java) 是全局统一响应对象: ```java public class R { int code; // 状态码 String message; // 描述 T data; // 数据 } ``` - 成功:`R.succeed(data)` → `{code: 200, message: "success", data: ...}` - 失败:`R.failure(500, "error")` → `{code: 500, message: "error", data: null}` [GlobalResultInterceptor](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/config/GlobalResultInterceptor.java) 是 Solon 全局路由拦截器,自动将 Controller 返回值包装为 `R`: - 正常返回 → `R.succeed(data)` - `ServerException` → `R.failure(code, msg)` - 其他异常 → `R.failure(500)` ### 6.2 请求参数基类 [Query.java](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/base/Query.java) 是所有请求参数的基类: | 字段 | 类型 | 默认值 | 说明 | |---|---|---|---| | `page` | int | 1 | 页码 | | `limit` | int | 20 | 每页条数 | | `startDate` | String | — | 开始日期 | | `endDate` | String | — | 结束日期 | | `orderKey` | String | — | 排序字段 | | `sort` | String | ASC | 排序方式 | | `tableName` | String | — | 表名 | | `personName` | String | — | 人员名称 | | `personNames` | List\ | — | 多人员名称 | | `keyword` | String | — | 关键词 | ### 6.3 全局缓存 #### GlobalCache [GlobalCache.java](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/cache/GlobalCache.java) 管理平台级元数据缓存: - **表/字段元数据**:`TABLE_INFO_LIST`、`TABLE_FIELD_CACHE`、`ALL_FIELDS_MAP`、`ALL_MATCHED_FIELDS_MAP`、`TABLE_HEAD` - **功能正则**:`FUN_REGULAR_MAP`(邮箱/手机/身份证/金额/车牌) - **治理配置**:`DIRECTION_CONF_MAP`、`TEMPLATE_TABLE_MAP`、`ID_NAME_ZH_MAP` - **RocksDB 初始化**:`initCaseRocksDbData()` 将银行卡 BIN 和运营商数据写入 RocksDB #### CaseDataCache [CaseDataCache.java](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/cache/CaseDataCache.java) 管理案件级数据缓存: - **核心数据结构**:`PERSON_LIB_NO_MAP`(`HashMultimap`),人员名+类型 → 卡号/手机号集合 - **查询方法**:`getPerPhones(personName)`、`getPerCards(personName)`、`findPersonVal(val)` - **维护方法**:`putPerLib()`、`cleanCache()` ### 6.4 RocksDB 封装 RocksDB 封装采用三层架构: ``` RocksDBHelper (全局单例) └── DefaultRocksDBStorage (3 列簇实现) └── AbstractRocksDBStorage (抽象基类,生命周期/CRUD) ``` **三个列簇**: | 列簇 | 用途 | Key 格式 | |---|---|---| | `default` | 银行卡 BIN → 发卡行 | `bank_{cardPrefix6}` | | `ispData` | 手机号前7位 → 运营商 | `phone_{phonePrefix7}` | | `cacheData` | Solon CacheService 存储 | 自定义 key | ### 6.5 Elasticsearch 全文检索 全文检索已由嵌入式 Lucene 迁移到 **Elasticsearch**(spring-data-elasticsearch 5.5, 服务端 8.18 单节点,`spring.elasticsearch.uris` 默认 `http://127.0.0.1:9200`),代码位于 `module/search/`: - **实体**:`SearchDoc`(`@Document(indexName = "zsjz_search")`)——所有案件共用单索引, 以 `caseId` 关键字字段隔离(写入端强制从 `CaseContextHolder` 注入,查询/删除必带过滤); `_id = caseId|tableName|rowId`(行主键字段不能命名为 id,见 SearchDoc 注释),同表同记录重复清洗幂等覆盖 - **分词**:`es/search-settings.json` 定义 `unit_separator` 分词器(按 `\x1F` 切段 + 小写), 对应各实体 `str()` 的全字段拼接约定;检索为对分词后 token 的 wildcard 子串匹配 - **写入**:`SearchDocIndexService.add(SearchDoc)` 按 caseId 攒批 1000 条经 `saveAll` 走 `_bulk`; 全量清洗完成回调 `flush()` 提交尾批并 refresh;`deleteByFileIds` / `deleteByCase` 走 delete_by_query - **查询**:`SearchQueryService.search(keyword)`(`GET /search/list` 出入参契约不变), dateTime 降序 top 500 + unified 高亮 `` - **索引初始化**:`SearchIndexInitializer` 启动时 `createWithMapping`(幂等,ES 不可用不阻断启动) ### 6.6 全局线程池 [GlobalPool.java](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/utils/GlobalPool.java): - `CORE_NUM`:核心线程数(最大 16) - `EXC_POOL`:`ThreadPoolExecutor`(容量 100000 阻塞队列) - `blockingCoefficient(float)`:根据阻塞系数计算推荐线程数 ### 6.7 案件状态管理 [StateManager.java](file:///e:/workspace/qingjian-server/src/main/java/com/qingjian/common/utils/StateManager.java) (静态内部类单例): - `getCaseId()` — 获取当前案件 ID(默认 888888) - `caseWorkspace()` — 获取当前案件工作空间路径 - `cleanCase()` — 清除案件信息 --- ## 7. 业务模块详解 (module) ### 7.1 AI 智能模块 (ai) **职责**:提供 AI 对话、知识库管理、向量检索能力 | 类 | 说明 | |---|---| | `AiController` | AI 聊天接口(SSE 流式响应)、模型配置 | | `AiService` | 核心聊天逻辑:构建 AgentRuntime → 流式推理 → 映射 chunk | | `AgentRuntime` | ReAct Agent 运行时(加载 AGENTS.md、注册 Skill/Tool) | | `SseEmitterManager` | SSE 连接管理(24h 超时) | | `ChatModelService` | AI 模型配置 CRUD | | `ChatSessionService` | 会话管理 | | `ChatMessageService` | 消息管理 | | `KnowledgeInfoService` | 知识库信息管理(plat 数据源) | | `VectorStoreService` | 向量检索(预留) | | `SkillService` | AI Skill 管理 | **AI Skill 体系**: | Skill | 说明 | |---|---| | `Tools` | 内置工具集 | | `CliSkillProvider` | CLI 技能提供者 | | `ExpertSkill` | 专家技能 | | `Text2SqlSkill` | 自然语言转 SQL(支持 DuckDB 方言) | **AI 角色**:系统提示词定义在 `resources/prompt/AGENTS.md`,角色为"智能反贪数据研判专家",核心能力包括资金流向穿透分析、通讯关系图谱构建、时空轨迹碰撞、新型腐败识别。 ### 7.2 数据管理模块 (dm) **职责**:文件导入、数据预处理、ETL 清洗 | 类 | 说明 | |---|---| | `DmController` | 文件预览、清洗执行、文件管理、治理配置 | | `DmService` | 核心业务:文件预处理(虚拟线程并发)+ 数据清洗(CleanTask)+ 文件删除 | | `CleanFactory` | 工厂类:根据 StrategyType 创建清洗器和数据加载器(20 种策略) | | `CleanTask` | 单文件 ETL 任务(FesodSheet 流式读取 → StreamingEtlListener 处理) | | `AbstractDataCleaner` | 清洗器抽象基类(模板方法:doClean → 校验 → after) | | `AbstractDataLoader` | 数据加载器抽象基类(DuckDB Appender 批量写入) | | `PreTask` / `PreDataListener` | 文件预处理(Excel/CSV 解析 + MD5 去重) | | `CleanErrorLogService` | 清洗错误日志 | **支持的清洗策略** (StrategyType): | 策略 | Code | 说明 | |---|---|---| | CALL_DATA | 1 | 通话数据清洗 | | TRANS_STD | 2 | 标准交易数据清洗 | | CALL_OPEN_INFO | 3 | 通话开通信息清洗 | | TRANS_OPEN | 4 | 交易开通数据清洗 | | TRANS_IN_OUT | 5/6/7 | 交易进出数据清洗(3 种变体) | | TRANS_CORE_HIS_BALANCE | 8 | 交易核心历史余额清洗 | | TRANS_AMOUNT_BALANCE_CALC | 9 | 交易金额余额计算清洗 | | CALL_IN_OUT | 10 | 通话进出数据清洗 | | TRANS_DEBIT_CREDIT | 11 | 交易借贷数据清洗 | | RESIDENT_POPULATION_INFO | 12 | 常住人口信息清洗 | | EXPRESS_INFO | 13 | 快递信息清洗 | | HOTEL_STAY_INFO | 14 | 酒店住宿信息清洗 | | FLIGHT_TICKET_INFO | 15 | 机票信息清洗 | | FLIGHT_DEPARTURE_INFO | 16 | 航班出发信息清洗 | | TOGETHER_FLIGHT_INFO | 17 | 同行航班信息清洗 | | TOGETHER_TRAIN_INFO | 18 | 同行火车信息清洗 | | TRAIN_TICKET_INFO | 19 | 火车票信息清洗 | | TOGETHER_LIVE_INFO | 20 | 同行住宿信息清洗 | ### 7.3 数据治理模块 (govern) **职责**:关系边计算、去重、时序图生成、分类树统计 | 类 | 说明 | |---|---| | `GovernController` | SSE 连接 + 治理计算任务启动 | | `GovernService` | 三阶段治理编排(详见下方) | | `GovernCache` | 时序图 UNION ALL SQL 生成(10 种数据源) | | `GovernConfService` | 治理配置管理(动态表名/去重/开户信息开关) | | `GovernTreeService` | 分类树管理 | | `CallResultHandler` | 通话关系结果处理 | | `TransResultHandler` | 交易关系结果处理 | **三阶段治理流程**: ``` 阶段一:数据清洗去重 ├── 交易去重(并行) ├── 通话去重(并行) └── 关系边清理(并行) └── [可选] 开户信息任务:清理开户表 → 复制手机号数据 阶段二:时序图 + 分类树(并行) ├── 创建时序图(10 种数据源 UNION ALL) └── 分类树统计 阶段三:关系边计算(并行) ├── 通话关系边 ├── 交易关系边 └── 其他关系边 └── 合并关系边 ``` ### 7.4 图谱模块 (graph) **职责**:关系图谱生成与查询 | 类 | 说明 | |---|---| | `GraphController` | 图谱视图、交易/通话明细、历史记录 | | `GraphService` | 核心图谱生成:创建临时核心节点表 → 动态拼接 DuckDB SQL → 提取节点/边 | **图谱类型**: | 类型 | Code | 说明 | |---|---|---| | 全部交易 | 101 | 查询与目标人员有关的所有交易 | | 共同交易 | 102 | 多人共同交易对方 | | 相互交易 | 103 | 双方直接交易 | | 全部通话 | 201 | 查询与目标人员有关的所有通话 | | 共同通话 | 202 | 多人共同通话对方 | | 相互通话 | 203 | 双方直接通话 | ### 7.5 人员信息模块 (person) **职责**:人员基本信息管理、证件号/手机号/银行卡号库、亲密度分析 | 类 | 说明 | |---|---| | `PersonBasicInfoController` | 人员 CRUD、排序 | | `PersonBasicInfoService` | 人员管理:创建/删除/更新/刷新缓存 | | `IntimacyController` | 亲密度分析 | | `IntimacyService` | 8 维度亲密度评估(大额转账/特殊日期/异卡同号/频繁联系/夜间通话/特殊日期联系/频繁交易/同行) | | `PersonGroupController` | 人员分组管理 | | `PersonGroupService` | 分组 CRUD | | `PersonLibNoService` | 人员库编号管理 | | `DataProfileController` | 数据画像 | | `DataProfileService` | 数据画像服务 | | `ResidentPopulationController` | 常住人口管理 | ### 7.6 平台模块 (plat) **职责**:SSE 推送、表/字段元数据管理、案件管理、系统配置 | 类 | 说明 | |---|---| | `TableInfoController` | 模板 CRUD | | `TableInfoService` | 模板管理(plat 数据源) | | `SseService` | SSE 连接管理(6h 超时,单例持有 SseEmitter) | | `CacheManageService` | 缓存管理(按前缀批量删除) | | `CaseInfoController` | 案件信息管理 | | `CaseInfoService` | 案件 CRUD | | `SystemController` | 系统信息 | | `SystemService` | 系统服务 | | `LunarService` | 农历计算 | ### 7.7 通话记录模块 (call) **职责**:通话记录查询与分析 | 类 | 说明 | |---|---| | `CallRecordController` | 通话记录分页查询 | | `CallRecordService` | 通话记录查询(动态表名) | | `CallContinuousController` | 连续通话分析 | | `CallContinuousService` | 连续通话统计 | | `CallNightController` | 夜间通话分析 | | `CallNightService` | 夜间通话统计 | | `CallSensitiveController` | 敏感通话分析 | | `CallSensitiveService` | 敏感通话统计 | | `CallOpenInfoService` | 开户信息查询 | ### 7.8 交易记录模块 (trans) **职责**:交易记录查询与多维度分析 | 类 | 说明 | |---|---| | `TransRecordController` | 交易记录分页查询 | | `TransRecordService` | 交易记录查询(多条件筛选,金额支持 5 种比较符) | | `TransFundFlowController` | 资金流向分析 | | `TransFundFlowService` | 资金流向图/追溯(树形结构,最多 4 层,向前/向后追溯) | | `TransBigController/Service` | 大额交易分析 | | `TransContinuousController/Service` | 连续交易分析 | | `TransFrequencyController/Service` | 频繁交易分析 | | `TransSensitiveController/Service` | 敏感交易分析 | | `TransCashFlowController/Service` | 现金流分析 | | `TransFastFundFlowController/Service` | 快速资金流向分析 | | `TransFinancialController/Service` | 财务分析 | | `TransCardHoldController/Service` | 持卡分析 | | `TransFixedDepositController/Service` | 定期存款分析 | | `CashKeyController/Service` | 现金关键字分析 | | `FinancialKeyController/Service` | 财务关键字分析 | | `FixedDepositKeyController/Service` | 定期存款关键字分析 | ### 7.9 轨迹追踪模块 (track) **职责**:多维度轨迹追踪与碰面分析 | 类 | 说明 | |---|---| | `TrackMeetController/Service` | 轨迹碰面分析(基站区域同时出现) | | `TrackCellTowerController/Service` | 基站轨迹查询 | | `TrackEnLocalController/Service` | 入境轨迹查询 | | `TrackTogetherLiveController/Service` | 同住分析 | | `TrackTogetherTravelController/Service` | 同行出行分析 | | `ExpressInfoController/Service` | 快递信息查询 | ### 7.10 对象轨迹生成模块 (otg) **职责**:时序数据查询与展示 | 类 | 说明 | |---|---| | `TimeSeriesController` | 时序数据列表/详情/跳转查询 | | `TimeSeriesService` | 时序数据查询(支持时间范围/人员/类型过滤) | | `SearchController` | 全文检索 | | `SearchQueryService` | Elasticsearch 查询(module/search) | | `SpecialDateController/Service` | 特殊日期管理 | --- ## 8. 数据流与模块依赖 ### 模块依赖关系 ``` ┌─────────┐ │ Plat │ ← SSE/缓存/模板/案件管理 └────┬────┘ │ 被依赖 ┌──────────────┼──────────────┐ ▼ ▼ ▼ ┌─────────┐ ┌──────────┐ ┌─────────┐ │ DM │ │ Govern │ │ Person │ │(导入清洗)│──▶│(治理计算) │──▶│(人员管理)│ └─────────┘ └────┬─────┘ └────┬────┘ │ │ ┌────────────┼──────────────┤ ▼ ▼ ▼ ┌─────────┐ ┌─────────┐ ┌──────────┐ │ Graph │ │ Trans │ │ Call │ │(关系图谱)│ │(交易分析)│ │(通话分析) │ └─────────┘ └─────────┘ └──────────┘ ▲ ▲ ▲ │ │ │ └────────────┼──────────────┘ ▼ ┌──────────┐ │ Track │ ← 轨迹追踪 │ OTG │ ← 对象轨迹 └──────────┘ ┌─────────┐ │ AI │ ← 独立模块,通过 Text2SqlSkill 查询 case 库 └─────────┘ ``` ### 核心数据流 1. **导入阶段**:用户上传 Excel/CSV → `DmService.preProcessFile()` → 虚拟线程并发解析 → MD5 去重 2. **清洗阶段**:`DmService.doClean()` → `CleanTask` → `StreamingEtlListener` → `AbstractDataCleaner.doClean()` → `AbstractDataLoader.write()` → DuckDB Appender 批量写入 3. **治理阶段**:`GovernService.calcTask()` → 三阶段编排(去重 → 时序图 → 关系边) 4. **分析阶段**:各查询模块从 DuckDB 读取治理后的数据,提供多维度分析 5. **AI 阶段**:`AiService.chat()` → `AgentRuntime.stream()` → ReAct Agent 推理 → Text2SqlSkill 查询 --- ## 9. AI Agent 架构 ### 架构概览 ``` AiController └── AiService ├── ChatModelService (模型配置) ├── ChatMessageService (消息管理) ├── KnowledgeInfoService (知识库) ├── VectorStoreService (向量检索) └── AgentRuntime (ReAct Agent 运行时) ├── AGENTS.md (系统提示词) ├── Tools (内置工具) ├── CliSkillProvider (CLI 技能) ├── Text2SqlSkill (NL2SQL) │ ├── SqlDialectManager (方言管理) │ └── SqlDialect (DuckDB 方言) ├── CompositeSummarizationStrategy (上下文摘要) └── AgentSessionProvider (文件持久化会话) ``` ### ReAct Agent 工作流 1. **接收用户输入** → 构建上下文消息 2. **推理循环**:思考(Reason) → 选择工具(Act) → 执行 → 观察结果(Observe) → 循环 3. **输出流式映射**:区分 `reason`(推理过程)、`action`(工具调用)、`text`(最终回复)三种 chunk 类型 4. **会话管理**:通过 `AgentSessionProvider` 持久化到文件 ### Text2SqlSkill 将自然语言查询转换为 DuckDB SQL,支持: - Schema 模式感知 - DuckDB 方言适配 - 安全沙箱执行 --- ## 10. ETL 清洗管线 ### 管线架构 ``` Excel/CSV 文件 │ ▼ PreTask (预处理) │ FesodSheet 流式读取 │ MD5 去重 │ 信号量限流 ▼ CleanTask (清洗任务) │ ▼ StreamingEtlListener (流式 ETL 监听器) │ ├── CleanFactory.createDataCleaner(strategyType) ──▶ AbstractDataCleaner │ │ │ ▼ │ doClean(rowData, lineNo) ← 子类实现具体清洗逻辑 │ │ │ ▼ │ after() ← 后置处理 │ └── CleanFactory.createDataLoader(strategyType) ──▶ AbstractDataLoader │ ▼ convert() → fillData() → doWrite() ← DuckDB Appender 批量写入 ``` ### 清洗器核心方法 `AbstractDataCleaner` 提供的模板方法: - `init(funcRegx, sheet)` — 初始化规则列表和函数映射 - `prepare(pre)` — 从表头前文本提取本方姓名/卡号/证件号 - `getContent(rowData, fileNameIndex, columnRule)` — 根据规则获取单元格内容 - `resolveContent()` — 解析映射规则(REF引用/NEW_DATA新数据)和拆分规则 ### 数据加载器 `AbstractDataLoader` 使用 DuckDB Appender 实现高性能批量写入: - 构造器初始化 DuckDB 连接和 Appender - `write(record)` 模板方法:`convert()` → `fillData()` → `doWrite()` - `close()` 关闭 Appender 和 Connection --- ## 11. 缓存体系 ### 三级缓存架构 ``` ┌─────────────────────────────────────────────────┐ │ L1: 内存缓存 (GlobalCache / CaseDataCache) │ │ - 表/字段元数据、正则、人员-卡号映射 │ │ - 启动时加载,运行时读写 │ ├─────────────────────────────────────────────────┤ │ L2: RocksDB (本地 KV 存储) │ │ - 银行卡 BIN → 发卡行 (default 列簇) │ │ - 手机号前7位 → 运营商 (ispData 列簇) │ │ - Solon CacheService (cacheData 列簇) │ ├─────────────────────────────────────────────────┤ │ L3: Elasticsearch 索引 (全文检索) │ │ - 人员/卡号/手机号搜索 │ │ - unit_separator 分词( 分段+小写) │ └─────────────────────────────────────────────────┘ ``` ### 缓存初始化流程 ``` AppLoadEndEventListener ├── initRocks() │ ├── DefaultRocksDBStorage.start() → 打开 RocksDB 实例 │ └── GlobalCache.initCaseRocksDbData() → 写入银行卡/运营商数据 ├── GlobalCache.initCache() │ ├── initRegCache() → 初始化正则缓存 │ ├── initTableInfoCache() → 从 plat 库加载表信息 │ └── initFieldCache() → 从 plat 库加载字段信息 └── CaseDataCache.initCache() → 从 case 库加载人员库编号 ``` --- ## 12. SSE 实时推送 ### 两个 SSE 服务 | 服务 | 超时 | 用途 | |---|---|---| | `SseService` (plat) | 6 小时 | 治理进度、清洗进度推送 | | `SseEmitterManager` (ai) | 24 小时 | AI 聊天流式响应 | ### SSE 事件类型 - **进度推送**:`SlowDownProgress` 生成渐进式进度值(三阶段策略:>50% 大步、>20% 中步、<20% 小步) - **AI 流式**:区分 `reason`/`action`/`text` 三种 chunk 类型 - **心跳**:`ssePut()` 定期发送心跳保持连接 --- ## 13. 枚举与常量 ### 核心枚举 | 枚举 | 说明 | |---|---| | `FieldTypeEnum` | 字段类型:TEXT/MONEY/NUMBER/DATE_TIME/DATE/TIME/PHONE/ACCT_CARD_NO/DEBIT_CREDIT_FLAG/CALL_DIRECTION/ID_CARD_NUMBER | | `StrategyType` | 清洗策略类型(20 种,见 7.2 节) | | `MappingEnum` | 映射类型 | | `DelimiterEnum` | 分隔符类型 | | `PersonTagEnum` | 人员标签 | | `TimeSeriesTypeEnum` | 时序类型 | | `ErrorEnum` | 错误码:400/401/403/404/500 | ### 核心常量 | 常量类 | 关键常量 | |---|---| | `StrConsts` | `DS_KEY_PLAT`("plat")、`DS_KEY_CASE`("case")、`CACHE_TAG_QING_JIAN` | | `PathConst` | `RUN_QINGJIAN_PATH`、`WORKSPACE`、`SYSTEM_ROCKSDB_PATH`、`BASE_CASE_DB_PATH`、`PLAT_DB_PATH`、`CELL_TOWER_PATH`、`LICENSE_PATH` | | `RegConst` | 邮箱/手机/身份证/金额/车牌正则 | --- > **文档生成时间**: 2026-05-15 > **项目版本**: 1.0 > **Solon 版本**: 3.10.3 > **Java 版本**: 25