1.pom

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>com.sunwoda</groupId>
    <artifactId>pg-test</artifactId>
    <version>2.0-SNAPSHOT</version>

    <properties>
        <java.version>1.8</java.version>
        <maven.compiler.source>${java.version}</maven.compiler.source>
        <maven.compiler.target>${java.version}</maven.compiler.target>
        <fastjson.vsersion>2.0.52</fastjson.vsersion>
        <postgresql.version>42.2.12</postgresql.version>
        <mysql.version>8.0.33</mysql.version>
        <logback.version>1.2.11</logback.version>
        <lombok.version>1.18.20</lombok.version>
    </properties>

    <dependencies>
        <dependency>
            <groupId>org.postgresql</groupId>
            <artifactId>postgresql</artifactId>
            <version>${postgresql.version}</version>
        </dependency>
        <dependency>
            <groupId>com.mysql</groupId>
            <artifactId>mysql-connector-j</artifactId>
            <version>${mysql.version}</version>
        </dependency>
        <dependency>
            <groupId>ch.qos.logback</groupId>
            <artifactId>logback-classic</artifactId>
            <version>${logback.version}</version>
        </dependency>
        <dependency>
            <groupId>ch.qos.logback</groupId>
            <artifactId>logback-core</artifactId>
            <version>${logback.version}</version>
        </dependency>
        <dependency>
            <groupId>com.alibaba</groupId>
            <artifactId>fastjson</artifactId>
            <version>${fastjson.vsersion}</version>
        </dependency>
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <version>${lombok.version}</version>
            <scope>provided</scope>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>3.0.0</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies</descriptorRef>
                    </descriptorRefs>
                    <archive>
                        <manifest>
                            <mainClass>com.test.cdc.PgCdcUltimateMonitor</mainClass>
                        </manifest>
                    </archive>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

2.java代码主类

package com.test.cdc;

import com.alibaba.fastjson.JSON;
import lombok.extern.slf4j.Slf4j;

import java.sql.*;
import java.util.*;
import java.util.Date;
@Slf4j
public class PgCdcUltimateMonitor {

    // PostgreSQL 配置(多个数据库)
    private static final String PG_HOST = "ip";
    private static final int PG_PORT = 5432;
    private static final String PG_USER = "";
    private static final String PG_PASSWORD = "";
    private static final List<String> PG_DATABASES = Arrays.asList(
            "dbName1",
            "dbName2",
            "dbName3"
    );

    // Doris 配置
    private static final String DORIS_URL = "jdbc:mysql://ip:9030/ods";
    private static final String DORIS_USER = "";
    private static final String DORIS_PASSWORD = "";

    // 采集间隔(3 分钟)
    private static final long INTERVAL = 60 * 3000;

    static {
        try {
            Class.forName("com.mysql.cj.jdbc.Driver");
        } catch (ClassNotFoundException e) {
            log.error("找不到 MySQL 驱动:{}", e.getMessage());
        }
    }

    public static void main(String[] args) {
        new PgCdcUltimateMonitor().start();
    }

    public void start() {
        log.info("===== PostgreSQL CDC 终极监控程序启动 =====");
        new Timer().scheduleAtFixedRate(new MonitorTask(), 0, INTERVAL);
    }


    static class MonitorTask extends TimerTask {

