返回博客

指标开发指南

指标(Metric)是 coomia-dip 平台中衡量业务表现的核心数据单元。本文介绍如何定义、计算、存储和可视化业务指标,包括基础指标、派生指标、复合指标三种类型。涵盖指标的 YAML 声明、Python 自定义计算、实时/批量计算模式、以及指标监控告警。

Coomia发布于 2026年1月18日10 分钟阅读
分享本文Twitter / X

系列:S12 开发者教程 · 第 9 篇 | 难度:中级 | 阅读时间:15 分钟

指标开发指南

#TL;DR

指标(Metric)是 coomia-dip 平台中衡量业务表现的核心数据单元。本文介绍如何定义、计算、存储和可视化业务指标,包括基础指标、派生指标、复合指标三种类型。涵盖指标的 YAML 声明、Python 自定义计算、实时/批量计算模式、以及指标监控告警。

#1. 指标体系概述

#1.1 什么是指标?

指标是对 Ontology 对象属性的量化度量。coomia-dip 的指标系统运行在 Control Layer(Control Layer)与 Data Layer(Data Layer)之上,通过 gRPC 提供服务。

Code
指标定义 (Schema)
    │
    ▼
指标计算 (Engine)
    ├── 实时计算 (Streaming)
    ├── 批量计算 (Batch)
    └── 按需计算 (On-Demand)
    │
    ▼
指标存储 (Doris OLAP)
    │
    ▼
指标消费 (Dashboard / API / Alert)

#1.2 指标类型

类型说明示例
基础指标(Base)直接从数据聚合销售总额、客户数量
派生指标(Derived)基于基础指标计算客单价 = 销售总额 / 订单数
复合指标(Composite)多维度组合部门人均产值

#1.3 指标维度

维度是指标的分析切面:

YAML
dimensions:
  - time: [day, week, month, quarter, year]
  - region: [country, province, city]
  - department: [business_unit, team]
  - product: [category, subcategory]
  - customer: [tier, industry, region]

#2. 环境准备

Python
from ontology_sdk import OntoPlatform

platform = OntoPlatform(
    control_plane_url="localhost:50051",
    data_plane_url="localhost:50052"
)

metric_manager = platform.metrics

#3. 定义基础指标

#3.1 YAML 声明

YAML
# metrics/total_revenue.yaml
name: total_revenue
display_name: 销售总收入
description: 所有已完成订单的总金额
category: financial
unit: CNY
precision: 2

source:
  object_type: Order
  filter:
    status: completed
  aggregation:
    function: SUM
    field: total_amount

dimensions:
  - name: time
    field: completed_at
    granularities: [day, week, month, quarter, year]
  - name: region
    field: customer.region
  - name: product_category
    field: items.category

schedule:
  compute_mode: batch
  cron: "0 1 * * *"    # 每天凌晨 1 点计算
  timezone: Asia/Shanghai

cache:
  ttl: 3600             # 缓存 1 小时
  warm_on_compute: true  # 计算后自动预热缓存

alerts:
  - name: revenue_drop
    condition: "current < previous * 0.8"
    comparison: day_over_day
    channel: dingtalk
    message: "日收入较昨日下降超过 20%"

#3.2 Python API 定义

Python
from ontology_sdk.metrics import MetricBuilder, Aggregation, Dimension

metric = (
    MetricBuilder("total_revenue")
    .display_name("销售总收入")
    .description("所有已完成订单的总金额")
    .category("financial")
    .unit("CNY")
    .precision(2)
    .source(
        object_type="Order",
        filter={"status": "completed"},
        aggregation=Aggregation.SUM("total_amount")
    )
    .dimension(Dimension.time("completed_at", ["day", "week", "month", "quarter", "year"]))
    .dimension(Dimension.category("customer.region", name="region"))
    .dimension(Dimension.category("items.category", name="product_category"))
    .schedule(compute_mode="batch", cron="0 1 * * *")
    .cache(ttl=3600)
    .build()
)

metric_manager.register(metric)

#3.3 更多基础指标

Python
# 客户数量指标
customer_count = (
    MetricBuilder("active_customer_count")
    .display_name("活跃客户数")
    .source(
        object_type="Customer",
        filter={"status": "active"},
        aggregation=Aggregation.COUNT()
    )
    .dimension(Dimension.category("tier"))
    .dimension(Dimension.category("industry"))
    .dimension(Dimension.time("last_active_at", ["month", "quarter"]))
    .build()
)

