指标系统的分层架构设计与实现

指标系统的核心在于其分层架构设计。让我深入分析这个系统的技术实现细节,从架构设计理念到实现。

系统整体架构

架构概览

指标系统采用了分层架构设计,整个系统分为数据源层、指标管理层、处理层、存储层和服务层五个核心层次,其中指标管理层进一步细分为原子指标、派生指标和复合指标三个层次,形成了完整的指标分层体系。

架构流程

1. 指标分层设计

  • 原子指标层(Atomic Metrics): 系统的基础层,直接对应业务事件的最小度量单元

  • 派生指标层(Derived Metrics): 通过组合原子指标、修饰词和时间约束形成更具体的业务指标

  • 复合指标层(Composite Metrics): 支持多个指标的复杂运算,实现业务逻辑的灵活组合

2. 数据处理流程优化

  • 数据抽取阶段: Datax 负责从各种数据源抽取原始数据到 Hive 数据仓库

  • 计算处理阶段: Spark 计算引擎根据原子指标、派生指标、复合指标的配置进行分布式计算

  • 数据转换阶段: Datax 在 Spark 计算完成后,根据业务定义的指标组合和调度周期,将计算结果转换并更新到业务数据库

3. 调度驱动机制

  • 调度配置驱动 Datax 数据转换的执行时机

  • 指标组合配置定义数据转换的业务逻辑

  • 支持天调度、月调度、年调度等多种调度周期

4. 多数据源支持

  • 后端数据库:结构化业务数据

  • 日志数据:用户行为日志

  • API 数据:外部接口数据

技术栈

  • 数据抽取: Datax(原始数据同步)

  • 数据转换: Datax(指标结果转换)

  • 计算引擎: Spark(分布式指标计算)

  • 数据存储: Hive(原始数据)、Elasticsearch(指标结果)、关系数据库(业务数据)

  • API 服务: RESTful API(实时查询)

  • 可视化: 数据可视化面板

  • 配置管理: 指标定义、指标组合、调度配置

  • 调度系统: 基于的 Quartz 定时调度系统

分层架构的设计理念

架构层次划分

指标系统的分层架构主要包含以下几个层次:

  1. 原子指标层(Atomic Metrics) - 系统的基础层

  2. 派生指标层(Derived Metrics) - 业务逻辑层

  3. 复合指标层(Composite Metrics) - 复杂计算层

  4. 应用层(Application Layer) - 业务应用层

原子指标层设计

原子指标是系统的基础,它直接对应业务事件的最小度量单元。

核心数据结构

1
2
3
4
5
6
7
8
9
10
11
// 原子指标的核心属性
public class MetricsBaseDto {
    private String nameCN;        // 指标中文名
    private String nameEN;        // 指标英文名  
    private String dataType;      // 数据类型
    private String formula;       // 计算公式
    private String dimension;     // 指标维度
    private String cron;          // 调度表达式
    private Integer metricsType;  // 指标类型(1-原子指标)
}