        @Override
        public void run() {
            // 遍历所有数据库
            for (String dbName : PG_DATABASES) {
                String pgUrl = String.format("jdbc:postgresql://%s:%d/%s", PG_HOST, PG_PORT, dbName);

                try (Connection pgConn = DriverManager.getConnection(pgUrl, PG_USER, PG_PASSWORD)) {
                    log.info("[{}] 开始采集数据库: {}", new Date(), dbName);

                    // --------------------------
                    // 1. 采集 PostgreSQL 指标
                    // --------------------------

                    long checkpointsTimed = 0;
                    long checkpointsReq = 0;
                    long walBuffersFull = 0;
                    double cpuUsage = 0.0;
                    double memoryUsage = 0.0;
                    int activeConn = 0;
                    int idleConn = 0;
                    int lockedConn = 0;
                    long xactCommit = 0;
                    long xactRollback = 0;
                    double cacheHitRatio = 0.0;
                    String longTransactionJson = "";

                    // --------------------------
                    // 新增:IO 核心指标
                    // --------------------------
                    long blksRead = 0;
                    long blksHit = 0;
                    long heapBlksRead = 0;
                    long heapBlksHit = 0;
                    double heapCacheHitRatio = 0.0;
                    long idxScan = 0;
                    long idxTupRead = 0;
                    long idxTupFetch = 0;
                    long walWrite = 0;
                    long walWriteTime = 0;
                    long walSync = 0;
                    long walSyncTime = 0;
                    long checkpointWriteTime = 0;
                    long checkpointSyncTime = 0;
                    // 复制槽监控
                    String slotState = "";
                    String writeLagTime = "";
                    String flushLagTime = "";
                    String replayLagTime = "";
                    String writeBacklogSize = "";
                    String flushBacklogSize = "";
                    String replayBacklogSize = "";
                    String slotTotalBacklogSize = "";
                    String unflushBacklogSize = "";
                    String restartLsn = "";
                    String confirmedFlushLsn = "";

                    // WAL & Checkpoint
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT checkpoints_timed, checkpoints_req FROM pg_stat_bgwriter"
                    ).executeQuery()) {
                        if (rs.next()) {
                            checkpointsTimed = rs.getLong("checkpoints_timed");
                            checkpointsReq = rs.getLong("checkpoints_req");
                        }
                    }

                    // WAL Buffers Full
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT wal_buffers_full FROM pg_stat_wal"
                    ).executeQuery()) {
                        if (rs.next()) {
                            walBuffersFull = rs.getLong("wal_buffers_full");
                        }
                    }

                    // --------------------------
                    // 新增:WAL IO 详细指标
                    // --------------------------
                    try (ResultSet rs = pgConn.prepareStatement(
                            "select wal_write, wal_write_time, wal_sync, wal_sync_time FROM pg_stat_wal"
                    ).executeQuery()) {
                        if (rs.next()) {
                            walWrite = rs.getLong("wal_write");
                            walWriteTime = rs.getLong("wal_write_time");
                            walSync = rs.getLong("wal_sync");
                            walSyncTime = rs.getLong("wal_sync_time");
                        }
                    }

