APM说DB慢,DBA说没慢SQL?我用这套全链路诊断引擎,终结了TDSQL-PG的甩锅大会!
2026/7/23 17:48:27 网站建设 项目流程

“Java + TDSQL-PG 全链路性能诊断引擎”。
今天,我把这套方案的底裤都扒干净:
✅ TDSQL-PG分布式架构下pg_stat_statements的隐藏坑点
✅ Java APM耗时与DB耗时“对不上”的5大元凶
✅ 可直接复用的MyBatis TraceID注入拦截器(保姆级代码)
✅ Python全自动甩锅终结脚本(拉取CN/DN数据做关联分析)
✅ 3个让我熬秃头的真实翻车案例

💀 第一章:TDSQL-PG版的pg_stat_statements,水有多深?

老铁们,如果你用的是单机PostgreSQL,pg_stat_statements(简称PSS)闭着眼睛用就行。

但TDSQL-PG是分布式架构啊!
它由CN(Coordinator Node,协调节点) 和DN(Data Node,数据节点) 组成。

这俩节点的PSS,看到的“世界”是完全不一样的。

1.1 CN与DN的“罗生门”
维度 CN节点的PSS DN节点的PSS
记录的SQL 用户原始SQL(逻辑SQL) 被CN改写后的物理分片SQL

包含网络耗时 ✅ 包含DN间数据shuffle耗时 ❌ 只记录本地执行耗时

包含解析耗时 ✅ 包含语法解析+分布式计划生成 ❌ 只记录执行耗时

数据倾斜可见性 ❌ 看不到,只能看到总耗时 ✅ 能看到哪个分片跑得最慢

排查价值 找“整体慢”的SQL 找“数据倾斜/分片不均”的SQL

⚠️ 重点:DBA如果只查DN的PSS,就会完美漏掉那些“DN跑得快,但CN网络传输慢”的幽灵SQL!开头那个甩锅大会,就是这么来的。

1.2 TDSQL-PG PSS的正确配置姿势

在TDSQL-PG里,PSS必须在CN和DN上都开启,而且参数有讲究。

– ============================================================
– 📌 TDSQL-PG 协调节点(CN) postgresql.conf 核心配置
– 💡 设计思想:CN是流量入口,必须记录全量SQL和完整耗时
– ============================================================

– 1️⃣ 必须加到共享预加载库(改完要重启!)
shared_preload_libraries = ‘pg_stat_statements’

– 2️⃣ 追踪级别:必须设为 ‘all’
– ⚠️ 坑点:如果设为 ‘top’,存储过程内部的SQL不会被记录!
pg_stat_statements.track = ‘all’

– 3️⃣ 最大记录SQL数:建议调大
– 💡 默认5000太小,大促期间高并发,热点SQL很快就会被挤出内存
pg_stat_statements.max = 20000

– 4️⃣ 是否记录规划时间(TDSQL-PG分布式计划生成很耗时,必须开!)
pg_stat_statements.track_planning = on

– 5️⃣ 是否落盘(防止重启丢失历史数据)
pg_stat_statements.save = on

– ============================================================
– 📌 启用插件(在CN和DN的业务库上都要执行)
– ============================================================
CREATE EXTENSION IF NOT EXISTS pg_stat_statements;

1.3 揪出“幽灵慢查询”的CN专属SQL

别再用网上那些单机版的PSS查询模板了!
在TDSQL-PG的CN节点,你需要关注分布式计划耗时和网络IO。

– ============================================================
– 📌 TDSQL-PG CN节点:排查“网络传输/重分布”导致的慢SQL
– 💡 核心逻辑:总耗时 - (规划耗时 + DN执行耗时) ≈ 网络Shuffle耗时
– ============================================================

SELECT
– 🎯 基础信息
queryid,
LEFT(query, 100) AS sql_snippet, – 截取前100个字符,防止刷屏

