如何用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异步请求兜底:
- 应用A的高流量端点先查Redis缓存,命中直接返回数据;
- 缓存未命中时,A通过Redis Pub/Sub给B发查询请求;
- B收到请求后从自己的数据库查数据,把数据存入Redis缓存,再通过Pub/Sub给A返回结果;
- 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
相关产品推荐
相关产品推荐