# 平均订单金额
avg_order_amount = (
    MetricBuilder("avg_order_amount")
    .display_name("平均订单金额")
    .source(
        object_type="Order",
        filter={"status": "completed"},
        aggregation=Aggregation.AVG("total_amount")
    )
    .dimension(Dimension.time("completed_at", ["day", "month"]))
    .build()
)

# 订单数量
order_count = (
    MetricBuilder("order_count")
    .display_name("订单数量")
    .source(
        object_type="Order",
        filter={"status": "completed"},
        aggregation=Aggregation.COUNT()
    )
    .dimension(Dimension.time("completed_at", ["day", "month"]))
    .build()
)

metric_manager.register_batch([customer_count, avg_order_amount, order_count])

#4. 定义派生指标

#4.1 基于公式的派生

YAML
# metrics/customer_unit_price.yaml
name: customer_unit_price
display_name: 客单价
description: 平均每位客户的消费金额
category: financial
type: derived

formula:
  expression: "total_revenue / active_customer_count"
  dependencies:
    - total_revenue
    - active_customer_count

dimensions:
  - name: time
    granularities: [month, quarter]
  - name: region
Python
from ontology_sdk.metrics import DerivedMetricBuilder

customer_unit_price = (
    DerivedMetricBuilder("customer_unit_price")
    .display_name("客单价")
    .formula("total_revenue / active_customer_count")
    .dependencies(["total_revenue", "active_customer_count"])
    .unit("CNY")
    .precision(2)
    .build()
)

metric_manager.register(customer_unit_price)

#4.2 复杂公式

Python
# 毛利率
gross_margin = (
    DerivedMetricBuilder("gross_margin_rate")
    .display_name("毛利率")
    .formula("(total_revenue - total_cost) / total_revenue * 100")
    .dependencies(["total_revenue", "total_cost"])
    .unit("%")
    .precision(1)
    .build()
)

# 同比增长率
yoy_growth = (
    DerivedMetricBuilder("revenue_yoy_growth")
    .display_name("收入同比增长率")
    .formula("(current_period - same_period_last_year) / same_period_last_year * 100")
    .time_comparison(
        current="total_revenue",
        offset="1 year",
        granularity="month"
    )
    .unit("%")
    .precision(1)
    .build()
)

# 环比增长率
mom_growth = (
    DerivedMetricBuilder("revenue_mom_growth")
    .display_name("收入环比增长率")
    .formula("(current_period - previous_period) / previous_period * 100")
    .time_comparison(
        current="total_revenue",
        offset="1 month",
        granularity="month"
    )
    .unit("%")
    .precision(1)
    .build()
)

#5. 自定义计算逻辑

#5.1 Python 自定义指标

Python
from ontology_sdk.metrics import custom_metric, MetricContext
from pydantic import BaseModel
from typing import Optional


class HealthScoreResult(BaseModel):
    score: float
    grade: str      # A / B / C / D / F
    factors: dict


@custom_metric(
    name="project_health_score",
    display_name="项目健康度评分",
    description="基于多维度数据综合计算的项目健康度",
    compute_mode="on_demand",
    cache_ttl=1800,
)
def compute_project_health(ctx: MetricContext) -> list[HealthScoreResult]:
    """计算所有进行中项目的健康度评分"""

    projects = ctx.oql.execute("""
        FIND Project
        WHERE status = 'in_progress'
        INCLUDE
            TRAVERSE has_task -> Task
            AGGREGATE
                COUNT(*) AS total_tasks,
                COUNT(CASE WHEN Task.status = 'done' THEN 1 END) AS done_tasks,
                COUNT(CASE WHEN Task.due_date < NOW() AND Task.status != 'done' THEN 1 END) AS overdue_tasks
        SELECT name, budget, end_date, total_tasks, done_tasks, overdue_tasks
    """)

    results = []
    for project in projects:
        # 进度因子 (0-30 分)
        completion = project.done_tasks / max(project.total_tasks, 1)
        time_elapsed = (ctx.now() - project.start_date).days
        expected_duration = (project.end_date - project.start_date).days
        time_ratio = time_elapsed / max(expected_duration, 1)
        progress_score = max(0, 30 * (1 - abs(completion - time_ratio)))

        # 逾期因子 (0-30 分)
        overdue_ratio = project.overdue_tasks / max(project.total_tasks, 1)
        overdue_score = 30 * (1 - overdue_ratio)

        # 预算因子 (0-20 分)
        budget_used = ctx.get_metric("project_budget_used", filter={"project_rid": project.rid})
        budget_ratio = budget_used / max(project.budget, 1) if budget_used else 0
        budget_score = 20 if budget_ratio < time_ratio * 1.1 else max(0, 20 * (1 - (budget_ratio - time_ratio)))

        # 团队因子 (0-20 分)
        team_velocity = ctx.get_metric("team_velocity", filter={"project_rid": project.rid})
        team_score = min(20, 20 * (team_velocity / 10)) if team_velocity else 10

        total_score = progress_score + overdue_score + budget_score + team_score
        grade = "A" if total_score >= 85 else "B" if total_score >= 70 else "C" if total_score >= 55 else "D" if total_score >= 40 else "F"

        results.append(HealthScoreResult(
            score=round(total_score, 1),
            grade=grade,
            factors={
                "progress": round(progress_score, 1),
                "overdue": round(overdue_score, 1),
                "budget": round(budget_score, 1),
                "team": round(team_score, 1),
            }
        ))

    return results

