跳转至

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 执行流程

  1. QualityTaskExecutorServiceImpl.executeTask(taskId) 提交异步执行(ThreadPoolTaskExecutor + Redis 去重锁)
  2. 查询任务 → 生成批次号 → 获取数据源对象 → 按数据源 ID 分组
  3. 对每个数据源:连接 → 加载规则 → 通过 RuleTaskExecutor 逐条执行规则
  4. 错误数据存入 MongoDB(CheckErrorData 文档)
  5. 支持 3 种错误数据修改方式:
  6. updateType=1:修改数据 + 更新物理表
  7. updateType=2:修改备注
  8. 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 注解 + ApiJwtUtil JWT 验证

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 个空桩接口,使其可脱离其他微服务独立运行。