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

如何基于oat++实现SSE或HTTP长轮询?求相关示例

在oat++中实现Server-Sent Events(SSE)与HTTP长轮询

当然可以在oat中实现这两种服务器推送技术,oat的异步IO模型天然适配这类需要长时间保持连接的场景,以下是具体实现示例:


Server-Sent Events(SSE) 实现示例

SSE核心是服务器向客户端发送text/event-stream格式的响应,保持连接持续推送数据:

#include "oatpp/web/server/api/ApiController.hpp"
#include "oatpp/core/macro/codegen.hpp"
#include "oatpp/core/macro/component.hpp"
#include <chrono>
#include <thread>

#include OATPP_CODEGEN_BEGIN(ApiController)

class SSEController : public oatpp::web::server::api::ApiController {
public:
  SSEController(OATPP_COMPONENT(std::shared_ptr<ObjectMapper>, objectMapper))
    : ApiController(objectMapper) {}

  ENDPOINT_ASYNC("GET", "/sse", SSEEndpoint) {

    ENDPOINT_ASYNC_INIT(SSEEndpoint)

    Action act() override {
      // 设置SSE标准响应头
      auto response = createResponse(Status::CODE_200, "");
      response->putHeader(Header::CONTENT_TYPE, "text/event-stream");
      response->putHeader(Header::CACHE_CONTROL, "no-cache");
      response->putHeader(Header::CONNECTION, "keep-alive");

      auto outputStream = response->getOutputStream();

      // 持续推送事件,直到客户端断开连接
      return outputStream->writeSimple("event: heartbeat\n")->callbackTo([this, outputStream]() {
        return outputStream->writeSimple("data: oat++ SSE push at " + std::to_string(time(nullptr)) + "\n\n")->callbackTo([this]() {
          std::this_thread::sleep_for(std::chrono::seconds(2));
          return act();
        });
      });
    }
  };

};

#include OATPP_CODEGEN_END(ApiController)

关键说明:

  • 使用ENDPOINT_ASYNC定义异步端点,避免长时间连接阻塞线程池
  • 严格遵循SSE格式:event: [事件名]\ndata: [内容]\n\n
  • 客户端断开时,oat++会自动清理连接资源

HTTP长轮询实现示例

长轮询逻辑是客户端发起请求后,服务器保持连接直到有数据返回或超时,客户端需在收到响应后重新发起请求:

#include "oatpp/web/server/api/ApiController.hpp"
#include "oatpp/core/macro/codegen.hpp"
#include "oatpp/core/macro/component.hpp"
#include <chrono>
#include <mutex>
#include <condition_variable>
#include <string>

#include OATPP_CODEGEN_BEGIN(ApiController)

// 模拟业务消息队列
class MessageQueue {
private:
  std::mutex m_mutex;
  std::condition_variable m_cv;
  std::string m_pendingMsg;
  bool m_hasNewMsg = false;
public:
  void push(const std::string& msg) {
    std::lock_guard<std::mutex> lock(m_mutex);
    m_pendingMsg = msg;
    m_hasNewMsg = true;
    m_cv.notify_one();
  }

  std::string wait(std::chrono::seconds timeout) {
    std::unique_lock<std::mutex> lock(m_mutex);
    if(m_cv.wait_for(lock, timeout, [this](){ return m_hasNewMsg; })) {
      auto msg = m_pendingMsg;
      m_hasNewMsg = false;
      return msg;
    }
    return "";
  }
};

class LongPollController : public oatpp::web::server::api::ApiController {
private:
  std::shared_ptr<MessageQueue> m_msgQueue;
public:
  LongPollController(OATPP_COMPONENT(std::shared_ptr<ObjectMapper>, objectMapper),
                     std::shared_ptr<MessageQueue> msgQueue)
    : ApiController(objectMapper), m_msgQueue(msgQueue) {}

  ENDPOINT_ASYNC("GET", "/long-poll", LongPollEndpoint) {

    ENDPOINT_ASYNC_INIT(LongPollEndpoint)

    Action act() override {
      // 等待30秒,超时则返回空响应
      auto msg = m_msgQueue->wait(std::chrono::seconds(30));
      if(!msg.empty()) {
        return createResponse(Status::CODE_200, msg)->send();
      } else {
        return createResponse(Status::CODE_204, "Timeout")->send();
      }
    }
  };

  // 模拟业务触发消息推送的接口
  ENDPOINT("POST", "/trigger", triggerMsg,
           BODY_STRING(String, message)) {
    m_msgQueue->push(message);
    return createResponse(Status::CODE_200, "Message queued");
  };

};

#include OATPP_CODEGEN_END(ApiController)

关键说明:

  • 用条件变量实现阻塞等待,避免空轮询浪费资源
  • 超时时间可根据业务需求调整,防止连接长期闲置
  • 客户端需在收到响应后立即重新发起请求,维持长轮询链路

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 04:00:32