QDATA 数据治理与 ETL 报告¶
生成日期:2026-08-27 | 基于源码深度分析
一、ETL 数据集成¶
1.1 双引擎架构¶
qData 采用 DataX + Spark 双引擎架构:
| 引擎 | 适用场景 | 数据源支持 |
|---|---|---|
| DataX | RDBMS 间数据库同步(轻量级) | rdbmsreader / rdbmswriter |
| Spark | 复杂转换、CSV/Excel 导入、分布式处理 | DB / CSV / Excel Reader |
┌──────────────────────────────────────────────────────────┐
│ qdata-module-dpp │
│ (ETL 任务管理、调度、监控) │
├──────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────┐ ┌─────────────────────────┐ │
│ │ DataX 引擎 │ │ Spark 引擎 │ │
│ │ DataXExecutor │ │ EtlApplication │ │
│ │ DataXJsonBuilder │ │ ReaderFactory │ │
│ │ Python subprocess │ │ TransitionFactory │ │
│ │ │ │ WriterFactory │ │
│ │ · rdbmsreader │ │ │ │
│ │ · rdbmswriter │ │ · DBReader (3 模式) │ │
│ │ · VALUE_MAP │ │ · CsvReader │ │
│ │ · ADD_CONSTANT │ │ · ExcelReader │ │
│ │ · SELECT_FIELDS │ │ · DBWriter (3 模式) │ │
│ │ · FIELD_DERIVATION │ │ · 7 种 Transition │ │
│ └──────────┬──────────┘ └────────────┬────────────┘ │
│ │ │ │
│ └─────────┬──────────────────┘ │
│ ▼ │
│ ┌──────────────────┐ │
│ │ DolphinScheduler │ │
│ │ / Quartz 调度 │ │
│ └──────────────────┘ │
└──────────────────────────────────────────────────────────┘
1.2 DataX 引擎¶
核心类:
| 类 | 职责 |
|---|---|
DataXExecutor |
通过 Python 子进程执行 DataX 任务(python3 datax.py job.json) |
DataXJsonBuilder |
生成 DataX job.json 配置,支持 Reader/Writer/Processor 节点 |
DataXProperties |
配置项:datax.home、datax.python-command、datax.datax-py-path、datax.job-dir |
支持的 Processor:VALUE_MAP(值映射)、ADD_CONSTANT(添加常量列)、SELECT_FIELDS(列选择)、FIELD_DERIVATION(字段派生)。
1.3 Spark 引擎¶
入口:tech.qiantong.qdata.spark.etl.EtlApplication.main()
执行流程:
1. 接收 DolphinScheduler 传递的 Base64 编码 JSON 参数(reader + transition[] + writer + config)
2. 通过 RabbitMQ 发布 RUNNING 状态
3. 初始化 Redis(增量游标存储)
4. 创建 SparkSession
5. 执行管线:ReaderFactory → TransitionFactory(N 个转换)→ WriterFactory
6. 通过 RabbitMQ 上报 SUCCESS/FAILURE 状态
7. 成功写入后将增量游标存入 Redis
Reader 注册表:
| 类型码 | 类 | 说明 |
|---|---|---|
DB_READER |
DBReader |
Spark JDBC 读取 |
EXCEL_READER |
ExcelReader |
Excel 文件读取 |
CSV_READER |
CsvReader |
CSV 文件读取 |
DBReader 读取模式:
| 模式 | 说明 |
|---|---|
| 1 = 全量 | 读取全表 |
| 2 = ID 增量 | 基于 Redis 缓存的 ID 游标 |
| 3 = 时间范围增量 | 支持固定值、时间范围、SQL 表达式 |
Transition 注册表:
| 类型码 | 类 | 说明 |
|---|---|---|
SPARK_CLEAN |
CleanTransition |
数据清洗(31KB,最复杂的转换) |
SORT_RECORD |
SortTransition |
排序 |
FIELD_DERIVATION |
FieldDerivationTransition |
字段派生 |
DATA_DEDUPLICATION |
DataDeduplicationTransition |
数据去重 |
VALUE_MAP |
ValueMapTransition |
值映射 |
ADD_CONSTANT |
AddConstantTransition |
添加常量列 |
SELECT_FIELDS |
SelectFieldsTransition |
列选择 |
Writer 注册表:
| 类型码 | 类 | 说明 |
|---|---|---|
DB_WRITER |
DBWriter |
Spark JDBC 写入 |
DBWriter 写入模式:
| 模式 | 说明 |
|---|---|
| 1 = 全量 | 临时表 + 重命名交换 |
| 2 = 追加 | 直接追加写入 |
| 3 = 增量更新 | Upsert(PreparedStatement 批处理) |
支持 preSql / postSql 执行,方言适配 Kingbase8、SQL Server、Doris 的表重命名操作。
JDBC 驱动内置:MySQL 8、DM8、Oracle、KingBase、SQLServer、jtds。
1.4 调度集成¶
| 调度器 | 模式 | 说明 |
|---|---|---|
| Quartz | 进程内 | 轻量级,适合单机部署 |
| DolphinScheduler | 分布式 | 通过 HTTP API 集成,31 个端点 |
DolphinScheduler API 端点(定义于 QianTongDCApiType):
| 类别 | 端点数 | 示例 |
|---|---|---|
| 项目管理 | 4 | POST /v2/projects, GET /v2/projects/{code} |
| 流程定义 | 11 | 创建、更新、删除、发布、批量操作、版本管理 |
| 调度管理 | 6 | 创建、更新、上线、下线、预览 |
| 执行管理 | 2 | 启动流程实例、执行流程 |
| 监控查询 | 5 | 流程实例、任务实例、日志 |
| 任务编码 | 1 | 生成任务定义编码 |
回调机制:ETL Spark 任务不直接调用 DolphinScheduler REST API,而是通过 RabbitMQ 发布状态更新:
- 交换机:ds.exchange.processInstance、ds.exchange.taskInstance
- 路由键:ds.queue.processInstance、ds.queue.taskInstance.insert、ds.queue.taskInstance.update
二、数据开发¶
2.1 任务类型¶
| 类型码 | 说明 |
|---|---|
| 1 | 离线任务 |
| 2 | 实时任务 |
| 3 | 数据开发任务(SQL 脚本) |
| 4 | 作业任务 |
2.2 任务状态¶
| 状态码 | 说明 |
|---|---|
| -3 | 部分上线 |
| -2 | 草稿 |
| 0 | 下线 |
| 1 | 上线 |
2.3 核心功能¶
- 可视化 ETL 编排:基于 AntV X6 的拖拽式流程编辑器
- SQL 脚本开发:Monaco Editor 支持的 SQL 编辑与执行
- 任务版本管理:节点版本化,支持历史回溯
- 统计面板:运行中数量、今日错误数、今日执行数、成功率
三、数据质量¶
3.1 独立微服务¶
qdata-service-quality(端口 8083)独立部署,使用 MongoDB 存储错误数据。
3.2 质量规则引擎¶
工厂模式:QualitySqlGenerateFactory 根据规则类型分发到对应的 SQL 生成器。
9 种质量规则生成器:
| 生成器 | 维度 | 说明 |
|---|---|---|
CharacterValidationGenerator |
有效性 | 字符格式验证 |
CompositeUniquenessGenerator |
唯一性 | 组合字段唯一性 |
DecimalPrecisionGenerator |
有效性 | 小数精度验证 |
EnumValidationGenerator |
有效性 | 枚举值验证 |
GroupFieldCompletenessGenerator |
完整性 | 分组字段完整性 |
LengthValidationGenerator |
有效性 | 长度验证 |
NumericRangeValidationGenerator |
有效性 | 数值范围验证 |
TimeOrderValidationGenerator |
有效性 | 时间顺序验证 |
CommonGenerator |
通用 | 通用规则 |
3.3 数据库方言注册¶
ComponentRegistry 注册了 6 种数据库方言:
| 数据库 | 方言类 |
|---|---|
| MySQL | MySqlQuality |
| Oracle 12c | Oracle12cQuality |
| Oracle | OracleQuality |
| SQL Server | SQLServerQuality |
| DM8 | DM8Quality |
| 默认 | DefaultQuality |
3.4 执行流程¶
QualityTaskExecutorServiceImpl.executeTask(taskId)提交异步执行(ThreadPoolTaskExecutor+ Redis 去重锁)- 查询任务 → 生成批次号 → 获取数据源对象 → 按数据源 ID 分组
- 对每个数据源:连接 → 加载规则 → 通过
RuleTaskExecutor逐条执行规则 - 错误数据存入 MongoDB(
CheckErrorData文档) - 支持 3 种错误数据修改方式:
updateType=1:修改数据 + 更新物理表updateType=2:修改备注updateType=3:修改修复状态
四、元数据管理¶
4.1 采集架构¶
qdata-module-mc 提供自动化的元数据采集,支持增量/全量两种模式,使用 Neo4j 图数据库存储元数据关系。
4.2 采集流程¶
McTaskServiceImpl.runDaDiscoveryTask(taskId)
│
├── Redis 分布式锁(防重复执行)
├── 创建采集实例
│
└── executeTaskSafely(task, instance)
│
├── 准备数据源连接
├── 加载数据库范围(自定义 / 全部)
├── 对比数据库列表(新增/删除/变更)
│
└── 对每个数据库:
├── 加载表列表 → 对比表列表
├── 加载列信息 → 对比列信息
└── 记录变更日志
4.3 变更检测¶
| 变更类型 | 说明 |
|---|---|
| "1" | 表注释变更 |
| "2" | 列新增/删除/修改 |
| "3" | 索引字段变更 |
| "4" | 存储大小变更 |
| "1"-"9"(列级) | 注释、类型、长度、精度、小数位、默认值、主键、外键、可空 |
4.4 数据库方言¶
DatabaseDialectFactory 注册的方言:
| 数据库 | 方言类 | 状态 |
|---|---|---|
| MySQL | MySqlDialect |
完整实现 |
| Hive | HiveDialect |
完整实现 |
| DM8 | DamengDialect |
完整实现 |
| Oracle | OracleDialect |
占位实现 |
| PostgreSQL | PostgreSqlDialect |
占位实现 |
| SQL Server | AbstractDialect |
占位实现 |
方言接口方法:getStorageEngine(), getTableRowCount(), getTableIndexes(), getTablePartitionFields(), isColumnAutoIncrement(), getDbMetadata(), getTableMetadata(), getColumnMetadata()。
五、数据资产¶
5.1 资产编目¶
DaAssetServiceImpl 是核心服务,集成 11 个子服务:
| 子服务 | 职责 |
|---|---|
IDaAssetColumnService |
资产字段编目 |
IDaAssetApiService |
资产 API 定义 |
IDaAssetApiParamService |
API 参数管理 |
IDaAssetGeoService |
地理空间资产 |
IDaAssetGisService |
GIS 资产 |
IDaAssetVideoService |
视频资产 |
IDaAssetFilesService |
文件资产 |
IDaAssetThemeRelService |
资产-主题关联 |
IDaAssetProjectRelService |
资产-项目关联 |
IDaDiscoveryTableService |
发现表管理 |
IDaDiscoveryTaskService |
发现任务管理 |
5.2 数据血缘¶
通过 Neo4j 图数据库追踪 ETL 管线的数据血缘:
| 节点类型 | 标签 | 说明 |
|---|---|---|
TableNode |
@Node("Table") |
数据表节点 |
TaskNode |
@Node("Task") |
ETL 任务节点 |
| 关系类型 | 方向 | 说明 |
|---|---|---|
TABLE_TO_TASK |
Table → Task | 表作为任务输入 |
TASK_TO_TABLE |
Task → Table | 任务输出到表 |
关系属性:taskId, taskCode, datasourceHostPort, tableName。
LineageDataService 支持 Cypher 查询遍历上下游 Table → Task → Table 链路。
六、数据服务(API 发布)¶
6.1 动态 API 注册¶
MappingHandlerMapping 使用 Spring 的 RequestMappingHandlerMapping 在运行时动态注册/注销 REST API:
- URL 模式:
/services/{version}/{path}(如/services/v1.0.0/user/1) - 内部存储:
ConcurrentHashMap<String, DsApiDO> mappings - 认证:
@DsCheckClientToken注解 +ApiJwtUtilJWT 验证
6.2 API 服务类型¶
| 类型码 | 说明 |
|---|---|
| 1 | 数据服务(数据表/SQL 查询) |
| 2 | 模型数据服务 |
| 3 | 第三方 API 代理 |
| 4 | 文件服务 |
6.3 响应类型¶
| 类型码 | 说明 |
|---|---|
| 1 | 详情(单条记录) |
| 2 | 列表 |
| 3 | 分页 |
6.4 核心功能¶
- SQL 解析:
DsApiServiceImpl.sqlParse()使用 JSqlParser 解析 SQL,提取请求参数和响应字段 - SQL 生成:
sqlJdbcNamedParameterBuild()从表单配置生成带命名参数的 SQL - 多数据库 SQL 生成:支持 Kingbase8、PostgreSQL、SQL Server、MySQL、Oracle
- 在线测试:
serviceTesting()支持分页/列表/详情响应测试 - 异步日志:
AsyncTask.doTask(DsApiLogDO)异步记录 API 调用日志
七、数据建模¶
7.1 模型类型¶
| 类型码 | 说明 |
|---|---|
| 1 | 逻辑模型 |
| 2 | 物理模型(从数据源导入) |
7.2 模型分层¶
| 层 | 目标用户 | 说明 |
|---|---|---|
| 公共层(Public Layer) | 数据开发者 | 明细模型 |
| 应用层(Application Layer) | 业务/分析师 | 聚合模型 |
7.3 核心功能¶
- 逻辑模型管理:
DpModelServiceImpl提供 CRUD、发布状态管理、版本控制 - 数据元管理:数据元标准定义、码表映射、数据元-资产关联
- 物化视图:
DpModelMaterializedController管理物化视图 - 树形结构:公共层(业务分类)+ 应用层(主题域)
八、数据治理¶
8.1 数据分类分级¶
qdata-module-dg 提供完整的数据分类分级体系:
| 功能 | Controller | 说明 |
|---|---|---|
| 数据分类 | DgDataCategoryController |
数据分类目录 |
| 分类子类 | DgDataCategoryCatController |
分类细分 |
| 数据级别 | DgDataLevelController |
数据分级 |
| 敏感级别 | DgSensitiveLevelController |
敏感数据分级 |
8.2 数据脱敏¶
| 功能 | Controller | 说明 |
|---|---|---|
| 脱敏规则 | DgDesensitizeRuleController |
脱敏规则定义 |
| 脱敏区间 | DgDesensitizeIntervalController |
脱敏区间配置 |
| 脱敏字段 | DgDesensitizeAssetcolumnController |
资产字段脱敏关联 |
| 白名单 | DgDesensitizeWhitelistController |
脱敏白名单 |
| 用户关联 | DgDesensitizeUserRelController |
用户-脱敏规则关联 |
九、AI 智能能力¶
9.1 AI 平台支持¶
AiPlatformEnum 支持 19 个 AI 平台:
国内(11 个):通义千问(DashScope)、文心一言(百度)、DeepSeek、智谱、星火(讯飞)、豆包(字节)、混元(腾讯)、SiliconFlow、MiniMax、Moonshot(Kimi)、百川。
海外(8 个):OpenAI、Azure OpenAI、Anthropic(Claude)、Gemini(Google)、Ollama、StableDiffusion、Midjourney、Suno、Grok。
9.2 Text2SQL¶
StatisticsPromptBuilder 实现自然语言到 SQL 的转换:
- 星型模型支持:事实表 + 维度表 + 事实-维度关联
- 数据库方言感知:MySQL、DM8/Oracle、Kingbase8、PostgreSQL、SQL Server、Doris
- 双回复模式:
CHART模式:返回 SQL + 维度 + 度量 + 时间粒度(用于图表渲染)QA模式:返回文本或 SQL(用于问答)
9.3 事实-维度表匹配¶
MatchPromptBuilder 使用 AI 自动匹配事实表与维度表的关联关系:
- 输入:事实表(表名、别名、描述、列、主键)+ 维度表列表
- 输出:JSON 格式的关联关系(维度表、事实列名、维度列名、匹配原因)
9.4 独立部署¶
qdata-service-ai(端口 8087)独立部署,要求 JDK 17,通过 DualServiceImpl 实现 12 个空桩接口,使其可脱离其他微服务独立运行。