                    // --------------------------
                    // 新增:检查点 IO 指标
                    // --------------------------
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT checkpoint_write_time, checkpoint_sync_time FROM pg_stat_bgwriter"
                    ).executeQuery()) {
                        if (rs.next()) {
                            checkpointWriteTime = rs.getLong("checkpoint_write_time");
                            checkpointSyncTime = rs.getLong("checkpoint_sync_time");
                        }
                    }

                    // CPU / Memory(最终稳定版,PostgreSQL 15 已验证)
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT " +
                                    "ROUND( " +
                                    " (COUNT(*) FILTER (WHERE state = 'active'))::NUMERIC / COUNT(*) * 100, " +
                                    " 2 " +
                                    ") AS cpu_usage, " +
                                    "ROUND( " +
                                    " (SELECT setting::NUMERIC FROM pg_settings WHERE name = 'shared_buffers') / 1024 / 1024 / 1024, " +
                                    " 2 " +
                                    ") AS memory_usage " +
                                    "FROM pg_stat_activity " +
                                    "LIMIT 1"
                    ).executeQuery()) {
                        if (rs.next()) {
                            cpuUsage = rs.getDouble("cpu_usage");
                            memoryUsage = rs.getDouble("memory_usage");
                        }
                    }

                    // Connections
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT " +
                                    "COUNT(*) FILTER (WHERE state = 'active') AS active, " +
                                    "COUNT(*) FILTER (WHERE state IN ('idle','idle in transaction')) AS idle, " +
                                    "COUNT(*) FILTER (WHERE wait_event_type = 'LOCK') AS locked " +
                                    "FROM pg_stat_activity"
                    ).executeQuery()) {
                        if (rs.next()) {
                            activeConn = rs.getInt("active");
                            idleConn = rs.getInt("idle");
                            lockedConn = rs.getInt("locked");
                        }
                    }

                    // Transaction & Cache
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT " +
                                    "xact_commit, " +
                                    "xact_rollback, " +
                                    "(blks_hit::FLOAT/(blks_hit + blks_read))*100 AS cache_hit_ratio, " +
                                    "blks_read, " +
                                    "blks_hit " +
                                    "FROM pg_stat_database WHERE datname = current_database()"
                    ).executeQuery()) {
                        if (rs.next()) {
                            xactCommit = rs.getLong("xact_commit");
                            xactRollback = rs.getLong("xact_rollback");
                            cacheHitRatio = rs.getDouble("cache_hit_ratio");
                            blksRead = rs.getLong("blks_read");
                            blksHit = rs.getLong("blks_hit");
                        }
                    }

                    // --------------------------
                    // 表 IO 统计
                    // --------------------------
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT " +
                                    "blks_read AS heap_blks_read, " +
                                    "blks_hit AS heap_blks_hit, " +
                                    "(blks_hit::FLOAT / (blks_hit + blks_read)) * 100 AS heap_cache_hit_ratio " +
                                    "FROM pg_stat_database " +
                                    "WHERE datname = current_database()"
                    ).executeQuery()) {
                        if (rs.next()) {
                            heapBlksRead = rs.getLong("heap_blks_read");
                            heapBlksHit = rs.getLong("heap_blks_hit");
                            heapCacheHitRatio = rs.getDouble("heap_cache_hit_ratio");
                        }
                    }

                    // --------------------------
                    // 索引 IO 统计(你已确认可用)
                    // -------------------------
                    try (ResultSet rs = pgConn.prepareStatement(
                            "SELECT " +
                                    "SUM(idx_scan) AS total_idx_scan, " +
                                    "SUM(idx_tup_read) AS total_idx_tup_read, " +
                                    "SUM(idx_tup_fetch) AS total_idx_tup_fetch " +
                                    "FROM pg_stat_user_indexes"
                    ).executeQuery()) {
                        if (rs.next()) {
                            idxScan = rs.getLong("total_idx_scan");
                            idxTupRead = rs.getLong("total_idx_tup_read");
                            idxTupFetch = rs.getLong("total_idx_tup_fetch");
                        }
                    }

                    //采集PostgreSQL中耗时超3分钟的大事务
                    List<LongTransaction> longTransactions = new ArrayList<>();
                    try (PreparedStatement pstmt = pgConn.prepareStatement(
                            "SELECT " +
                                    "pid AS transaction_id, " +
                                    "EXTRACT(EPOCH FROM (NOW() - xact_start)) AS duration_seconds, " +
                                    "query AS sql, " +
                                    "state " +
                                    "FROM pg_stat_activity " +
                                    "WHERE " +
                                    "xact_start IS NOT NULL " + // 存在活跃事务
                                    "AND EXTRACT(EPOCH FROM (NOW() - xact_start)) > ? " + // 事务时长超阈值
                                    "AND state <> 'idle'")) {
                        pstmt.setLong(1, 180);
                        try (ResultSet rs = pstmt.executeQuery()) {
                            while (rs.next()) {
                                LongTransaction transaction = new LongTransaction();
                                transaction.setTransactionId(rs.getString("transaction_id"));
                                transaction.setDurationSeconds(rs.getLong("duration_seconds"));
                                transaction.setSql(rs.getString("sql"));
                                transaction.setState(rs.getString("state"));
                                longTransactions.add(transaction);
                                System.out.println(longTransactions);

                                longTransactionJson =  JSON.toJSONString(longTransactions);
                            }
                        }
                    }

                    // --------------------------
                    // Flink CDC复制槽延迟监控
                    // --------------------------
                    try (PreparedStatement pstmt = pgConn.prepareStatement(
                            "SELECT " +
                                    "psr.state, " +
                                    "COALESCE(to_char(psr.write_lag, 'HH24:MI:SS.ms'), '00:00:00.000') AS write_lag_time, " +
                                    "COALESCE(to_char(psr.flush_lag, 'HH24:MI:SS.ms'), '00:00:00.000') AS flush_lag_time, " +
                                    "COALESCE(to_char(psr.replay_lag, 'HH24:MI:SS.ms'), '00:00:00.000') AS replay_lag_time, " +
                                    "pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), psr.write_lsn)) AS write_backlog_size, " +
                                    "pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), psr.flush_lsn)) AS flush_backlog_size, " +
                                    "pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), psr.replay_lsn)) AS replay_backlog_size, " +
                                    "pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), prs.restart_lsn)) AS slot_total_backlog_size, " +
                                    "pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), COALESCE(prs.confirmed_flush_lsn, prs.restart_lsn))) AS unflush_backlog_size, " +
                                    "prs.restart_lsn, " +
                                    "prs.confirmed_flush_lsn " +
                                    "FROM pg_stat_replication psr " +
                                    "LEFT JOIN pg_replication_slots prs ON psr.pid = prs.active_pid " +
                                    "WHERE prs.slot_name IS NOT NULL " +
                                    "AND prs.slot_name LIKE ?" // 用?占位符,替代字符串拼接
                    )){
                        // 设置模糊查询参数,手动拼接%,避免SQL注入
                        pstmt.setString(1, "%" + dbName + "%");
                      try(ResultSet rs = pstmt.executeQuery()) {
                        while (rs.next()) {
                            // 状态&时间延迟
                             slotState = rs.getString("state");
                             writeLagTime = rs.getString("write_lag_time");
                             flushLagTime = rs.getString("flush_lag_time");
                             replayLagTime = rs.getString("replay_lag_time");
                            // 数据量积压(带单位:B/KB/MB/GB)
                             writeBacklogSize = rs.getString("write_backlog_size");
                             flushBacklogSize = rs.getString("flush_backlog_size");
                             replayBacklogSize = rs.getString("replay_backlog_size");
                             slotTotalBacklogSize = rs.getString("slot_total_backlog_size");
                             unflushBacklogSize = rs.getString("unflush_backlog_size");
                            // 原始LSN位点(字符串类型,用于位点定位/排查)
                             restartLsn = rs.getString("restart_lsn");
                             confirmedFlushLsn = rs.getString("confirmed_flush_lsn");
                        }
                    } catch (Exception e) {
                        // 按需添加异常处理(如日志打印)
                        log.error("复制槽监控:{}",e.getMessage());
                    }
                }
                    // --------------------------
                    // 2. CDC 健康度分析
                    // --------------------------
                    CdcHealthResult health = CdcHealthResult.analyze(
                            dbName,
                            walBuffersFull,
                            checkpointsReq,
                            lockedConn,
                            cpuUsage,
                            cacheHitRatio,
                            activeConn,
                            heapBlksRead,
                            heapCacheHitRatio,
                            idxScan,
                            walWrite,
                            walWriteTime,
                            checkpointWriteTime,
                            unflushBacklogSize
                    );

                    // --------------------------
                    // 3. 写入 Doris
                    // --------------------------

                    try (Connection dorisConn = DriverManager.getConnection(DORIS_URL, DORIS_USER, DORIS_PASSWORD)) {

                        String sql = "INSERT INTO t_pg_log (" +
                                "id, " +
                                "database_name, " +
                                "checkpoints_timed, checkpoints_req, wal_buffers_full, " +
                                "wal_writes, wal_write_time, wal_syncs, wal_sync_time, " +
                                "checkpoint_write_time, checkpoint_sync_time, " +
                                "blks_read, blks_hit, " +
                                "heap_blks_read, heap_blks_hit, heap_cache_hit_ratio," +
                                "idx_scan, idx_tup_read, idx_tup_fetch," +
                                "cpu_usage, memory_usage, active_connections, idle_connections, locked_connections, " +
                                "xact_commit, xact_rollback, cache_hit_ratio, " +
                                "slot_state, write_lag_time, flush_lag_time, replay_lag_time, write_backlog_size, flush_backlog_size,"+
                                "replay_backlog_size, slot_total_backlog_size, unflush_backlog_size, restart_lsn, confirmed_flush_lsn,"+
                                "cdc_health_level, cdc_health_reason, cdc_health_suggestion, longTransactionJson, create_time" +
                                ") VALUES (?,?,?,?,?,?,?,?,?,?,?,?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)";

                        try (PreparedStatement pstmt = dorisConn.prepareStatement(sql)) {
                            int idx = 1;
                            pstmt.setString(idx++, UUID.randomUUID().toString().replace("-", ""));
                            pstmt.setString(idx++, dbName);
                            pstmt.setLong(idx++, checkpointsTimed);
                            pstmt.setLong(idx++, checkpointsReq);
                            pstmt.setLong(idx++, walBuffersFull);

                            // 新增 IO 字段
                            pstmt.setLong(idx++, walWrite);
                            pstmt.setLong(idx++, walWriteTime);
                            pstmt.setLong(idx++, walSync);
                            pstmt.setLong(idx++, walSyncTime);
                            pstmt.setLong(idx++, checkpointWriteTime);
                            pstmt.setLong(idx++, checkpointSyncTime);
                            pstmt.setLong(idx++, blksRead);
                            pstmt.setLong(idx++, blksHit);
                            pstmt.setLong(idx++, heapBlksRead);
                            pstmt.setLong(idx++, heapBlksHit);
                            pstmt.setDouble(idx++, heapCacheHitRatio);
                            pstmt.setLong(idx++, idxScan);
                            pstmt.setLong(idx++, idxTupRead);
                            pstmt.setLong(idx++, idxTupFetch);

                            pstmt.setDouble(idx++, cpuUsage);
                            pstmt.setDouble(idx++, memoryUsage);
                            pstmt.setInt(idx++, activeConn);
                            pstmt.setInt(idx++, idleConn);
                            pstmt.setInt(idx++, lockedConn);
                            pstmt.setLong(idx++, xactCommit);
                            pstmt.setLong(idx++, xactRollback);
                            pstmt.setDouble(idx++, cacheHitRatio);
                            // 新增复制槽监控
                            pstmt.setString(idx++,slotState);
                            pstmt.setString(idx++,writeLagTime);
                            pstmt.setString(idx++,flushLagTime);
                            pstmt.setString(idx++,replayLagTime);
                            pstmt.setString(idx++,writeBacklogSize);
                            pstmt.setString(idx++,flushBacklogSize);
                            pstmt.setString(idx++,replayBacklogSize);
                            pstmt.setString(idx++,slotTotalBacklogSize);
                            pstmt.setString(idx++,unflushBacklogSize);
                            pstmt.setString(idx++,restartLsn);
                            pstmt.setString(idx++,confirmedFlushLsn);
                            pstmt.setString(idx++, health.level);
                            pstmt.setString(idx++, health.reason);
                            pstmt.setString(idx++, health.suggestion);
                            pstmt.setString(idx++, longTransactionJson);
                            pstmt.setTimestamp(idx++, new Timestamp(System.currentTimeMillis()));

                            pstmt.executeUpdate();
                            log.info("[{}] 数据库 {} 写入 Doris 成功:CDC 健康度={}",
                                    new Date(), dbName, health.level);
                        }
                    }

                } catch (SQLException e) {
                    log.error("[{}] 数据库 {} 采集失败: {}", new Date(), dbName, e.getMessage());
                }
            }
        }
    }
}

