数据开放平台架构设计(五):监控与治理
上一篇: 数据开放平台(四):数据转换与加工
数据开放平台上线之后,更怕两件事:一是出了问题不知道,二是知道了问题找不到原因。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);
}
}五、治理仪表盘
把上面的能力整合到一个仪表盘里,一屏掌握全局:
监控与治理看板可以按数据来源拆开,页面上再组合成同一张运营视图:
概览区:今日调用总量、成功率、平均耗时、活跃 API 数、活跃调用方数。
实时监控:QPS 曲线图、耗时分布图(P50/P90/P99)、错误率趋势图。异常时自动标红。
告警列表:近段时间的告警事件,按时间倒序,支持一键跳转到追踪详情。
慢查询 TOP10:更慢的 10 个 API,展示平均耗时、调用次数、SQL 摘要。点进去看完整 SQL 和优化建议。
数据血缘图:可视化展示表和 API 的关联关系。点击某个表,高亮所有使用它的 API。
仪表盘用 Grafana 搭建,数据源接 Prometheus(指标)和 Elasticsearch(日志和追踪)。前端嵌入 iframe,和管理后台统一入口。
上一篇: 数据开放平台(四):数据转换与加工
下一篇: 数据开放平台(六):落地实践与运维