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

WebSocket对接PostgreSQL时数据变更客户端未实时更新问题咨询

问题根因

原服务端仅在客户端首次建立WebSocket连接时查询一次数据库并返回数据,没有监听数据库变更并主动向客户端推送最新数据,所以数据更新后客户端无法感知。

修复方案

服务端修改(推荐用PostgreSQL LISTEN/NOTIFY 实时通知方案)

步骤1:给PostgreSQL表添加变更通知触发器

先执行SQL创建触发器函数和触发器,每次表数据增删改时自动发送通知:

-- 创建通知函数
CREATE OR REPLACE FUNCTION notify_temp_change()
RETURNS trigger AS $$
BEGIN
  PERFORM pg_notify('temp_data_changed', '');
  RETURN NULL;
END;
$$ LANGUAGE plpgsql;

-- 绑定触发器到温度表
CREATE TRIGGER trigger_temp_change
AFTER INSERT OR UPDATE OR DELETE ON my_temp_table
FOR EACH STATEMENT EXECUTE FUNCTION notify_temp_change();

步骤2:修改服务端代码

监听PostgreSQL的变更通知,收到通知后拉取最新数据推送给所有已连接的客户端:

const express = require('express')
const app = express()
const server = require('http').createServer(app);
// 注意:Node.js连接PostgreSQL官方包为pg,需先执行npm install pg安装
const { Pool } = require('pg');
const WebSocket = require('ws');

// 初始化PG连接池,替换为你自己的数据库配置
const pool = new Pool({
  host: 'localhost',
  user: '你的数据库用户名',
  password: '你的数据库密码',
  database: '你的数据库名',
  port: 5432
});

const wss = new WebSocket.Server({ server:server });

const getTempData = async () => {
  try {
    const tempData = await pool.query("select country, temp from my_temp_table");
    return JSON.stringify(tempData.rows)
  } catch(err) {
      // 修复原代码拼写错误:messasge → message
      console.error(err.message);
  }
}

// 启动PG变更监听
(async () => {
  const notifyClient = await pool.connect();
  try {
    await notifyClient.query('LISTEN temp_data_changed');
    notifyClient.on('notification', async () => {
      // 收到表变更通知,拉取最新数据
      const latestTemp = await getTempData();
      // 推送给所有在线的WebSocket客户端
      wss.clients.forEach(client => {
        if (client.readyState === WebSocket.OPEN) {
          client.send(latestTemp);
        }
      })
    })
  } finally {
    notifyClient.release();
  }
})()

wss.on('connection', async (webSocketClient) => {
  console.log('A new client Connected!');   
  const tempDetails = await getTempData();
  // 首次连接返回全量数据
  webSocketClient.send(tempDetails);      
  webSocketClient.on('message', (message) => {
    console.log('received: %s', message);    
  });
});           

server.listen(3000, () => console.log(`Listening on port :3000`))

轻量替代方案(无需修改数据库,用定时轮询)

如果不想配置PG触发器,可在服务端添加定时任务定期拉取数据库,对比数据变化后推送,修改服务端代码添加如下内容即可:

// 缓存上一次推送的温度数据
let lastSentData = '';
// 每3秒拉取一次数据库,可按需调整间隔
setInterval(async () => {
  const currentData = await getTempData();
  // 只有数据发生变化才推送,节省带宽
  if (currentData && currentData !== lastSentData) {
    lastSentData = currentData;
    wss.clients.forEach(client => {
      if (client.readyState === WebSocket.OPEN) client.send(currentData);
    })
  }
}, 3000)

客户端代码优化(可选)

你现有的客户端逻辑可以正常接收推送更新,仅做小优化避免重复注册事件:

useEffect(() => {
    if (!ws.current) return;
    const handleMessage = e => {
        if (isPaused) return;
        console.log("getting temp data....");
        const data = JSON.parse(e.data);
        setTempData(data)          
        console.log("data: ",data);
    };
    ws.current.addEventListener('message', handleMessage);
    return () => ws.current?.removeEventListener('message', handleMessage);
}, [isPaused]);

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 10:27:04