3.输出优化建议

package com.test.cdc;

import org.apache.commons.lang3.StringUtils;
public class CdcHealthResult {
    public String level;
    public String reason;
    public String suggestion;

    public CdcHealthResult(String level, String reason, String suggestion) {
        this.level = level;
        this.reason = reason;
        this.suggestion = suggestion;
    }

    public static CdcHealthResult analyze(
            String dbName,
            long walBuffersFull,
            long checkpointsReq,
            int lockedConn,
            double cpuUsage,
            double cacheHitRatio,
            int activeConn,

            // --------------------------
            // 新增:IO 指标
            // --------------------------
            long heapBlksRead,
            double heapCacheHitRatio,
            long idxScan,
            long walWrite,
            long walSyncTime,
            long checkpointWriteTime
    ) {

        // --------------------------
        // 1. 原有增量计算
        // --------------------------
        long walBuffersFullIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "wal_buffers_full", walBuffersFull);

        long checkpointsReqIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "checkpoints_req", checkpointsReq);

        // --------------------------
        // 2. 新增:IO 增量计算
        // --------------------------

        long heapBlksReadIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "heap_blks_read", heapBlksRead);

        long idxScanIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "idx_scan", idxScan);

        long walWriteIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "wal_write", walWrite);

        long walSyncTimeIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "wal_sync_time", walSyncTime);

        long checkpointWriteTimeIncrement = MetricCache.getInstance()
                .getIncrement(dbName, "checkpoint_write_time", checkpointWriteTime);

        if (walBuffersFullIncrement > 100) {
            return new CdcHealthResult(
                    "ALERT",
                    "WAL 缓冲区写满,CDC 会被阻塞,延迟必然上升",
                    "增大 wal_buffers、提升磁盘 IO、减少大事务"
            );
        }

        if (lockedConn > 5) {
            return new CdcHealthResult(
                    "ALERT",
                    "锁等待过多(" + lockedConn + "),事务提交变慢,CDC 延迟会上升",
                    "检查慢 SQL、长事务、锁竞争"
            );
        }

        if (checkpointsReqIncrement > 10 || cpuUsage > 80) {
            return new CdcHealthResult(
                    "WARNING",
                    "checkpoint 频繁或 CPU 过高,可能影响 CDC 性能",
                    "增大 shared_buffers、优化 SQL、降低写入压力"
            );
        }

        if (cacheHitRatio < 95) {
            return new CdcHealthResult(
                    "WARNING",
                    "缓存命中率低(" + String.format("%.2f", cacheHitRatio) + "%),磁盘 IO 升高",
                    "增大 shared_buffers、优化索引"
            );
        }

        if (activeConn > 200) {
            return new CdcHealthResult(
                    "WARNING",
                    "活跃连接过多(" + activeConn + "),PostgreSQL 压力增大",
                    "优化连接池、减少长连接"
            );
        }

        // --------------------------
        // 3. 新增:IO 健康规则
        // --------------------------

        // 表 IO 过高(磁盘读太多,参数自己调整)
        if (heapBlksReadIncrement > 50000) {
            return new CdcHealthResult(
                    "WARNING",
                    "表数据块磁盘读取量突增(" + heapBlksReadIncrement + "次),可能引发 CDC 延迟",
                    "优化 SQL、增加索引、提高 shared_buffers"
            );
        }

        // 表缓存命中率过低
        if (heapCacheHitRatio < 95) {
            return new CdcHealthResult(
                    "WARNING",
                    "表缓存命中率低(" + String.format("%.2f", heapCacheHitRatio) + "%)",
                    "增加 shared_buffers、优化热点数据访问"
            );
        }

        // 索引扫描突增(参数自己调整)
        if (idxScanIncrement > 7000000) {
            return new CdcHealthResult(
                    "WARNING",
                    "索引扫描次数突增(" + idxScanIncrement + "次),可能影响 CDC 性能",
                    "优化索引、检查慢查询"
            );
        }

        // WAL 写入量突增
        if (walWriteIncrement > 1024 * 1024 * 100) { // 100MB
            return new CdcHealthResult(
                    "WARNING",
                    "WAL 写入量突增(" + walWriteIncrement / (1024 * 1024) + "MB),CDC 可能延迟",
                    "减少大事务、优化写入模式"
            );
        }

        // WAL 同步耗时过长
        if (walSyncTimeIncrement > 5000) { // 5 秒
            return new CdcHealthResult(
                    "ALERT",
                    "WAL 同步耗时过长(" + walSyncTimeIncrement + "ms),CDC 会被阻塞",
                    "提升磁盘 IO 性能(建议 SSD)、降低 wal_sync 频率"
            );
        }

        // Checkpoint 写入耗时过长
        if (checkpointWriteTimeIncrement > 10000) { // 10 秒
            return new CdcHealthResult(
                    "WARNING",
                    "Checkpoint 写入耗时过长(" + checkpointWriteTimeIncrement + "ms)",
                    "增大 shared_buffers、提升磁盘性能"
            );
        }

        // --------------------------
        // 4. 健康状态
        // --------------------------
        return new CdcHealthResult(
                "HEALTHY",
                "PostgreSQL 状态良好,CDC 运行稳定",
                "继续保持当前配置"
        );
    }
}

