← 返回

数据开放平台架构设计(五):监控与治理

上一篇: 数据开放平台(四):数据转换与加工

数据开放平台上线之后,更怕两件事:一是出了问题不知道,二是知道了问题找不到原因。API 响应慢了,是 SQL 没走索引还是数据库负载高?某个调用方报错,是参数传错了还是数据源挂了?某个字段被改了,谁改的?什么时候改的?

这些问题靠翻代码加日志是搞不定的,需要系统化的监控和治理能力。这一篇聊四个核心能力:调用链追踪慢查询告警数据血缘审计日志


一、调用链追踪

每个 API 调用都要记录完整的执行链路:从请求进来到响应返回,中间经历了哪些步骤,每步耗时多少。

TraceId 机制

每个请求分配一个全局唯一的 TraceId,贯穿整个调用链。后续所有日志、指标、告警都带上这个 TraceId,方便关联查询。

@Component
@Order(-100) // 最先执行
public class TraceHandler implements PipelineHandler {

    private static final String TRACE_ID = "X-Trace-Id";

    @Override
    public void before(SqlApiConfig config, Map<String, Object> params,
                        PipelineContext context) {
        // 优先用调用方传来的 TraceId(链路透传)
        String traceId = RequestContext.getHeader(TRACE_ID);
        if (traceId == null) {
            traceId = generateTraceId();
        }
        context.setTraceId(traceId);

        // 记录请求开始时间
        context.setStartTime(System.currentTimeMillis());

        // TraceId 放入 MDC,后续日志自动带上
        MDC.put("traceId", traceId);
        MDC.put("apiId", config.getApiId());
        MDC.put("apiPath", config.getApiPath());
    }

    @Override
    public void after(SqlApiConfig config, Object result,
                       PipelineContext context) {
        long elapsed = System.currentTimeMillis() - context.getStartTime();
        log.info("API 调用完成 traceId={} api={} elapsed={}ms rows={}",
            context.getTraceId(), config.getApiPath(), elapsed,
            result instanceof Collection ? ((Collection<?>) result).size() : 1);
        MDC.clear();
    }

    private String generateTraceId() {
        return UUID.randomUUID().toString().replace("-", "").substring(0, 16);
    }
}

调用链记录

每个步骤的耗时和结果都记录下来,形成完整的调用链:

@Data
public class TraceSpan {
    private String traceId;
    private String spanId;
    private String parentSpanId;
    private String step;          // 步骤名称
    private long startTime;
    private long endTime;
    private long elapsed;         // 耗时 ms
    private String status;        // SUCCESS / ERROR
    private String errorMessage;  // 错误信息
    private Map<String, Object> tags; // 附加信息
}

Pipeline 的每个 Handler 执行前后都记录 Span:

public class TracedPipeline {

    public R<?> execute(SqlApiConfig config, Map<String, Object> params) {
        PipelineContext context = new PipelineContext();
        List<TraceSpan> spans = new ArrayList<>();

        for (PipelineHandler handler : handlers) {
            TraceSpan span = new TraceSpan();
            span.setStep(handler.name());
            span.setStartTime(System.currentTimeMillis());

            try {
                handler.before(config, params, context);
                handler.execute(config, params, context);
                span.setStatus("SUCCESS");
            } catch (Exception e) {
                span.setStatus("ERROR");
                span.setErrorMessage(e.getMessage());
                throw e;
            } finally {
                span.setEndTime(System.currentTimeMillis());
                span.setElapsed(span.getEndTime() - span.getStartTime());
                spans.add(span);
            }
        }

        // 异步写入追踪数据
        traceSink.submit(context.getTraceId(), spans);
        return context.getResult();
    }
}

追踪数据写入 Elasticsearch,支持按 TraceId 查询完整链路,也支持按步骤名称、耗时范围、状态等条件聚合分析。Grafana 面板上可以看到每个 API 的 P50/P90/P99 耗时分布,哪个步骤是瓶颈一目了然。

这里有个细节:参数和结果不要原样全量进日志。手机号、身份证、Token 这类字段要先脱敏,大结果集只记录行数和摘要,否则监控系统反而会变成新的数据泄露点。