-- ⏱️ 耗时拆解(毫秒) calls AS exec_count, -- 执行次数 round(total_plan_time::numeric, 2) AS total_plan_ms, -- 总规划耗时 round(total_exec_time::numeric, 2) AS total_exec_ms, -- 总执行耗时 round(mean_time::numeric, 2) AS avg_total_ms, -- 平均总耗时 -- 🚨 核心指标:单次阅读块数(判断是否发生了大面积Broadcast) shared_blks_read AS phy_read_blks, -- 💡 辅助判断:如果临时块(temp_blks)极高,说明CN在做大规模Hash Join的内存溢出 temp_blks_read + temp_blks_written AS temp_io_blks, -- 🧮 缓存命中率(低于90%就要警惕了) round( (shared_blks_hit::numeric / NULLIF(shared_blks_hit + shared_blks_read, 0)) * 100, 2 ) AS hit_ratio_pct

FROM pg_stat_statements
WHERE
– 🚫 过滤掉系统自身的统计查询和PSS重置操作
query NOT ILIKE ‘%pg_stat_statements%’
AND query NOT ILIKE ‘%pg_catalog%’
– 🎯 只找平均耗时 > 100ms 的SQL
AND mean_time > 100
ORDER BY
– 💡 排序策略:按“总耗时”降序,找出真正的“时间黑洞”
total_exec_time DESC
LIMIT 20;

💡 金句:在分布式数据库里,慢的往往不是计算,而是“搬家”(数据Shuffle)。

🔥 第二章:Java端的“盲区”——为什么APM和DB对不上?

APM(如SkyWalking/Pinpoint)显示DB调用耗时500ms。
DBA查PSS,平均耗时只有20ms。
中间那480ms去哪了?被狗吃了吗?

老铁,这480ms,通常藏在Java到DB的“黑盒”里。

2.1 耗时差异的5大元凶

graph LR
J[☕ Java应用] -->|1. 获取连接| P[🏊 HikariCP连接池]
P -->|2. 网络传输| N[🌐 网络层]
N -->|3. 解析+规划| C[🧠 TDSQL-PG CN]
C -->|4. 执行+Shuffle| D[💾 TDSQL-PG DN]
D -->|5. 结果集返回| J

style P fill:#ff6b6b,color:#fff style N fill:#ffd43b,color:#333 style C fill:#845ef7,color:#fff

元凶 现象 怎么查

  1. 连接池等待 APM耗时高,DB没记录 查HikariCP的pending指标

  2. 结果集过大 DB执行快,但网络传了半分钟 查PSS的rows字段和Java端内存

  3. 事务未提交 SQL早跑完了,但事务挂起 查DB的pg_stat_activity状态

  4. JDBC隐式转换 索引失效,全表扫描 查执行计划(EXPLAIN)

  5. 批量插入未开启 1000条Insert变成1000次网络IO 查JDBC URL是否加了reWriteBatchedInserts

🚫 血泪教训:第5点我踩过巨坑!MyBatis-Plus的saveBatch默认是假批量,如果不加JDBC参数,TDSQL-PG会被打出屎。一定要在URL加上reWriteBatchedInserts=true!

💻 第三章:全链路诊断引擎(复制就能用,终结甩锅)

既然两边都有盲区,那我们就把Java和DB打通。
核心思路:在Java端给SQL注入TraceID,在DB端通过TraceID反查PSS。

3.1 Java端:MyBatis拦截器注入TraceID

我们要在SQL执行前,把SkyWalking的TraceID以SQL注释的形式拼接到SQL末尾。
这样,DBA在PSS里就能直接看到这条SQL对应的链路ID!

// ============================================================
// 📌 文件:TraceIdSqlInterceptor.java
// 📌 用途:MyBatis拦截器,自动在SQL尾部注入TraceID
// 📌 设计思想:
// 1. 使用 /* trace_id=xxx */ 注释,不影响SQL语法和执行计划
// 2. 兼容TDSQL-PG的CN节点解析(CN会保留注释记录到PSS)
// 3. 性能开销极低(仅字符串拼接)
// ============================================================