4.对部分指标进行增量计算

package com.test.cdc;

import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

// 历史数据缓存类,单例模式
public class MetricCache {
    private static final MetricCache INSTANCE = new MetricCache();
    // 存储各数据库的历史指标(key: databaseName,value: 指标Map)
    private final Map<String, Map<String, Long>> historyMetrics = new ConcurrentHashMap<>();

    private MetricCache() {}

    public static MetricCache getInstance() {
        return INSTANCE;
    }

    // 保存当前指标,返回增量
    public long getIncrement(String dbName, String metricName, long currentValue) {
        // 初始化数据库对应的指标Map
        historyMetrics.computeIfAbsent(dbName, k -> new HashMap<>());
        Map<String, Long> dbMetrics = historyMetrics.get(dbName);
        if (!dbMetrics.containsKey(metricName)) {
            dbMetrics.put(metricName, currentValue); // 缓存当前值作为下次历史值
            return -1; // 标识首次采集,不做增量计算,只初始值
        }
        // 计算增量(首次采集时增量为0)
        long lastValue = dbMetrics.getOrDefault(metricName, 0L);
        long increment = currentValue - lastValue;

        // 更新历史值
        dbMetrics.put(metricName, currentValue);
        return increment;
    }
}

5.创建个事务类,用来记录监控耗时长的sql