#6. 查询与消费指标

#6.1 基本查询

Python
# 查询单一维度
revenue = metric_manager.query(
    "total_revenue",
    time_range=("2025-01-01", "2025-12-31"),
    granularity="month"
)

for point in revenue.data_points:
    print(f"{point.time}: ¥{point.value:,.2f}")

# 多维度交叉查询
breakdown = metric_manager.query(
    "total_revenue",
    time_range=("2025-01-01", "2025-12-31"),
    granularity="month",
    dimensions=["region", "product_category"]
)

for group in breakdown.groups:
    print(f"\n{group.dimension_values}:")
    for point in group.data_points:
        print(f"  {point.time}: ¥{point.value:,.2f}")

#6.2 对比查询

Python
# 同比对比
comparison = metric_manager.compare(
    "total_revenue",
    current_period=("2025-01-01", "2025-12-31"),
    previous_period=("2024-01-01", "2024-12-31"),
    granularity="month"
)

for point in comparison.data_points:
    change = ((point.current - point.previous) / point.previous * 100
              if point.previous else 0)
    print(f"{point.time}: ¥{point.current:,.0f} "
          f"(去年: ¥{point.previous:,.0f}, 变化: {change:+.1f}%)")

#6.3 排名查询

Python
# 区域收入排名
ranking = metric_manager.rank(
    "total_revenue",
    dimension="region",
    time_range=("2025-01-01", "2025-03-31"),
    top_n=10,
    order="desc"
)

for i, entry in enumerate(ranking.entries, 1):
    print(f"  #{i} {entry.dimension_value}: ¥{entry.value:,.0f}")

#7. 实时指标

#7.1 配置实时计算

YAML
# metrics/realtime_order_rate.yaml
name: realtime_order_rate
display_name: 实时下单速率
type: base
compute_mode: streaming

source:
  type: kafka
  topic: order-events
  filter:
    event_type: order_created
  aggregation:
    function: COUNT
    window:
      type: sliding
      size: 5m
      slide: 1m

dimensions:
  - name: region
    field: customer.region

#7.2 流式指标消费

Python
# 订阅实时指标
async def monitor_orders():
    async for update in metric_manager.subscribe("realtime_order_rate"):
        print(f"[{update.timestamp}] 订单速率: {update.value}/分钟")
        if update.value < 10:
            print("  告警: 下单速率异常低!")

import asyncio
asyncio.run(monitor_orders())

#8. 指标告警

#8.1 告警规则

Python
from ontology_sdk.metrics import AlertRule, AlertCondition

# 绝对值告警
alert1 = AlertRule(
    name="revenue_floor",
    metric="total_revenue",
    condition=AlertCondition.less_than(100000),
    granularity="day",
    channel="dingtalk",
    recipients=["sales-team"],
    message="日收入低于 10 万元"
)

# 环比告警
alert2 = AlertRule(
    name="revenue_drop",
    metric="total_revenue",
    condition=AlertCondition.period_over_period_drop(threshold=0.2),
    comparison="day_over_day",
    channel="email",
    recipients=["cfo@company.com"],
    message="日收入环比下降超过 20%"
)

# 异常检测告警
alert3 = AlertRule(
    name="order_anomaly",
    metric="order_count",
    condition=AlertCondition.anomaly_detection(
        method="z_score",
        threshold=3.0,
        training_window="30d"
    ),
    channel="dingtalk",
    recipients=["ops-team"],
    message="订单量出现统计异常"
)

metric_manager.register_alerts([alert1, alert2, alert3])

#9. 指标目录与元数据

#9.1 指标浏览

Python
# 列出所有指标
all_metrics = metric_manager.list(category="financial")
for m in all_metrics:
    print(f"{m.name}: {m.display_name} [{m.type}]")
    print(f"  单位: {m.unit}, 维度: {[d.name for d in m.dimensions]}")