import org.apache.ibatis.executor.statement.StatementHandler;
import org.apache.ibatis.mapping.BoundSql;
import org.apache.ibatis.plugin.*;
import org.apache.skywalking.apm.toolkit.trace.TraceContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;

import java.lang.reflect.Field;
import java.sql.Connection;
import java.util.Properties;

@Component
@Intercepts({
@Signature(
type = StatementHandler.class,
method = “prepare”,
args = {Connection.class, Integer.class}
)
})
public class TraceIdSqlInterceptor implements Interceptor {

private static final Logger log = LoggerFactory.getLogger(TraceIdSqlInterceptor.class); // 📌 反射获取BoundSql的Field(缓存起来,避免每次反射影响性能) private static final Field BOUND_SQL_FIELD; static { try { // 💡 技巧:MyBatis的RoutingStatementHandler内部包装了真实的Handler // 需要拿到delegate里的boundSql BOUND_SQL_FIELD = BoundSql.class.getDeclaredField("sql"); BOUND_SQL_FIELD.setAccessible(true); } catch (NoSuchFieldException e) { throw new RuntimeException("初始化TraceId拦截器失败", e); } } @Override public Object intercept(Invocation invocation) throws Throwable { try { // 1️⃣ 获取StatementHandler StatementHandler handler = (StatementHandler) invocation.getTarget(); // 💡 边界处理:如果是RoutingStatementHandler,需要解包 // (这里为了代码简洁,假设直接拿到了BaseStatementHandler) BoundSql boundSql = handler.getBoundSql(); String originalSql = boundSql.getSql(); // 2️⃣ 获取SkyWalking的TraceID // ⚠️ 易错点:如果没引入SkyWalking依赖,这里会报错 // 建议用try-catch兜底,或者换成MDC.get("traceId") String traceId = TraceContext.traceId(); if (traceId == null || traceId.isEmpty() || "N/A".equals(traceId)) { traceId = "unknown"; } // 3️⃣ 清理SQL:去掉多余的换行和空格,防止注释被截断 String cleanSql = originalSql.replaceAll("[\s]+", " ").trim(); // 4️⃣ 拼接TraceID注释 // 🎯 格式:SELECT ... /* trace_id=abc123, app=order-service */ // 💡 为什么放尾部?因为有些ORM框架会截断SQL头部 String enhancedSql = String.format( "%s /* trace_id=%s, app=%s */", cleanSql, traceId, "order-service" // 📌 建议从环境变量读取应用名 ); // 5️⃣ 反射替换原始SQL BOUND_SQL_FIELD.set(boundSql, enhancedSql); } catch (Exception e) { // 🛡️ 兜底:拦截器绝对不能影响主流程!出错就放行原始SQL log.warn("⚠️ TraceID注入失败,降级使用原始SQL: {}", e.getMessage()); } // 6️⃣ 继续执行MyBatis后续流程 return invocation.proceed(); } @Override public Object plugin(Object target) { return Plugin.wrap(target, this); } @Override public void setProperties(Properties properties) { // 暂无额外配置 }

}

💡 技巧:加上这个拦截器后,DBA在TDSQL-PG的CN节点查PSS,query字段就会带上/* trace_id=xxx */。出了慢SQL,DBA直接把TraceID甩给Java开发:“自己去看APM链路!”甩锅闭环,完美!

3.2 DB端:Python全自动“甩锅终结”脚本

DBA不可能天天盯着PSS看。
我写了个Python脚本,定时拉取TDSQL-PG的PSS数据,结合Java APM的阈值,自动生成“诊断报告”。

============================================================
📌 文件:tdsql_pg_diagnostic.py
📌 用途:TDSQL-PG + Java 全链路性能自动诊断脚本
📌 依赖:pip install psycopg2-binary requests