二、慢查询告警

SQL 执行是整个 Pipeline 里更耗时的步骤。慢查询不处理,早晚把数据库拖垮。

慢查询检测

在 SQL 执行 Handler 里记录耗时,超过阈值自动标记:

@Component
public class SlowQueryDetector {

    // 默认阈值:3 秒
    private static final long DEFAULT_THRESHOLD = 3000;

    /**
     * 检测并记录慢查询
     */
    public void check(String apiId, String sql, long elapsed,
                       Map<String, Object> params) {
        long threshold = getThreshold(apiId);

        if (elapsed > threshold) {
            SlowQueryRecord record = SlowQueryRecord.builder()
                .apiId(apiId)
                .sql(sql)
                .elapsed(elapsed)
                .threshold(threshold)
                .params(JSONUtil.toJsonStr(maskParams(params)))
                .timestamp(System.currentTimeMillis())
                .build();

            // 写入数据库
            slowQueryMapper.insert(record);

            // 触发告警
            alertService.send(buildAlert(record));

            // 记录指标
            Metrics.counter("slow_query_total",
                "apiId", apiId).increment();
        }
    }

    private long getThreshold(String apiId) {
        // 支持按 API 自定义阈值
        String key = "slow:threshold:" + apiId;
        Long custom = redisTemplate.opsForValue().get(key);
        return custom != null ? custom : DEFAULT_THRESHOLD;
    }
}

告警通知

告警推送到钉钉/飞书群,包含关键信息:

@Component
public class AlertService {

    @Autowired
    private DingTalkClient dingTalkClient;

    public void send(AlertMessage alert) {
        // 告警去重:同一个 API 5 分钟内只告一次
        String dedupKey = "alert:dedup:" + alert.getApiId();
        if (Boolean.TRUE.equals(redisTemplate.hasKey(dedupKey))) {
            return;
        }
        redisTemplate.opsForValue().set(dedupKey, "1", 5, TimeUnit.MINUTES);

        // 构建消息
        String message = String.format(
            "🐌 慢查询告警\n" +
            "API: %s\n" +
            "耗时: %dms (阈值: %dms)\n" +
            "时间: %s\n" +
            "TraceId: %s",
            alert.getApiPath(),
            alert.getElapsed(),
            alert.getThreshold(),
            LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")),
            alert.getTraceId()
        );

        dingTalkClient.send(message);
    }
}

告警消息示例:

🐌 慢查询告警
API: /api/report/user-orders
耗时: 5230ms (阈值: 3000ms)
时间: 2026-04-24 14:30:15
TraceId: a1b2c3d4e5f6g7h8

收到告警后,用 TraceId 去追踪系统查完整链路,定位是 SQL 本身慢还是数据库负载高。

慢查询优化建议

系统还能自动分析慢 SQL,给出优化建议:

@Component
public class SqlAnalyzer {

    public List<String> analyze(String sql, long elapsed) {
        List<String> suggestions = new ArrayList<>();

        String upperSql = sql.toUpperCase();

        // 没有 WHERE 条件的全表扫描
        if (!upperSql.contains("WHERE") && upperSql.contains("SELECT")) {
            suggestions.add("SQL 没有 WHERE 条件,可能是全表扫描,建议加筛选条件");
        }

        // SELECT * 问题
        if (upperSql.contains("SELECT *")) {
            suggestions.add("使用了 SELECT *,建议只查询需要的字段");
        }

        // LIKE '%xxx' 左模糊
        if (upperSql.contains("LIKE '%")) {
            suggestions.add("LIKE 左模糊无法走索引,建议改为右模糊或使用全文索引");
        }

        // 没有 LIMIT 的大结果集
        if (!upperSql.contains("LIMIT") && !upperSql.contains("ROWNUM")) {
            suggestions.add("没有 LIMIT 限制,结果集可能很大,建议加分页");
        }

        // 子查询
        if (upperSql.contains("SELECT") &&
            sql.toUpperCase().split("SELECT").length > 2) {
            suggestions.add("包含子查询,考虑是否可以改写为 JOIN 提升性能");
        }

        return suggestions;
    }
}