package com.test.cdc;

import lombok.Data;

@Data
public class LongTransaction {
    // 事务ID
    private String transactionId;
    // 事务时长(秒)
    private long durationSeconds;
    // 事务执行的SQL语句
    private String sql;
    // 事务状态
    private String state;
}

6.doris表创建,创建个存放指标结果的日志表

 CREATE TABLE `t_pg_log` (
  `id` varchar(200) NULL COMMENT "日志唯一标识",
  `database_name` varchar(150) NULL COMMENT "数据库名称",

  -- Checkpoint
  `checkpoints_timed` bigint NULL COMMENT "定时检查点次数",
  `checkpoints_req` bigint NULL COMMENT "请求触发检查点次数",

  -- WAL
  `wal_buffers_full` bigint NULL COMMENT "WAL 缓冲区写满次数",
  `wal_writes` bigint NULL COMMENT "WAL 写入次数(累计)",
  `wal_write_time` bigint NULL COMMENT "WAL 写入总耗时(微秒,累计)",
  `wal_syncs` bigint NULL COMMENT "WAL fsync 次数(累计)",
  `wal_sync_time` bigint NULL COMMENT "WAL fsync 总耗时(微秒,累计)",

  -- Checkpoint IO
  `checkpoint_write_time` bigint NULL COMMENT "检查点写入耗时(微秒,累计)",
  `checkpoint_sync_time` bigint NULL COMMENT "检查点同步耗时(微秒,累计)",

  -- 块 IO
  `blks_read` bigint NULL COMMENT "从磁盘读取的数据块数(累计)",
  `blks_hit` bigint NULL COMMENT "从缓存命中的数据块数(累计)",

  -- 表 IO
  `heap_blks_read` bigint NULL COMMENT '表数据块磁盘读取数(累计)',
  `heap_blks_hit` bigint NULL COMMENT '表数据块缓存命中数(累计)',
  `heap_cache_hit_ratio` double NULL COMMENT '表缓存命中率(%)',

  -- 索引 IO
  `idx_scan` bigint NULL COMMENT '索引总扫描次数(累计)',
  `idx_tup_read` bigint NULL COMMENT '索引返回总行数(累计)',
  `idx_tup_fetch` bigint NULL COMMENT '索引获取表数据总行数(累计)',

  -- 系统指标
  `cpu_usage` double NULL COMMENT "CPU 使用率(%)",
  `memory_usage` double NULL COMMENT "shared_buffers 内存占用(GB)",

  -- 连接
  `active_connections` int NULL COMMENT "活跃连接数",
  `idle_connections` int NULL COMMENT "空闲连接数",
  `locked_connections` int NULL COMMENT "被锁连接数",

  -- 事务
  `xact_commit` bigint NULL COMMENT "事务提交次数(累计)",
  `xact_rollback` bigint NULL COMMENT "事务回滚次数(累计)",
  `cache_hit_ratio` double NULL COMMENT "缓存命中率(%)",
  
  -- 复制槽日志
   `slot_state` varchar(65533) NULL COMMENT "同步状态",
  `write_lag_time` varchar(150) NULL COMMENT "写盘延迟时间",
  `flush_lag_time` varchar(150) NULL COMMENT "刷盘延迟时间",
  `replay_lag_time` varchar(150) NULL COMMENT "回放延迟时间",
  `write_backlog_size` varchar(150) NULL COMMENT "写盘积压大小",
  `flush_backlog_size` varchar(150) NULL COMMENT "刷盘积压大小",
  `replay_backlog_size` varchar(150) NULL COMMENT "回放积压大小",
  `slot_total_backlog_size` varchar(150) NULL COMMENT "复制槽总滞后大小",
  `unflush_backlog_size` varchar(150) NULL COMMENT "未刷盘滞后大小",
  `restart_lsn` varchar(150) NULL COMMENT "复制槽重启LSN(追赶起始位点)",
  `confirmed_flush_lsn` varchar(150) NULL COMMENT "已确认刷盘LSN(CDC持久化位点)",
  
  -- CDC 健康度
  `cdc_health_level` varchar(20) NULL COMMENT "CDC 健康等级(HEALTHY/WARN/CRITICAL)",
  `cdc_health_reason` varchar(500) NULL COMMENT "CDC 健康异常原因",
  `cdc_health_suggestion` varchar(500) NULL COMMENT "CDC 健康优化建议",

  -- 长事务
  `longTransactionJson` text NULL COMMENT "长事务json信息",

  `create_time` datetime NULL DEFAULT CURRENT_TIMESTAMP COMMENT "记录创建时间"
) ENGINE=OLAP
UNIQUE KEY(`id`)
COMMENT 'PostgreSQL 监控日志表'
DISTRIBUTED BY HASH(`id`) BUCKETS 3
PROPERTIES (
  "replication_allocation" = "tag.location.default: 1"
);

