You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Node.js实时检测PostgreSQL数据库的变更?

高效实现PostgreSQL数据库变更实时检测(替代定时轮询)

问题背景

原先使用node-schedule每分钟轮询PostgreSQL,对比MongoDB存储的统计数据,更新用户、分类、文章的数量统计。希望彻底减少无效查询,实现数据实时更新,当前技术栈为Knex(查询构建器)+ MongoDB(统计存储)。

原定时轮询代码:

const schedule = require('node-schedule')

module.exports = app => {
    schedule.scheduleJob('*/1 * * * *', async function () {
        const usersCount = await app.db('users').count('id').first()
        const categoriesCount = await app.db('categories').count('id').first()
        const articlesCount = await app.db('articles').count('id').first()

        const { Stat } = app.api.stat

        const lastStat = await Stat.findOne({}, {},
            { sort: { 'createdAt': -1 } })
        
        const stat = new Stat({
            users: usersCount.count,
            categories: categoriesCount.count,
            articles: articlesCount.count,
            createdAt: new Date()
        })

        const hasChanged = !lastStat
            || stat.users !== lastStat.users
            || stat.categories !== lastStat.categories
            || stat.articles !== lastStat.articles;
        
        if (hasChanged) {
            stat.save().then(() => console.log('[STATS] Estatísticas atualizadas!'))
        }
    })
}

已尝试PostgreSQL触发器+LISTEN/NOTIFY方案,但仅针对categories表配置,需完善实现。


解决方案:PostgreSQL触发器+实时监听

1. 完善数据库触发器与通知函数

给users、categories、articles三个表添加语句级触发器(避免批量操作发送大量通知),当表发生增/删/改时,向Node.js服务发送变更通知。

执行以下SQL(建议放在Knex迁移文件中,而非每次启动执行):

CREATE OR REPLACE FUNCTION public.notify_stat_change() RETURNS trigger AS $$
BEGIN
    -- 传递表名作为通知内容,便于排查,但核心只需知道有变更即可
    PERFORM pg_notify('stat_change', TG_TABLE_NAME);
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

-- 幂等创建users表触发器
DO $$
BEGIN
    IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'users_stat_change_trigger') THEN
        CREATE TRIGGER users_stat_change_trigger
        AFTER INSERT OR UPDATE OR DELETE ON users
        FOR EACH STATEMENT EXECUTE FUNCTION public.notify_stat_change();
    END IF;
END $$;

-- 幂等创建categories表触发器
DO $$
BEGIN
    IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'categories_stat_change_trigger') THEN
        CREATE TRIGGER categories_stat_change_trigger
        AFTER INSERT OR UPDATE OR DELETE ON categories
        FOR EACH STATEMENT EXECUTE FUNCTION public.notify_stat_change();
    END IF;
END $$;

-- 幂等创建articles表触发器
DO $$
BEGIN
    IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'articles_stat_change_trigger') THEN
        CREATE TRIGGER articles_stat_change_trigger
        AFTER INSERT OR UPDATE OR DELETE ON articles
        FOR EACH STATEMENT EXECUTE FUNCTION public.notify_stat_change();
    END IF;
END $$;

2. 重构Node.js监听逻辑

移除定时轮询,改用PG客户端监听数据库通知,收到变更后再查询统计数据并更新MongoDB:

const config = require('../knexfile')
const knex = require('knex')(config)
const { Client } = require('pg')

// 执行Knex迁移
knex.migrate.latest([config])

// 复用Knex的数据库配置创建PG客户端
const pgClient = new Client({
    connectionString: knex.client.connectionSettings.connectionString
})

// 初始化触发器(可选:若未通过迁移执行SQL,可在启动时执行)
async function initTriggers() {
    try {
        await knex.raw(`
            CREATE OR REPLACE FUNCTION public.notify_stat_change() RETURNS trigger AS $$
            BEGIN
                PERFORM pg_notify('stat_change', TG_TABLE_NAME);
                RETURN NEW;
            END;
            $$ LANGUAGE plpgsql;

            DO $$
            BEGIN
                IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'users_stat_change_trigger') THEN
                    CREATE TRIGGER users_stat_change_trigger
                    AFTER INSERT OR UPDATE OR DELETE ON users
                    FOR EACH STATEMENT EXECUTE FUNCTION public.notify_stat_change();
                END IF;

                IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'categories_stat_change_trigger') THEN
                    CREATE TRIGGER categories_stat_change_trigger
                    AFTER INSERT OR UPDATE OR DELETE ON categories
                    FOR EACH STATEMENT EXECUTE FUNCTION public.notify_stat_change();
                END IF;

                IF NOT EXISTS (SELECT 1 FROM pg_trigger WHERE tgname = 'articles_stat_change_trigger') THEN
                    CREATE TRIGGER articles_stat_change_trigger
                    AFTER INSERT OR UPDATE OR DELETE ON articles
                    FOR EACH STATEMENT EXECUTE FUNCTION public.notify_stat_change();
                END IF;
            END $$;
        `)
        console.log('数据库触发器初始化完成')
    } catch (err) {
        console.error('触发器初始化失败:', err)
    }
}

// 监听数据库变更并更新统计数据
async function setupStatListener(app) {
    await pgClient.connect()
    console.log('PG通知客户端连接成功')

    // 监听指定通知频道
    await pgClient.query('LISTEN stat_change')

    pgClient.on('notification', async () => {
        console.log('收到数据库变更通知,开始更新统计数据')
        try {
            // 并行查询三个表的统计数,提升效率
            const [usersCount, categoriesCount, articlesCount] = await Promise.all([
                knex('users').count('id').first(),
                knex('categories').count('id').first(),
                knex('articles').count('id').first()
            ])

            const { Stat } = app.api.stat
            const lastStat = await Stat.findOne({}, {}, { sort: { 'createdAt': -1 } })
            
            const stat = new Stat({
                users: usersCount.count,
                categories: categoriesCount.count,
                articles: articlesCount.count,
                createdAt: new Date()
            })

            const hasChanged = !lastStat
                || stat.users !== lastStat.users
                || stat.categories !== lastStat.categories
                || stat.articles !== lastStat.articles;
            
            if (hasChanged) {
                await stat.save()
                console.log('[STATS] 统计数据已更新!')
            }
        } catch (err) {
            console.error('更新统计数据失败:', err)
        }
    })

    // 处理客户端连接错误,可选添加重连逻辑
    pgClient.on('error', (err) => {
        console.error('PG客户端连接异常:', err)
        setTimeout(() => setupStatListener(app), 5000)
    })
}

// 导出Knex实例与初始化函数
module.exports = {
    knex,
    initTriggers,
    setupStatListener
}

3. 应用启动时初始化

在应用入口文件中调用初始化逻辑:

const { knex, initTriggers, setupStatListener } = require('./path/to/your/knex-config-file')

// 假设app为Express/Koa实例
async function startApp(app) {
    await initTriggers()
    await setupStatListener(app)
    app.listen(3000, () => console.log('应用启动成功'))
}

startApp(app)

核心优化点

  • 消除无效轮询:仅在数据库发生实际变更时执行统计查询,彻底避免每分钟空查
  • 减少冗余通知:使用FOR EACH STATEMENT触发器,批量操作仅发送一条通知
  • 幂等性保障:触发器创建时添加存在性判断,避免重复执行报错
  • 稳定性提升:添加PG客户端错误监听与自动重连逻辑
  • 查询效率优化:用Promise.all并行查询三个表的统计数

内容的提问来源于stack exchange,提问作者Vanortton

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 19:07:08