CODE_WIKI.md 35 KB

Server — Code Wiki

项目名称: qingjian-server(清鉴)
定位: 智能反贪数据研判平台后端服务
技术栈: Solon 3.10.3 / Java 25 / MyBatis-Plus 3.5.12 / DuckDB + SQLite + RocksDB / Apache Lucene 10.3.2 / JGraphT 1.5.2


目录

  1. 项目整体架构
  2. 技术栈与依赖
  3. 目录结构
  4. 启动与运行
  5. 多数据源机制
  6. 公共层详解 (common)
  7. 业务模块详解 (module)
  8. 数据流与模块依赖
  9. AI Agent 架构
  10. ETL 清洗管线
  11. 缓存体系
  12. SSE 实时推送
  13. 枚举与常量
  14. API 端点索引

1. 项目整体架构

qingjian Server 是一个面向纪检监察和经侦部门的智能数据研判平台 ,核心能力是处理和分析银行账单、第三方支付流水、通讯话单(CDR)、基站轨迹等多源异构数据,通过图算法、时序分析和异常检测模型挖掘隐蔽的贿赂线索、利益输送链条和攻守同盟关系。

架构分层

┌─────────────────────────────────────────────────────────┐
│                     Controller 层                        │
│    (Solon MVC, 路由注解 @Controller @Mapping)            │
├─────────────────────────────────────────────────────────┤
│                      Service 层                          │
│    (业务逻辑, @Inject 注入, MyBatis-Plus ServiceImpl)     │
├─────────────────────────────────────────────────────────┤
│                      Mapper 层                           │
│    (MyBatis-Plus BaseMapper<T>, XML Mapper)              │
├─────────────────────────────────────────────────────────┤
│                   数据访问层                              │
│    DynamicDataSource → SQLite(plat) / DuckDB(case)       │
│    RocksDB → 本地 KV (运营商/银行卡/缓存)                 │
│    Lucene → 全文检索索引                                 │
└─────────────────────────────────────────────────────────┘

核心数据流

文件导入 → 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 存储

搜索与图计算

依赖 版本 用途
lucene-core / queryparser / analysis-common / highlighter 10.3.2 全文检索引擎
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/Lucene)
│   ├── 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)
│       ├── LuceneManager.java        # Lucene 索引管理(IK 分词/读写锁保护)
│       ├── 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

构建命令

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 + 关闭 Lucene 索引

关键配置 (app.yml)

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 是自定义的 DuckDB 数据源实现:

  • 单进程共享:所有连接通过 duckDBConnection.duplicate() 复制,共享同一 DuckDB 实例
  • 自动配置:根据 CPU 核数和物理内存自动设置 memory_limit 和 threads
  • 嵌入式:DuckDB 以嵌入式模式运行,数据文件位于案件工作空间

数据源切换方式

// 编程式切换
DynamicDsKey.with(StrConsts.DS_KEY_PLAT, () -> {
    // 此代码块内使用 SQLite 数据源
    return tableInfoMapper.selectList(null);
});

// 注解式切换(Service 类级别)
@DynamicDs(StrConsts.DS_KEY_PLAT)
public class KnowledgeInfoService extends ServiceImpl<KnowledgeInfoMapper, KnowledgeInfo> { ... }

6. 公共层详解 (common)

6.1 统一响应封装

R.java 是全局统一响应对象:

public class R<T> {
    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 是 Solon 全局路由拦截器,自动将 Controller 返回值包装为 R:

  • 正常返回 → R.succeed(data)
  • ServerException → R.failure(code, msg)
  • 其他异常 → R.failure(500)

6.2 请求参数基类

Query.java 是所有请求参数的基类:

字段 类型 默认值 说明
page int 1 页码
limit int 20 每页条数
startDate String — 开始日期
endDate String — 结束日期
orderKey String — 排序字段
sort String ASC 排序方式
tableName String — 表名
personName String — 人员名称
personNames List<String> — 多人员名称
keyword String — 关键词

6.3 全局缓存

GlobalCache

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 管理案件级数据缓存:

  • 核心数据结构:PERSON_LIB_NO_MAP(HashMultimap<String, String>),人员名+类型 → 卡号/手机号集合
  • 查询方法: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 Lucene 全文检索

LuceneManager.java 封装 Lucene 索引管理:

  • 分词器:IKAnalyzer(中文分词)
  • 线程安全:resourceLock(ReentrantReadWriteLock)保护 IndexWriter 和 SearcherManager
  • 核心方法:add(SearchDTO)、addDocument(Document)、deleteDocument(Term)、updateDocument(Term, Document)
  • 搜索器管理:acquireSearcher() / releaseSearcher() 配对使用

6.6 全局线程池

GlobalPool.java:

  • CORE_NUM:核心线程数(最大 16)
  • EXC_POOL:ThreadPoolExecutor(容量 100000 阻塞队列)
  • blockingCoefficient(float):根据阻塞系数计算推荐线程数

6.7 案件状态管理

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 全文检索
SearchService Lucene 搜索
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: Lucene 索引 (全文检索)                      │
│  - 人员/卡号/手机号搜索                           │
│  - IK 中文分词                                   │
└─────────────────────────────────────────────────┘

缓存初始化流程

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