7.加个logback.xml辅助日志打印

<?xml version="1.0" encoding="UTF-8"?>
<configuration scan="true" scanPeriod="60 seconds">
    <!-- 关键修改:相对路径(与 Jar 包同级的 logs 目录) -->
    <property name="LOG_BASE_PATH" value="logs" /> <!-- 无需绝对路径,直接写目录名 -->
    <property name="LOG_FILE_NAME" value="pg-cdc-monitor" />
    <property name="MAX_HISTORY" value="30" />
    <property name="FILE_ENCODING" value="UTF-8" />

    <!-- 控制台输出(可选,服务器部署可注释) -->
    <appender name="CONSOLE" class="ch.qos.logback.core.ConsoleAppender">
        <encoder>
            <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
            <charset>${FILE_ENCODING}</charset>
        </encoder>
        <filter class="ch.qos.logback.classic.filter.ThresholdFilter">
            <level>INFO</level>
        </filter>
    </appender>

    <!-- 文件输出:相对路径滚动(核心不变) -->
    <appender name="FILE_ROLLING" class="ch.qos.logback.core.rolling.RollingFileAppender">
        <file>${LOG_BASE_PATH}/${LOG_FILE_NAME}.log</file>
        <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
            <fileNamePattern>${LOG_BASE_PATH}/${LOG_FILE_NAME}.%d{yyyy-MM-dd}.%i.log</fileNamePattern>
            <maxHistory>${MAX_HISTORY}</maxHistory>
            <timeBasedFileNamingAndTriggeringPolicy class="ch.qos.logback.core.rolling.SizeAndTimeBasedFNATP">
                <maxFileSize>200MB</maxFileSize>
            </timeBasedFileNamingAndTriggeringPolicy>
        </rollingPolicy>
        <encoder>
            <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
            <charset>${FILE_ENCODING}</charset>
        </encoder>
        <filter class="ch.qos.logback.classic.filter.ThresholdFilter">
            <level>INFO</level>
        </filter>
    </appender>

    <!-- 错误日志单独输出(相对路径) -->
    <appender name="ERROR_FILE" class="ch.qos.logback.core.rolling.RollingFileAppender">
        <file>${LOG_BASE_PATH}/${LOG_FILE_NAME}-error.log</file>
        <rollingPolicy class="ch.qos.logback.core.rolling.TimeBasedRollingPolicy">
            <fileNamePattern>${LOG_BASE_PATH}/${LOG_FILE_NAME}-error.%d{yyyy-MM-dd}.log</fileNamePattern>
            <maxHistory>${MAX_HISTORY}</maxHistory>
        </rollingPolicy>
        <encoder>
            <pattern>%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{50} - %msg%n</pattern>
            <charset>${FILE_ENCODING}</charset>
        </encoder>
        <filter class="ch.qos.logback.classic.filter.LevelFilter">
            <level>ERROR</level>
            <onMatch>ACCEPT</onMatch>
            <onMismatch>DENY</onMismatch>
        </filter>
    </appender>

    <root level="INFO">
        <appender-ref ref="FILE_ROLLING" />
        <appender-ref ref="ERROR_FILE" />
         <appender-ref ref="CONSOLE" />
    </root>

    <!-- 第三方依赖日志控制 -->
    <logger name="org.postgresql" level="WARN" />
    <logger name="com.mysql.cj" level="WARN" />
    <logger name="java.sql" level="WARN" />
</configuration>

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