Redis Stream与WebSocket投票系统渲染速度优化咨询
优化投票系统POST请求后得分渲染速度的方案
当前流程与问题
当前投票流程:
- API将投票信息写入Redis Stream
- checkinprocessor.js读取流数据
- 计算新的投票得分并更新数据库
- checkinreceiver.js等待处理完成后,将新的平均得分返回给UI
问题:POST请求后新投票得分的渲染速度过慢,现有实现逻辑存在优化空间。
现有代码与配置
系统配置文件 conf.json
{ "logLevel": "debug", "session": { "secret": "red15alltheTh1ngZ!", "appName": "checkinapp", "keyPrefix": "session" }, "application": { "port": 8081 }, "auth": { "port": 8083 }, "checkinReceiver": { "port": 8082, "maxStreamLength": 500000 }, "redis": { "host": "localhost", "port": 6379, "password": null, "keyPrefix": "ncc" } }
Express.js server.js 代码
const config = require("better-config"); const express = require("express"); const morgan = require("morgan"); const cors = require("cors"); const routes = require("./routes"); const logger = require("./utils/logger"); const port = config.get("application.port"); const helmet = require("helmet"); const { Server } = require("ws"); // Importing the 'ws' package const { createServer } = require("http"); // you can use https as well const pgClient = require("./utils/db"); // Import PostgreSQL client // Load the configuration config.set(`../${process.env.CRASH_COURSE_CONFIG_FILE || "config.json"}`); const app = express(); const server = createServer(app); app.use(helmet()); app.use(morgan("combined", { stream: logger.stream })); app.use(cors()); app.use("/api", routes); app.get("/", (_, res) => { res.json({}); }); // Create a WebSocket server const wss = new Server({ noServer: true }); // Handle WebSocket connection wss.on("connection", (ws) => { console.log("Client connected"); // Send the latest items to the client when they connect broadcastLatestItems(ws); // Handle incoming messages from the client ws.on("message", (message) => { console.log(`Received: ${message}`); // Handle messages from the client here }); ws.on("close", () => { console.log("Client disconnected"); }); }); // Function to broadcast latest items to all connected clients async function broadcastLatestItems() { try { const result = await pgClient.query( "SELECT * FROM items ORDER BY numvotes DESC LIMIT 100" ); const latestItems = JSON.stringify(result.rows); wss.clients.forEach((client) => { if (client.readyState === client.OPEN) { client.send(latestItems); } }); } catch (error) { logger.error("Failed to fetch items from SQL", error); } } setInterval(broadcastLatestItems, 2000); // Send updates every 2 seconds // Upgrade HTTP server to handle WebSocket requests server.on("upgrade", (request, socket, head) => { wss.handleUpgrade(request, socket, head, (ws) => { wss.emit("connection", ws, request); }); }); // Start the server server.listen(port, () => { logger.info(`Application listening on port ${port}.`); }); // Check for required environment variables. if (process.env.WEATHER_API_KEY === undefined) { console.warn("Warning: Environment variable WEATHER_API_KEY is not set!"); }
Redis API接收器 checkinreceiver.js 代码
const config = require("better-config"); const express = require("express"); const { body } = require("express-validator"); const morgan = require("morgan"); const cors = require("cors"); const logger = require("./utils/logger"); const apiErrorReporter = require("./utils/apierrorreporter"); const useAuth = process.argv[2] === "auth"; const pgClient = require("./utils/db"); // Import PostgreSQL client config.set(`../${process.env.CRASH_COURSE_CONFIG_FILE || "config.json"}`); const redis = require("./utils/redisclient"); const session = require("express-session"); const RedisStore = require("connect-redis").default; const redisClient = redis.getClient(); const pubsubClient = redis.getClient(); // Separate Redis client for Pub/Sub const app = express(); app.use(morgan("combined", { stream: logger.stream })); app.use(cors()); app.use(express.json()); // Initialize store. let redisStore = new RedisStore({ client: redisClient, prefix: redis.getKeyName(`${config.session.keyPrefix}:`), }); if (useAuth) { logger.info("Authentication enabled, checkins require a valid user session."); app.use( session({ secret: config.session.secret, store: redisStore, name: config.session.appName, resave: false, saveUninitialized: true, }) ); } else { logger.info( "Authentication disabled, checkins do not require a valid user session." ); } const votesStreamKey = redis.getKeyName("votes"); const maxStreamLength = config.get("checkinReceiver.maxStreamLength"); app.post( "/api/checkin", (req, res, next) => { if (useAuth && !req.session.user) { logger.debug("Rejecting checkin - no valid user session found."); return res.status(401).send("Authentication required."); } return next(); }, [ body().isObject(), body("userId").isInt({ min: 1 }), body("itemId").isInt({ min: 1 }), body("starRating").isInt({ min: 0, max: 5 }), apiErrorReporter, ], async (req, res) => { const vote = req.body; const pipeline = redisClient.pipeline(); pipeline.xadd( votesStreamKey, "MAXLEN", "~", maxStreamLength, "*", ...Object.entries(vote).flat(), (err, result) => { if (err) { logger.error("Error adding checkin to stream:"); logger.error(err); } else { logger.debug(`Received checkin, added to stream as ${result}`); incomingId = result; } } ); await pipeline.exec(); // Wait for the notification from checkinprocessor.js and fetch the updated data try { const items = await waitForCheckinCompletion(); return res.status(200).json(items); } catch (error) { logger.error("Error fetching new data:", error); return res.status(500).send("Internal Server Error"); } } ); // Pub/Sub based function to fetch new data when notified by checkinprocessor.js async function waitForCheckinCompletion() { return new Promise((resolve, reject) => { const pubChannel = redis.getKeyName("checkin-complete"); // Subscribe to the Redis channel pubsubClient.subscribe(pubChannel); // Listen for published messages pubsubClient.on("message", async (channel, message) => { console.log("message", message); console.log("channel", channel); console.log("pubChannel", pubChannel); if (message === result) { try { // Fetch data from PostgreSQL when notification is received const result = await pgClient.query( "SELECT * FROM items ORDER BY numvotes DESC LIMIT 100" ); resolve(result.rows); } catch (error) { reject(error); } finally { // Unsubscribe after receiving the message to prevent multiple triggers pubsubClient.unsubscribe(pubChannel); } } }); // Handle errors in Redis Pub/Sub pubsubClient.on("error", (err) => { reject(err); }); }); } const port = config.get("checkinReceiver.port"); app.listen(port, () => { logger.info(`Checkin receiver listening on port ${port}.`); });
UI端Next.js POST请求代码
const upvoteItem = async (itemId: number) => { try { const response = await fetch("http://localhost:8082/api/checkin", { method: "POST", headers: { "Content-Type": "application/json", }, body: JSON.stringify({ userId: Math.ceil(Math.random() * (5 - 1) + 1), itemId: itemId, starRating: Math.ceil(Math.random() * (4 - 0) + 0), }), }); if (!response.ok) { throw new Error("Failed to upvote item"); } // Parse the JSON response const data = await response.json(); // Set the data (assuming setData is defined in your component) setData(data); // Optionally display modal or handle UI updates console.log("Upvote successful:", data); } catch (error) { console.error("Error upvoting item:", error); } };
优化方案
1. 异步响应+WebSocket实时推送
当前POST请求等待处理完成才返回,导致响应延迟。优化方式:
- POST请求写入Redis Stream后立即返回成功,不等待处理结果
- 当checkinprocessor完成计算并更新数据库后,通过WebSocket主动推送更新后的得分到所有客户端
- UI端收到WebSocket推送后自动更新数据,无需等待POST响应
修改checkinreceiver.js的POST接口:
app.post( "/api/checkin", // 中间件保持不变 async (req, res) => { const vote = req.body; const pipeline = redisClient.pipeline(); pipeline.xadd( votesStreamKey, "MAXLEN", "~", maxStreamLength, "*", ...Object.entries(vote).flat(), (err, result) => { if (err) { logger.error("Error adding checkin to stream:"); logger.error(err); return res.status(500).send("Failed to submit vote"); } else { logger.debug(`Received checkin, added to stream as ${result}`); // 立即返回成功,不等待处理 return res.status(202).json({ message: "Vote submitted successfully", streamId: result }); } } ); await pipeline.exec(); } );
修改checkinprocessor.js,完成数据库更新后触发WebSocket推送:
// 当完成计算并更新数据库后 async function processVote(vote) { // 计算得分、更新数据库逻辑 // 通过Redis Pub/Sub通知server.js触发推送 const pubChannel = redis.getKeyName("checkin-complete"); await redisClient.publish(pubChannel, "updated"); }
在server.js中监听Redis Pub/Sub,收到通知后立即推送:
// 在server.js中添加Redis Pub/Sub监听 const pubsubClient = redis.getClient(); const pubChannel = redis.getKeyName("checkin-complete"); pubsubClient.subscribe(pubChannel); pubsubClient.on("message", async (channel, message) => { if (channel === pubChannel) { await broadcastLatestItems(); } }); // 可选:保留固定轮询作为兜底,或直接移除 // setInterval(broadcastLatestItems, 2000);
2. 本地UI乐观更新
用户投票后,先在本地临时更新该项目的得分,等到WebSocket推送真实数据后再替换,让用户立刻看到反馈:
const upvoteItem = async (itemId: number) => { try { // 1. 乐观更新UI const tempRating = Math.ceil(Math.random() * (4 - 0) + 0); setData(prevData => prevData.map(item => { if (item.itemId === itemId) { // 临时计算新得分(根据实际逻辑调整) const newNumVotes = item.numvotes + 1; const newTotalRating = item.totalrating + tempRating; return { ...item, numvotes: newNumVotes, averagerating: newTotalRating / newNumVotes }; } return item; })); // 2. 发送POST请求 const response = await fetch("http://localhost:8082/api/checkin", { method: "POST", headers: { "Content-Type": "application/json", }, body: JSON.stringify({ userId: Math.ceil(Math.random() * (5 - 1) + 1), itemId: itemId, starRating: tempRating, }), }); if (!response.ok) { throw new Error("Failed to upvote item"); } console.log("Upvote submitted successfully"); } catch (error) { console.error("Error upvoting item:", error); // 乐观更新失败,恢复原数据 // fetchLatestData(); } };
3. 优化Redis Stream与处理逻辑
- 使用Redis Stream的
xreadgroup消费组模式,避免重复处理,提高消费效率 - 用Redis Hash存储每个项目的总评分和投票数,更新时直接在Redis中累加,定期同步到数据库;或处理投票时先更新Redis缓存,再异步更新数据库,UI优先从Redis获取数据
4. 修复现有Pub/Sub的bug
当前waitForCheckinCompletion函数存在变量作用域问题,if (message === result)中的result未正确传递,导致无法匹配。修复方式:
async (req, res) => { const vote = req.body; let incomingId; const pipeline = redisClient.pipeline(); pipeline.xadd( votesStreamKey, "MAXLEN", "~", maxStreamLength, "*", ...Object.entries(vote).flat(), (err, result) => { if (err) { logger.error("Error adding checkin to stream:"); logger.error(err); } else { logger.debug(`Received checkin, added to stream as ${result}`); incomingId = result; } } ); await pipeline.exec(); try { // 传递incomingId到函数中 const items = await waitForCheckinCompletion(incomingId); return res.status(200).json(items); } catch (error) { logger.error("Error fetching new data:", error); return res.status(500).send("Internal Server Error"); } } async function waitForCheckinCompletion(expectedId) { return new Promise((resolve, reject) => { const pubChannel = redis.getKeyName("checkin-complete"); pubsubClient.subscribe(pubChannel); pubsubClient.on("message", async (channel, message) => { if (channel === pubChannel && message === expectedId) { try { const result = await pgClient.query( "SELECT * FROM items ORDER BY numvotes DESC LIMIT 100" ); resolve(result.rows); } catch (error) { reject(error); } finally { pubsubClient.unsubscribe(pubChannel); } } }); pubsubClient.on("error", (err) => { reject(err); }); // 添加超时处理,避免请求挂起 setTimeout(() => { pubsubClient.unsubscribe(pubChannel); reject(new Error("Wait for checkin completion timed out")); }, 5000); }); }
内容的提问来源于stack exchange,提问作者Canadian Lonman
相关产品推荐
相关产品推荐

