慢 sql 捕获应下沉至 dao 层,通过 spring aop 拦截 mapper 方法或 datasource 代理实现无侵入识别;提取真实 sql 与参数并脱敏,结合 mdc 或链路追踪获取 consumer 上下文,异步批量上报至审计中台。

在 Consumer 内部记录慢 SQL 并同步到审计中台,核心不是“在消费逻辑里手动测耗时”,而是把慢 SQL 的识别、提取和上报行为,从业务 Consumer 解耦出来,做到无侵入、可复用、带上下文、异步可靠。
✅ 关键思路:拦截执行点,而非包裹 Runnable 或 Consumer
Consumer(比如 KafkaListener、RocketMQListener)本身只负责接收消息并触发业务处理。真正执行 SQL 的是 DAO 层(MyBatis/JDBC),所以慢 SQL 的捕获必须下沉到数据访问层,不能靠在 @KafkaListener 方法里加 System.nanoTime()——那样既漏报(如批量更新、事务内多条 SQL)、又难提取真实 SQL 和参数,还污染业务代码。
? 一、推荐方式:用 Spring AOP 拦截 DAO 方法(最精准)
适用于 Spring Boot + MyBatis/MyBatis-Plus 场景:
-
切点定义
匹配所有查询/更新方法,例如:@Pointcut("execution(* com.example.dao..*Mapper.*(..)) || " + "execution(* com.example.repository..*Repository.*(..))") -
环绕通知中做三件事
- 记录开始时间(
long start = System.nanoTime()) -
proceed()执行原方法 -
finally中计算耗时,超阈值(如 300ms)则构造审计事件
- 记录开始时间(
-
SQL 提取(关键细节)
- MyBatis:通过
MappedStatement.getBoundSql(param)获取预编译 SQL + 参数映射 - MyBatis-Plus:
LambdaQueryWrapper可通过getTargetSql()或sqlBuilder提取(需配合MybatisPlusMetaObjectHandler或自定义Interceptor) - 参数摘要:用
ToStringBuilder.reflectionToString(params, ToStringStyle.SIMPLE_STYLE)截断+脱敏(如手机号、身份证号掩码)
- MyBatis:通过
-
审计内容示例(JSON 结构)
{ "sql": "SELECT id,name FROM user WHERE status = ? AND create_time > ?", "params": ["ACTIVE", "2026-06-01"], "costMs": 482, "method": "UserMapper.selectActiveUsers", "topic": "user_event", // 来源 Consumer 主题(从 MDC 或 Listener 方法名推导) "traceId": "a1b2c3d4", // 从 ThreadLocal 或 SkyWalking 上下文获取 "host": "app-server-01", "timestamp": "2026-06-12T02:27:15.123Z" }
? 二、通用方式:DataSource 代理 + Statement 封装(跨框架)
适合非 Spring 环境,或需要全局生效(含 JdbcTemplate、原生 JDBC):
- 自定义
DelegatingDataSource,重写getConnection() - 返回
ConnectionWrapper,再包装PreparedStatementWrapper - 在
executeQuery()/executeUpdate()前后打点,用nanoTime精确计时 - SQL 还原:对
PreparedStatement,需用反射读取this.sql字段,并用参数替换?(可复用p6spy的SqlWithParams逻辑)
⚠️ 注意:不要用
toString()直接取 SQL,PreparedStatement 的toString()多数返回类名,不是真实语句。
? 三、同步到审计中台:解耦 + 容错 + 聚合
-
不走 HTTP 同步调用,改用内存队列(如
BlockingQueue<slowsqlevent></slowsqlevent>)+ 单独守护线程消费 - 消费线程负责:
- 批量打包(如每 10 条或 1s 一次)
- 序列化为 JSON,发往审计中台 HTTP 接口(带重试、超时、限流)
- 失败时落盘(本地文件或 RocksDB),避免丢失
-
防雪崩设计:
- 同一 SQL 模板(参数占位符化后)5 分钟内最多上报 3 次
- 支持按
md5(sqlTemplate)聚合去重,避免高频重复告警
? 四、Consumer 上下文怎么关联?
- 在
@KafkaListener方法开头,把 topic、partition、offset、key 等写入MDC(Mapped Diagnostic Context):MDC.put("kafka_topic", topic); MDC.put("kafka_offset", String.valueOf(record.offset())); - AOP 切面中从
MDC.get("kafka_topic")取值,填入审计事件 - 若用 SkyWalking 或 OpenTelemetry,直接取
Tracer.currentSpan().context().traceId()
不复杂但容易忽略
Java免费学习笔记:立即使用
解锁 Java 大师之旅:从入门到精通的终极指南