数据库设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
-- 原子指标表结构
CREATE TABLE `metrics_atomic` (
    `id` int unsigned NOT NULL AUTO_INCREMENT COMMENT '原子指标ID',
    `name_cn` varchar(100) NOT NULL COMMENT '指标名称',
    `name_en` varchar(40) NOT NULL COMMENT '指标英文名称',
    `data_type` varchar(40) NOT NULL COMMENT '数据类型',
    `desc` varchar(1024) DEFAULT NULL COMMENT '描述',
    `datasource_id` int(11) NOT NULL COMMENT '数据源ID',
    `db_name` varchar(255) DEFAULT NULL COMMENT '数据库名',
    `table_name` varchar(255) NOT NULL COMMENT '表名',
    `field_name` varchar(200) NOT NULL COMMENT '字段名',
    `field_type` varchar(40) NOT NULL COMMENT '字段类型',
    `formula` varchar(1024) DEFAULT NULL COMMENT '公式',
    `dimension` varchar(255) DEFAULT NULL COMMENT '指标维度',
    `constraint_id` int(11) DEFAULT NULL COMMENT '修饰词ID',
    `time_constraint_id` int(11) DEFAULT NULL COMMENT '时间修饰ID',
    `cron` varchar(20) DEFAULT NULL COMMENT '调度表达式',
    `cron_type` tinyint DEFAULT NULL COMMENT '调度类型',
    `job_id` int(11) DEFAULT NULL COMMENT 'etl作业ID',
    `status` tinyint DEFAULT 0 COMMENT '是否启用',
    `create_time` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP,
    `update_time` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='原子指标';

原子指标的特点

  1. 不可再拆分:原子指标是业务定义中不可再拆分的指标

  2. 明确业务含义:具有明确的业务含义和名称

  3. 直接数据源:直接对应数据源中的具体字段

  4. 基础计算单元:为派生指标和复合指标提供基础数据

派生指标层设计

派生指标通过组合原子指标、修饰词和时间约束,形成更具体的业务指标。

核心数据结构

1
2
3
4
5
6
7
// 派生指标DTO
public class DeriveMetricsDto extends MetricsBaseDto {
    private Integer atomicMetricsId;    // 原子指标ID
    private List<MetricsConstraintPo> constraints;  // 修饰词列表
    private Integer timeConstraintId;   // 时间修饰ID
}

数据库设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
-- 派生指标表结构
CREATE TABLE `metrics_derive` (
    `id` int unsigned NOT NULL AUTO_INCREMENT COMMENT '指标ID',
    `name_cn` varchar(100) NOT NULL COMMENT '指标名称',
    `name_en` varchar(40) NOT NULL COMMENT '指标英文名称',
    `data_type` varchar(40) NOT NULL COMMENT '数据类型',
    `desc` varchar(1024) DEFAULT NULL COMMENT '描述',
    `atomic_metrics_id` int unsigned NOT NULL COMMENT '原子指标ID',
    `formula` varchar(1024) DEFAULT NULL COMMENT '公式',
    `dimension` varchar(255) DEFAULT NULL COMMENT '指标维度',
    `constraint_id` int(11) DEFAULT NULL COMMENT '修饰词ID',
    `time_constraint_id` int(11) DEFAULT NULL COMMENT '时间修饰ID',
    `cron` varchar(20) DEFAULT NULL COMMENT '调度表达式',
    `cron_type` tinyint DEFAULT NULL COMMENT '调度类型',
    `job_id` int(11) DEFAULT NULL COMMENT 'etl作业ID',
    `status` tinyint DEFAULT 0 COMMENT '是否启用',
    `create_time` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP,
    `update_time` timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    PRIMARY KEY (`id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='衍生指标';

派生指标的计算公式

派生指标的计算公式遵循以下模式:

1
2
派生指标 = 1个原子指标 + 多个修饰词(可选) + 时间周期

例如:

  • 原子指标:支付金额

  • 修饰词:海外、买家

  • 时间周期:最近 7 天

  • 派生指标:最近 7 天海外买家支付金额

复合指标层设计

复合指标支持多个指标的复杂运算,实现业务逻辑的灵活组合。

核心数据结构

1
2
3
4
5
6
7
// 复合指标DTO
public class CompositeMetricsDto extends MetricsBaseDto {
    private List<Integer> dependentMetricsIds;  // 依赖指标ID列表
    private String compositeFormula;            // 复合计算公式
    private Map<String, Object> parameters;     // 计算参数
}

复合指标的计算逻辑

复合指标支持多种计算模式:

  1. 算术运算:加减乘除等基本运算

  2. 逻辑运算:与或非等逻辑运算

  3. 函数运算:自定义函数运算

  4. 条件运算:基于条件的条件运算

维度管理的设计

维度是指标系统的重要概念,它定义了指标的统计环境。

维度分层设计

1
2
3
4
5
6
7
// 维度定义
public class MetricsDimensionDto {
    private List<MetricsDimensionDto> dimensions;    // 具体维度
    private List<Tag> businessDomain;                // 业务域
    private List<Tag> dataDomain;                    // 数据域
}

维度分类

  1. 业务域维度:面向业务理解的维度分类
  • 用户维度:用户 ID、用户类型、用户等级

  • 产品维度:产品 ID、产品类别、产品品牌

  • 地域维度:国家、省份、城市

  1. 数据域维度:面向技术实现的维度分类
  • 时间维度:年、季、月、周、日

  • 渠道维度:PC 端、移动端、小程序

  • 设备维度:iOS、Android、Web

修饰词系统的设计

修饰词系统是指标系统的一个创新点,它允许对指标进行精细化控制。

修饰词类型

1
2
3
4
5
6
7
8
9
// 修饰词配置
public class MetricsConstraintPo {
    private Integer constraintId;      // 修饰词ID
    private String constraintName;     // 修饰词名称
    private String constraintType;     // 修饰词类型
    private String constraintValue;    // 修饰词值
    private String operator;           // 操作符
}

修饰词分类

  1. 业务修饰词:限定业务场景
  • 用户类型:新用户、老用户、VIP 用户

  • 产品类型:热门产品、新品、促销产品

  • 渠道类型:线上渠道、线下渠道

  1. 时间修饰词:限定时间范围
  • 相对时间:最近 7 天、最近 30 天、本月

  • 绝对时间:2024 年 1 月、Q1 季度

  • 特殊时间:工作日、周末、节假日

计算引擎的设计

指标系统的计算引擎采用了 Spark 作为核心,实现了高性能的分布式计算。

Spark 计算引擎初始化

1
2
3
4
5
6
7
8
9
10
11
12
13
// Spark计算引擎初始化
def initSparkContext(): SparkSession = {
  val sparkBuilder = SparkSession.builder()
    .appName("metrics-calculates")
    .config("hive.exec.dynamic.partition", "true")
    .config("spark.default.parallelism", 12)
    .config("spark.executor.cores", executorCpu.toInt)
    .config("spark.executor.memory", executorMem)
    .config("spark.cores.max", executorNum.toInt)
    .enableHiveSupport()
  sparkBuilder.getOrCreate()
}

计算引擎的特点

  1. 批量计算:支持大规模数据的批量处理

  2. 分区优化:通过动态分区提高计算效率

  3. 并行处理:通过并行度配置优化性能

  4. 内存管理:通过内存配置优化资源使用

调度系统的设计

调度系统是指标系统的关键组件,它负责管理指标的计算时机。

调度配置

1
2
3
4
5
6
7
8
9
// 调度配置
public class ScheduleConfig {
    private String cron;              // Cron表达式
    private Integer cronType;         // 调度类型
    private Integer jobId;            // 作业ID
    private String scheduleTime;      // 调度时间
    private String runTime;           // 运行时间
}

调度类型

  1. 天调度:每日执行的指标
  • 适用场景:日活跃用户数、日交易额

  • 执行时间:每天凌晨 1 点

  1. 月调度:每月执行的指标
  • 适用场景:月活跃用户数、月交易额

  • 执行时间:每月 1 号凌晨 2 点

  1. 年调度:每年执行的指标
  • 适用场景:年活跃用户数、年交易额

  • 执行时间:每年 1 月 1 号凌晨 3 点

数据存储的设计

指标结果采用多级存储策略,既保证了数据的完整性,又优化了查询性能。

存储表设计

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
-- 指标最近调度运算结果
CREATE TABLE IF NOT EXISTS metrics.ms_metrics_dimension_result (
  `id` string COMMENT '唯一ID流水号',
  `result_id` string COMMENT '数据唯一ID',
  `name_en` string COMMENT '指标英文名',
  `name_cn` string COMMENT '指标中文名',
  `metrics_type` int COMMENT '指标类型',
  `metrics_value` string COMMENT '指标值',
  `dimension_list` string COMMENT '维度列表',
  `dimension` string COMMENT '维度名称',
  `dimension_value` string COMMENT '维度值',
  `job_id` int COMMENT '作业ID',
  `run_or_execution_id` int COMMENT '运行ID或者执行ID',
  `is_schedule` int COMMENT '是否调度',
  `schedule_time` bigint COMMENT '计算调度时间',
  `run_time` bigint COMMENT '计算执行时间',
  `write_time` bigint COMMENT '计算写入时间',
  `dimension_values` string COMMENT '维度值字典'
) PARTITIONED BY (schedule_type_time string, metrics_id int)
STORED AS ORC
tblproperties("orc.compress"="SNAPPY");
-- 指标历史调度运算结果
CREATE TABLE IF NOT EXISTS metrics.ms_metrics_dimension_result (
  `id` string COMMENT '唯一ID流水号',
  `result_id`string COMMENT '数据唯一ID',
  `name_en` string COMMENT '指标英文名',
  `name_cn` string COMMENT '指标中文名',
  `metrics_type` int COMMENT '指标类型',
  `metrics_value` string COMMENT '指标值',
  `dimension_list` string COMMENT '维度列表',
  `dimension` string COMMENT '维度名称',
  `dimension_value` string COMMENT '维度值',
  `job_id` int COMMENT '作业ID',
  `run_or_execution_id` int COMMENT '运行ID或者执行ID',
  `is_schedule` int COMMENT '是否调度',
  `write_time` bigint COMMENT '计算写入时间',
  `schedule_type_time` string COMMENT '计算调度时间类型',
  `dimension_values` string COMMENT '维度值字典'
) PARTITIONED BY (schedule_time bigint, run_time bigint, metrics_id int)
STORED AS ORC
tblproperties("orc.compress"="SNAPPY");

存储策略

  1. 历史存储:保存所有历史计算结果,用于数据追溯

  2. 最新存储:只保存最新计算结果,用于快速查询

  3. 分区存储:按时间和指标 ID 分区,提高查询效率

总结:架构的可替换性与扩展性

指标系统的分层架构设计不仅实现了功能的完整性和性能的优化,更重要的是体现了高度的可替换性扩展性。这种设计使得系统能够适应不同的技术环境和业务需求变化。

数据源的可替换性

当前实现:hive 数据库 可替换方案

  • 关系型数据库:MySQL / PostgreSQL / Oracle / SQL Server

  • NoSQL 数据库:MongoDB / Cassandra / Redis / HBase

替换优势

  • 通过统一的数据源配置接口,实现数据源的即插即用

  • 支持多种数据格式和协议,适应不同的数据接入需求

  • 数据抽取逻辑与具体数据源解耦,降低迁移成本

调度引擎的可替换性

当前实现:自研的调度系统 可替换方案:Apache Airflow 、 Apache DolphinScheduler 、XXL-Job

替换优势

  • 调度配置与具体调度引擎解耦,支持多种调度策略

  • 支持分布式调度、容错调度、动态调度等高级特性

  • 可根据业务规模和复杂度选择合适的调度方案

计算引擎的可替换性

当前实现:Apache Spark 可替换方案

  • 批处理引擎:MapReduce / Tez

  • 实时计算引擎:Flink

  • SQL 计算引擎:Spark SQL / Clickhouse / Drios

替换优势

  • 计算逻辑与具体计算引擎解耦,支持多种计算模式

  • 可根据数据规模、实时性要求、计算复杂度选择合适的引擎

  • 支持混合计算模式,不同指标使用不同的计算引擎

通过这种可替换的架构设计,指标系统不仅能够满足当前的技术需求,更能够适应未来的技术发展和业务变化,实现长期的技术价值和业务价值。