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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 04:28:09