三、数据血缘

数据血缘记录每个 API 涉及哪些数据库、哪些表、哪些字段。有了血缘信息,做影响分析就方便了,“我要改 user 表的 name 字段,会影响哪些 API?”

血缘采集

在 SQL 解析阶段自动提取表名和字段名:

@Component
public class DataLineageCollector {

    /**
     * 从 SQL 模板中提取表名和字段名
     */
    public LineageInfo collect(String apiId, String sqlTemplate) {
        LineageInfo info = new LineageInfo();
        info.setApiId(apiId);

        // 提取表名(简化版,生产环境建议用 SQL AST 解析)
        Set<String> tables = extractTables(sqlTemplate);
        info.setTables(tables);

        // 提取字段名
        Set<String> fields = extractFields(sqlTemplate);
        info.setFields(fields);

        // 提取参数引用的字段
        Set<String> paramFields = extractParamFields(sqlTemplate);
        info.setParamFields(paramFields);

        return info;
    }

    private Set<String> extractTables(String sql) {
        Set<String> tables = new HashSet<>();
        // FROM 和 JOIN 后面跟的是表名
        Pattern pattern = Pattern.compile(
            "(?:FROM|JOIN)\\s+(\\w+)", Pattern.CASE_INSENSITIVE);
        Matcher matcher = pattern.matcher(sql);
        while (matcher.find()) {
            tables.add(matcher.group(1).toLowerCase());
        }
        return tables;
    }

    private Set<String> extractFields(String sql) {
        Set<String> fields = new HashSet<>();
        // 提取 #{param} 引用的字段
        Pattern pattern = Pattern.compile("#\\{(\\w+)}");
        Matcher matcher = pattern.matcher(sql);
        while (matcher.find()) {
            fields.add(matcher.group(1));
        }
        return fields;
    }
}

影响分析

当某个表或字段要变更时,查血缘数据就能知道影响范围:

@Service
public class LineageService {

    /**
     * 查询某个表被哪些 API 使用
     */
    public List<String> getApisByTable(String tableName) {
        return lineageMapper.selectApisByTable(tableName);
    }

    /**
     * 查询某个字段被哪些 API 使用
     */
    public List<String> getApisByField(String tableName, String fieldName) {
        return lineageMapper.selectApisByField(tableName, fieldName);
    }

    /**
     * 影响分析:改某个字段会影响哪些 API
     */
    public ImpactReport analyzeImpact(String tableName, String fieldName) {
        List<String> affectedApis = getApisByField(tableName, fieldName);

        ImpactReport report = new ImpactReport();
        report.setTable(tableName);
        report.setField(fieldName);
        report.setAffectedApis(affectedApis);
        report.setAffectedCount(affectedApis.size());
        report.setRiskLevel(affectedApis.size() > 10 ? "HIGH" :
            affectedApis.size() > 3 ? "MEDIUM" : "LOW");

        return report;
    }
}

影响分析报告示例:

📋 字段变更影响分析
表: user
字段: name
影响 API 数: 5
风险等级: MEDIUM

受影响 API:
  1. /api/user/list (GET)
  2. /api/user/detail (GET)
  3. /api/report/user-export (GET)
  4. /api/crm/customer-search (GET)
  5. /api/bi/user-analysis (GET)

四、审计日志

谁在什么时候做了什么操作,必须有记录。审计日志不只是为了出事后追查,也是合规要求,很多行业(金融、医疗、政务)对数据访问有强制审计要求。

审计事件

@Data
@Builder
public class AuditEvent {
    private String traceId;
    private String operator;       // 操作人(调用方标识)
    private String action;         // 操作类型
    private String resourceType;   // 资源类型(API / 配置 / 数据源)
    private String resourceId;     // 资源 ID
    private String detail;         // 操作详情
    private String ipAddress;      // 请求 IP
    private String userAgent;      // 调用方标识
    private Instant timestamp;
    private String result;         // SUCCESS / FAILURE
}

自动采集

