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

如何用Redis Pub/Sub结合Node.js async/await解决跨Express应用高流量数据查询的性能问题

如何用Redis Pub/Sub结合Node.js async/await解决跨Express应用高流量数据查询的性能问题

兄弟,我太懂你这种高流量接口的头疼事儿了——直接跨应用HTTP调用简直是性能杀手,用Redis Pub/Sub结合async/await来优化确实是个靠谱的路子,既避免了频繁HTTP请求的开销,还能让代码逻辑更清爽。结合你说的场景(应用A不能直接访问B的数据库,高流量端点需要查B的数据),我给你一步步拆解这个方案的落地方式:

整体流程梳理

核心思路就是缓存优先+Pub/Sub异步请求兜底:

  1. 应用A的高流量端点先查Redis缓存,命中直接返回数据;
  2. 缓存未命中时,A通过Redis Pub/Sub给B发查询请求;
  3. B收到请求后从自己的数据库查数据,把数据存入Redis缓存,再通过Pub/Sub给A返回结果;
  4. A收到B的响应后,返回数据给客户端,同时后续请求就能直接走Redis缓存了。

具体代码实现(基于Ioredis)

第一步:配置应用B的Redis订阅与处理逻辑

应用B需要订阅A的请求频道,收到请求后查库、存缓存、回传结果:

const Redis = require('ioredis');
const redis = new Redis(); // 连接到共用的Redis实例
const db = require('./your-b-db-connection'); // 替换成B的数据库连接

// 初始化Redis订阅器
async function setupRedisSubscriber() {
  const subscriber = new Redis();
  // 订阅A发送请求的频道
  await subscriber.subscribe('request:from:A');

  subscriber.on('message', async (channel, message) => {
    if (channel !== 'request:from:A') return;

    try {
      // 解析A发来的请求内容(包含唯一请求ID和要查询的数据标识)
      const { requestId, dataKey } = JSON.parse(message);
      
      // 从B的数据库查询目标数据
      const targetData = await db.query('SELECT * FROM your_table WHERE key = ?', [dataKey]);

      if (targetData) {
        // 把数据存入Redis,设置5分钟过期时间(可根据业务调整)
        await redis.setex(`b:data:${dataKey}`, 300, JSON.stringify(targetData));
        // 给A发送响应,带上请求ID做区分
        await redis.publish('response:to:A', JSON.stringify({
          requestId,
          data: targetData
        }));
      } else {
        // 没有查询到数据也通知A,避免A一直等待
        await redis.publish('response:to:A', JSON.stringify({
          requestId,
          data: null
        }));
      }
    } catch (err) {
      console.error('处理A的请求出错:', err);
      // 错误场景下也要给A发响应,避免请求挂起
      const { requestId } = JSON.parse(message) || {};
      if (requestId) {
        await redis.publish('response:to:A', JSON.stringify({
          requestId,
          error: 'Failed to fetch data from B\'s DB'
        }));
      }
    }
  });
}

// 启动订阅
setupRedisSubscriber().catch(console.error);

第二步:配置应用A的Redis订阅与高流量端点

应用A需要订阅B的响应频道,处理缓存未命中时的异步请求,同时用async/await让异步逻辑更直观:

const Redis = require('ioredis');
const { v4: uuidv4 } = require('uuid'); // 生成唯一请求ID,避免响应混淆
const express = require('express');
const app = express();
const redis = new Redis(); // 连接到共用的Redis实例

// 用Map存储待处理的请求,key是请求ID,value是Promise的resolve/reject方法
const pendingRequests = new Map();

// 初始化Redis订阅器,监听B的响应
async function setupRedisSubscriber() {
  const subscriber = new Redis();
  await subscriber.subscribe('response:to:A');

  subscriber.on('message', async (channel, message) => {
    if (channel !== 'response:to:A') return;

    try {
      const { requestId, data, error } = JSON.parse(message);
      const resolver = pendingRequests.get(requestId);
      if (resolver) {
        // 根据B的响应结果,触发Promise的成功或失败回调
        if (error) {
          resolver.reject(new Error(error));
        } else {
          resolver.resolve(data);
        }
        // 处理完后移除待处理请求
        pendingRequests.delete(requestId);
      }
    } catch (err) {
      console.error('处理B的响应出错:', err);
    }
  });
}

setupRedisSubscriber().catch(console.error);

// 高流量端点实现
app.get('/high-traffic/:dataKey', async (req, res) => {
  const { dataKey } = req.params;
  const redisCacheKey = `b:data:${dataKey}`;

  try {
    // 第一步:先查Redis缓存
    const cachedData = await redis.get(redisCacheKey);
    if (cachedData) {
      return res.json(JSON.parse(cachedData));
    }

    // 缓存未命中,生成唯一请求ID
    const requestId = uuidv4();
    // 给B发送查询请求
    await redis.publish('request:from:A', JSON.stringify({
      requestId,
      dataKey
    }));

    // 创建Promise等待B的响应,设置3秒超时(可根据业务调整)
    const responsePromise = new Promise((resolve, reject) => {
      pendingRequests.set(requestId, { resolve, reject });
      // 超时处理:避免请求一直挂着
      setTimeout(() => {
        reject(new Error('Request timeout waiting for B'));
        pendingRequests.delete(requestId);
      }, 3000);
    });

    // 等待B的响应
    const data = await responsePromise;
    if (data) {
      res.json(data);
    } else {
      res.status(404).json({ message: 'Data not found in B\'s DB' });
    }
  } catch (err) {
    console.error('端点处理出错:', err);
    // 兜底逻辑:超时或错误时,降级为HTTP调用B(保证可用性)
    try {
      const fallbackResponse = await fetch(`http://your-app-b-url/api/data/${dataKey}`);
      const fallbackData = await fallbackResponse.json();
      // 把兜底拿到的数据存入Redis,下次请求直接走缓存
      await redis.setex(redisCacheKey, 300, JSON.stringify(fallbackData));
      res.json(fallbackData);
    } catch (fallbackErr) {
      res.status(500).json({ message: 'Failed to fetch data from B' });
    }
  }
});

app.listen(3000, () => console.log('App A running on port 3000'));

关键注意事项

  • 请求ID的唯一性:必须用UUID这类唯一标识区分每个请求,否则高流量下多个请求的响应会混淆,这是核心细节;
  • 超时与兜底:绝对不能让请求无限等待,设置合理的超时时间,超时后用HTTP调用兜底,保证服务可用性;
  • 缓存过期策略:根据B的数据更新频率设置Redis键的过期时间,避免脏数据;如果B的数据主动更新了,也可以额外加一个Pub/Sub频道,让B通知A清除对应缓存;
  • 消息幂等性:Redis Pub/Sub可能会出现重复消息,要保证B的查库和缓存操作是幂等的(比如多次查询同一数据不会产生副作用);
  • Redis连接稳定性:Ioredis自带重连机制,但最好加上连接错误监听,避免订阅中断导致整个流程失效;
  • 监控与调优:可以加缓存命中率、请求超时率、兜底调用次数的监控,根据数据调整缓存过期时间、超时时间等参数。

备注:内容来源于stack exchange,提问作者Scrogneugneu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.16 11:22:58