import re
import psycopg2
import logging
from dataclasses import dataclass
from typing import Optional

logging.basicConfig(level=logging.INFO, format=“%(asctime)s [%(levelname)s] %(message)s”)

============================================================
📌 配置区域

TDSQL_CN_DSN = “host=10.0.0.10 port=15432 dbname=order_db user=readonly password=xxx”
TDSQL_DN_DSN = “host=10.0.0.11 port=15433 dbname=order_db user=readonly password=xxx”

🎯 诊断阈值(根据业务调整)
SLOW_SQL_THRESHOLD_MS = 200 # 慢SQL阈值(毫秒)
HIGH_TEMP_IO_THRESHOLD = 10000 # 高临时IO阈值(块)
LOW_HIT_RATIO_THRESHOLD = 80.0 # 低缓存命中率阈值(%)

@dataclass
class SlowQueryReport:
“”“慢查询诊断报告”“”
queryid: str
sql_snippet: str
trace_id: Optional[str]
avg_time_ms: float
calls: int
temp_io_blks: int
hit_ratio: float
diagnosis: str # 🩺 诊断结论

def extract_trace_id(sql: str) -> Optional[str]:
“”"
从SQL注释中提取TraceID
💡 正则匹配:/* trace_id=xxx */
“”"
match = re.search(r’/strace_id=([a-zA-Z0-9_.-]+)', sql)
return match.group(1) if match else None

def analyze_cn_slow_queries(conn) -> list[SlowQueryReport]:
“”"
分析CN节点的慢查询(侧重网络Shuffle和规划耗时)
“”"
query = “”"
SELECT
queryid::text,
query,
calls,
round(mean_time::numeric, 2) AS avg_time_ms,
temp_blks_read + temp_blks_written AS temp_io,
round((shared_blks_hit::numeric / NULLIF(shared_blks_hit + shared_blks_read, 0)) * 100, 2) AS hit_ratio
FROM pg_stat_statements
WHERE mean_time > %s
AND query NOT ILIKE ‘%pg_stat_statements%’
ORDER BY total_exec_time DESC
LIMIT 50;
“”"

reports = [] with conn.cursor() as cur: cur.execute(query, (SLOW_SQL_THRESHOLD_MS,)) for row in cur.fetchall(): queryid, sql, calls, avg_time, temp_io, hit_ratio = row trace_id = extract_trace_id(sql) # 🩺 自动诊断逻辑(核心干货!) diagnosis = [] # 1. 检查临时IO(判断是否发生大规模数据重分布/Hash溢出) if temp_io > HIGH_TEMP_IO_THRESHOLD: diagnosis.append( "🚨 [分布式Shuffle警告] 临时块IO极高({})," "大概率是分布式Join触发了Broadcast/Redistribute," "请检查Join字段是否为分布键(Shardkey)!".format(temp_io) ) # 2. 检查缓存命中率 if hit_ratio is not None and hit_ratio < LOW_HIT_RATIO_THRESHOLD: diagnosis.append( "⚠️ [缓存击穿警告] 命中率仅{}%," "可能存在全表扫描或大量冷数据查询,请检查索引。".format(hit_ratio) ) # 3. 检查执行频次(高频+慢 = 灾难) if calls > 10000 and avg_time > 500: diagnosis.append( "💀 [雪崩警告] 高频({})+高耗时({}ms)," "正在拖垮CN节点CPU,建议立即限流或加缓存!".format(calls, avg_time) ) if not diagnosis: diagnosis.append("🔍 常规慢查询,建议结合TraceID查看APM链路耗时。") reports.append(SlowQueryReport( queryid=queryid, sql_snippet=sql[:150].replace('n', ' '), trace_id=trace_id, avg_time_ms=avg_time, calls=calls, temp_io_blks=temp_io, hit_ratio=hit_ratio or 0.0, diagnosis=" | ".join(diagnosis) )) return reports

def generate_markdown_report(reports: list[SlowQueryReport]):
“”“生成Markdown格式的甩锅终结报告”“”
md = [“# 🩺 TDSQL-PG 慢查询全链路诊断报告n”]
md.append(f"> 扫描阈值: 平均耗时 > {SLOW_SQL_THRESHOLD_MS}msnn")

if not reports: md.append("✅ 恭喜!未发现超过阈值的慢查询,DBA今天可以准时下班。n") else: md.append(f"🚨 发现 {len(reports)} 个高危慢查询,请相关开发/DBA立即认领!nn") md.append("| 严重度 | TraceID | 平均耗时 | 频次 | 诊断结论 | SQL片段 |n") md.append("|---|---|---|---|---|---|n") for r in reports: severity = "🔴" if r.avg_time_ms > 1000 else "🟠" trace_link = f"{r.trace_id}" if r.trace_id else "❌ 未注入" md.append( f"| {severity} | {trace_link} | {r.avg_time_ms}ms | {r.calls} " f"| {r.diagnosis} | {r.sql_snippet[:80]}... |n" ) report_content = "n".join(md) with open("tdsql_slow_query_report.md", "w", encoding="utf-8") as f: f.write(report_content) print("📄 报告已生成: tdsql_slow_query_report.md")

if name == “main”:
try:
logging.info(“🔌 正在连接 TDSQL-PG CN节点…”)
cn_conn = psycopg2.connect(TDSQL_CN_DSN)
reports = analyze_cn_slow_queries(cn_conn)
generate_markdown_report(reports)
except Exception as e:
logging.error(f"诊断脚本执行失败: {e}")
finally:
if ‘cn_conn’ in locals():
cn_conn.close()

💡 金句:没有TraceID的慢SQL,就像没有指纹的凶器,谁都不认账。把TraceID刻进SQL里,让证据自己说话。

💀 第四章:3个真实翻车案例(血泪教训)

5.1 翻车1:JDBC隐式转换,把TDSQL-PG的索引干碎了

场景:订单表orders,order_no字段类型是VARCHAR(64),并且是分布键(Shardkey)+ 主键。

Java代码(MyBatis):

SELECT * FROM orders WHERE order_no = #{orderNo}

翻车点:
TDSQL-PG的CN节点在解析时,发现参数类型和字段类型不完全匹配,偷偷加了一个隐式转换:
– CN改写后的实际执行SQL
SELECT * FROM orders WHERE order_no = ($1)::text

后果:
因为带了类型转换函数,索引直接失效!
更恐怖的是,order_no是分布键。分布键失效意味着CN无法确定数据在哪个DN,只能向所有DN发起Broadcast(广播)全表扫描!
一个本来0.1ms的点查,变成了500ms的全集群广播。

修复:
在MyBatis中显式指定jdbcType:

SELECT * FROM orders WHERE order_no = #{orderNo, jdbcType=VARCHAR}

🚫 教训:在PostgreSQL/TDSQL-PG里,永远不要相信JDBC驱动的自动类型推断。VARCHAR、BIGINT、TIMESTAMP,统统给我显式写上jdbcType!

5.2 翻车2:分布式Join的“数据倾斜”惨案

场景:
– 用户表(按user_id分片) JOIN 订单表(按order_id分片)
SELECT u.name, COUNT(o.id)
FROM users u
JOIN orders o ON u.id = o.user_id
WHERE u.create_time > ‘2023-01-01’
GROUP BY u.name;

现象:
APM显示这个接口耗时3秒。DBA看CN的PSS,耗时3秒。看DN的PSS,发现DN-1跑了2.9秒,DN-2到DN-8只跑了0.01秒。

原因:
两个表的分布键不同(一个是user_id,一个是order_id)。
TDSQL-PG的CN为了做Join,必须把数据Redistribute(重分布) 到同一个节点。
偏偏有个“测试账号”(user_id=0),下面挂了200万条订单!
这200万条数据全部Shuffle到了DN-1,DN-1直接被CPU打满,其他DN在看戏。

修复:
业务层:清理脏数据,禁止测试账号参与真实统计。
SQL层:如果必须查,加上过滤条件剔除极端倾斜的Key。
架构层:对于高频Join的表,建表时尽量使用相同的分布键(Colocation Group),让Join在本地DN完成,避免网络Shuffle。

💡 魔性比喻:数据倾斜就像吃火锅,有人碗里全是肉(DN-1),有人碗里只有汤(DN-2)。建表时选好分布键,就是给每个人分个公平的漏勺。

5.3 翻车3:长事务+连接池,拖垮整个集群

场景:
Java端有个导出Excel的功能。
@Transactional // ❌ 大事务!
public void exportAllOrders() {
List list = orderMapper.selectAll(); // 查出50万条数据
// 然后开始用POI写Excel,写了整整2分钟…
excelWriter.write(list);
}

翻车点:
@Transactional开启了事务,查完50万条数据后,事务没有提交!
Java线程在慢吞吞地写Excel(耗时2分钟)。
这2分钟内,TDSQL-PG的这个连接一直处于idle in transaction(事务空闲)状态。
TDSQL-PG的MVCC机制导致这2分钟内,其他事务产生的旧版本数据无法被Vacuum回收!
并发一高,连接池瞬间耗尽,数据库表膨胀,全盘崩溃。

修复:
// ✅ 修复:剥离大事务,流式查询 + 分批处理
public void exportAllOrders() {
// 1. 去掉@Transactional
// 2. 使用MyBatis的ResultHandler流式读取,不占满JVM内存
orderMapper.streamSelectAll(new ResultHandler() {
@Override
public void handleResult(ResultContext<?> context) {
Order order = (Order) context.getResultObject();
excelWriter.writeRow(order); // 边查边写
}
});
}

🚫 教训:永远不要在事务里做耗时的非DB操作(如网络请求、写文件、发邮件)! 事务的生命周期必须像闪电一样快。

✅ 第五章:TDSQL-PG + Java 性能调优CheckList(建议打印贴工位)

📋 Java / MyBatis 层
[ ] 检查JDBC URL是否开启reWriteBatchedInserts=true(批量插入必开)
[ ] 检查MyBatis XML中的#{}是否显式指定了jdbcType
[ ] 检查HikariCP的maximumPoolSize是否合理(别盲目设大,推荐公式:CPU核数 * 2 + 磁盘数)
[ ] 检查是否存在@Transactional包裹了RPC调用或文件IO
[ ] 确认是否已部署TraceID注入拦截器

📋 TDSQL-PG 数据库层
[ ] 确认pg_stat_statements在CN和DN均已开启,且track=all
[ ] 检查高频Join的表,是否属于同一个Colocation Group(同分布键)
[ ] 检查是否存在大量idle in transaction的长连接(SELECT * FROM pg_stat_activity WHERE state = ‘idle in transaction’)
[ ] 检查DN节点是否存在严重的数据倾斜(通过pg_class和分片统计查看)
[ ] 确认work_mem配置是否足够(防止Hash Join溢出到磁盘临时文件)

🎯 结论:调优不是调参数,是懂全链路

回到开头那个甩锅大会。

为什么Java和DBA总是互相指责?
因为他们都只看到了大象的一条腿。

Java看到了APM的耗时,DBA看到了DN的耗时。
中间的网络Shuffle、连接池等待、隐式转换,成了没人管的“三不管地带”。

金句:在分布式时代,没有全链路视角的性能调优,都是盲人摸象。

把TraceID注入SQL,把CN和DN的PSS数据拉通,把Java的JDBC参数规范化。
这不是在写代码,这是在建立团队的“技术信任”。

当证据链闭环了,甩锅大会自然就变成了复盘大会。

这才是老码农该有的工程格局。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询