用 AOP 自动采集审计日志,业务代码不用手动埋点:

@Aspect
@Component
public class AuditAspect {

    @Autowired
    private AuditLogService auditLogService;

    @Around("@annotation(auditLog)")
    public Object audit(ProceedingJoinPoint pjp, AuditLog auditLog) throws Throwable {
        AuditEvent.AuditEventBuilder builder = AuditEvent.builder()
            .action(auditLog.action())
            .resourceType(auditLog.resourceType())
            .operator(getCurrentUser())
            .ipAddress(RequestContext.getClientIp())
            .timestamp(Instant.now());

        try {
            Object result = pjp.proceed();
            builder.result("SUCCESS")
                .resourceId(extractResourceId(result));
            return result;
        } catch (Exception e) {
            builder.result("FAILURE")
                .detail(e.getMessage());
            throw e;
        } finally {
            auditLogService.log(builder.build());
        }
    }
}

// 使用:在 Controller 方法上加注解
@AuditLog(action = "UPDATE_API_CONFIG", resourceType = "API")
@PutMapping("/{apiId}")
public R<Void> update(@PathVariable Long apiId,
                       @RequestBody SqlApiConfig config) {
    // ...
}

审计日志查询

管理后台提供审计日志查询界面,支持按时间、操作人、操作类型、资源等条件筛选:

@Service
public class AuditQueryService {

    /**
     * 查询审计日志,支持多条件筛选
     */
    public PageResult<AuditEvent> query(AuditQueryRequest request) {
        LambdaQueryWrapper<AuditEvent> wrapper = new LambdaQueryWrapper<>();

        if (request.getOperator() != null) {
            wrapper.eq(AuditEvent::getOperator, request.getOperator());
        }
        if (request.getAction() != null) {
            wrapper.eq(AuditEvent::getAction, request.getAction());
        }
        if (request.getStartTime() != null) {
            wrapper.ge(AuditEvent::getTimestamp, request.getStartTime());
        }
        if (request.getEndTime() != null) {
            wrapper.le(AuditEvent::getTimestamp, request.getEndTime());
        }

        wrapper.orderByDesc(AuditEvent::getTimestamp);
        return auditMapper.selectPage(
            new Page<>(request.getPageNum(), request.getPageSize()), wrapper);
    }
}

五、治理仪表盘

把上面的能力整合到一个仪表盘里,一屏掌握全局:

监控与治理看板可以按数据来源拆开,页面上再组合成同一张运营视图:

flowchart TB API["数据服务运行时"] --> Metrics["Prometheus 指标"] API --> Logs["Elasticsearch 日志"] API --> AuditDB["审计库"] Config["元数据与发布配置"] --> Lineage["血缘分析"] Metrics --> Overview["指标概览"] Metrics --> Trend["QPS、错误率、耗时趋势"] Logs --> Slow["慢查询 TOP"] Logs --> Trace["Trace 查询"] AuditDB --> Audit["审计查询"] Metrics --> Alerts["告警列表"] Lineage --> Impact["字段变更影响分析"] Overview --> Dashboard["治理仪表盘"] Trend --> Dashboard Slow --> Dashboard Trace --> Dashboard Audit --> Dashboard Alerts --> Dashboard Impact --> Dashboard

概览区:今日调用总量、成功率、平均耗时、活跃 API 数、活跃调用方数。

实时监控:QPS 曲线图、耗时分布图(P50/P90/P99)、错误率趋势图。异常时自动标红。

告警列表:近段时间的告警事件,按时间倒序,支持一键跳转到追踪详情。

慢查询 TOP10:更慢的 10 个 API,展示平均耗时、调用次数、SQL 摘要。点进去看完整 SQL 和优化建议。

数据血缘图:可视化展示表和 API 的关联关系。点击某个表,高亮所有使用它的 API。

仪表盘用 Grafana 搭建,数据源接 Prometheus(指标)和 Elasticsearch(日志和追踪)。前端嵌入 iframe,和管理后台统一入口。


上一篇: 数据开放平台(四):数据转换与加工
下一篇: 数据开放平台(六):落地实践与运维