# 查看指标详情
detail = metric_manager.describe("total_revenue")
print(f"名称: {detail.display_name}")
print(f"描述: {detail.description}")
print(f"计算方式: {detail.computation}")
print(f"依赖: {detail.dependencies}")
print(f"最近计算: {detail.last_computed_at}")
print(f"数据新鲜度: {detail.freshness}")

# 指标血缘
lineage = metric_manager.get_lineage("customer_unit_price")
print("依赖关系:")
for dep in lineage.upstream:
    print(f"  <- {dep.name} ({dep.type})")
print("被依赖:")
for dep in lineage.downstream:
    print(f"  -> {dep.name} ({dep.type})")

#10. 完整实战:构建销售指标体系

Python
from ontology_sdk import OntoPlatform
from ontology_sdk.metrics import MetricBuilder, DerivedMetricBuilder, Aggregation, Dimension

platform = OntoPlatform(
    control_plane_url="localhost:50051",
    data_plane_url="localhost:50052"
)
mm = platform.metrics

# === 基础指标 ===
base_metrics = [
    MetricBuilder("total_revenue")
        .display_name("销售总收入").unit("CNY")
        .source("Order", {"status": "completed"}, Aggregation.SUM("total_amount"))
        .dimension(Dimension.time("completed_at", ["day", "month", "quarter", "year"]))
        .dimension(Dimension.category("customer.region", name="region"))
        .build(),

    MetricBuilder("total_cost")
        .display_name("总成本").unit("CNY")
        .source("Order", {"status": "completed"}, Aggregation.SUM("cost"))
        .dimension(Dimension.time("completed_at", ["day", "month", "quarter"]))
        .build(),

    MetricBuilder("order_count")
        .display_name("订单数量").unit("单")
        .source("Order", {"status": "completed"}, Aggregation.COUNT())
        .dimension(Dimension.time("completed_at", ["day", "month"]))
        .build(),

    MetricBuilder("new_customer_count")
        .display_name("新增客户数").unit("个")
        .source("Customer", {}, Aggregation.COUNT())
        .dimension(Dimension.time("created_at", ["day", "month"]))
        .build(),
]

# === 派生指标 ===
derived_metrics = [
    DerivedMetricBuilder("gross_profit")
        .display_name("毛利润").unit("CNY")
        .formula("total_revenue - total_cost")
        .dependencies(["total_revenue", "total_cost"])
        .build(),

    DerivedMetricBuilder("gross_margin_rate")
        .display_name("毛利率").unit("%")
        .formula("(total_revenue - total_cost) / total_revenue * 100")
        .dependencies(["total_revenue", "total_cost"])
        .build(),

    DerivedMetricBuilder("avg_order_value")
        .display_name("平均订单金额").unit("CNY")
        .formula("total_revenue / order_count")
        .dependencies(["total_revenue", "order_count"])
        .build(),
]

# 批量注册
mm.register_batch(base_metrics + derived_metrics)

# === 查询示例 ===
print("=== 2025 年月度销售概览 ===")
revenue_data = mm.query("total_revenue", time_range=("2025-01-01", "2025-12-31"), granularity="month")
margin_data = mm.query("gross_margin_rate", time_range=("2025-01-01", "2025-12-31"), granularity="month")

for rev, margin in zip(revenue_data.data_points, margin_data.data_points):
    print(f"  {rev.time}: 收入 ¥{rev.value:,.0f} | 毛利率 {margin.value:.1f}%")

# === 区域排名 ===
print("\n=== Q1 区域收入排名 ===")
ranking = mm.rank("total_revenue", dimension="region",
                  time_range=("2025-01-01", "2025-03-31"), top_n=5)
for i, entry in enumerate(ranking.entries, 1):
    print(f"  #{i} {entry.dimension_value}: ¥{entry.value:,.0f}")

#Key Takeaways

  1. 三级指标体系:基础指标从数据聚合,派生指标用公式组合,复合指标多维度分析
  2. 声明式定义:优先使用 YAML 定义指标,复杂逻辑用 Python 自定义计算
  3. 多种计算模式:批量计算适合历史统计,实时计算适合监控,按需计算适合探索
  4. 维度灵活:时间、区域、业务类型等多维度自由交叉分析
  5. 告警驱动:指标异常自动触发告警,支持绝对值、环比、异常检测
  6. 血缘追踪:指标依赖关系可视化,上下游影响一目了然

#Next Article

下一篇:S12-10 Dashboard 开发指南 — 学习如何基于指标构建可视化仪表盘。

Tags: 指标 Metric KPI OLAP 数据分析 告警 coomia-dip