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

Google Pub/Sub C++ 有序订阅间歇性停止拉取消息无报错问题

问题根因

订阅随机停止拉取消息、无错误日志、死信队列为空的故障由代码中多处逻辑缺陷共同导致,和Pub/Sub服务本身无关:

  • 消息回调无异常兜底:你在订阅回调中直接调用std::stoi(m.attributes()["topic"])、messageCallback(msg),只要这两处抛出任何未捕获异常,会直接中断当前消息的处理流程,既不会ack也不会nack。由于你设置了并发数、最大待处理消息数均为1,卡住的消息会占满整个流的全部配额,后续所有消息都无法被拉取;叠加有序订阅的特性:同一个ordering key的消息只要前一个未处理完成,后续消息永远不会投递,会直接导致消费完全停滞。
  • 错误回调触发逻辑不符合预期:then()注册的回调只有在订阅流完全终止、返回最终状态时才会执行,单消息阻塞、流临时卡死的场景下,订阅连接本身没有断开,f.get()会一直阻塞,根本不会走到你打错误日志的逻辑,因此看不到任何报错输出。
  • 死信队列不生效的原因:死信队列仅在消息被明确nack、或投递次数超过配置阈值时才会接收消息,你的场景中消息既没被ack也没被nack,只是卡在回调流程中,永远达不到你设置的5次投递阈值,因此死信队列为空。
  • 资源重复创建存在稳定性隐患:你每次发送消息都重新创建Publisher连接、重复调用Topic/Subscription创建接口,会产生大量不必要的控制面API调用,触发Pub/Sub配额限流时会阻塞订阅流的后台工作线程,进一步加剧拉取中断问题。
  • 存在笔误bug:createOrderingSubscription函数中创建死信订阅失败时,抛出的是普通订阅的错误信息而非死信订阅的错误信息,会误导问题排查。
修复方案
  1. 给消息回调增加全链路异常捕获,保证所有分支明确返回ack/nack
    将原有订阅回调替换为如下实现,彻底避免消息悬而未决的状态:
auto session = subscriber.Subscribe([](pubsub::Message const& m, pubsub::AckHandler h) {
    // 投递次数超阈值直接nack,移交死信队列处理
    if (h.delivery_attempt().has_value() && h.delivery_attempt().value() >= 5) {
        LOGe << "Message exceed max delivery attempt, nack to DLQ, id: " << m.message_id();
        std::move(h).nack();
        return;
    }
    try {
        // 提前校验消息属性合法性,避免非法消息触发异常
        if (m.attributes().count("topic") == 0) {
            LOGe << "Invalid message without topic attribute, id: " << m.message_id();
            std::move(h).nack();
            return;
        }
        ProcessMsg::CommTopic topic = static_cast<ProcessMsg::CommTopic>(std::stoi(m.attributes().at("topic")));
        std::string msg = m.data();
        LOGi << "Received message, id " << m.message_id() << ", attempt " << h.delivery_attempt() <<", ordering key '" << m.ordering_key() << "', topic " << ProcessMsg::commTopic(topic) << " (" << topic << ")";
        messageCallback(msg);
        std::move(h).ack();
    } catch (std::exception const& e) {
        LOGe << "Process message failed, id: " << m.message_id() << ", err: " << e.what();
        // 处理失败明确nack,禁止消息长期占用配额
        std::move(h).nack();
    }
})
  1. 调整订阅配置避免单消息阻塞全流
    不要将max_outstanding_messages、max_concurrency设置为1,有序订阅场景下可以将该值调整为10~20,同一个ordering key的消息仍然会严格保证顺序不会乱序,但单个消息处理卡住不会完全堵死整个消费流。同时给messageCallback增加超时控制,业务处理超过30秒直接返回失败nack,避免长时间占用配额。
  2. 优化资源初始化逻辑
    Topic、Subscription、Publisher这类资源不需要在主发送循环中重复创建,程序启动时初始化一次即可长期复用,Publisher本身是线程安全的,复用既能提升性能也能避免多余的控制面调用触发限流。
  3. 修复笔误bug
    createOrderingSubscription中创建死信订阅dl_sub失败时,抛出的异常信息应该取dl_sub.status().message(),而非sub.status().message()。
  4. 增加健康检查兜底
    在主循环中记录最后一次收到消息的时间戳,如果连续2分钟没有收到任何消息、也没有触发错误回调,主动cancel当前订阅会话重新建立连接,实现故障自愈。
补充说明

有序订阅不支持恰好一次交付是官方明确的功能限制,当前C++客户端的有序订阅仅保证同ordering key的消息有序、至少一次交付,若需要恰好一次语义,需要在业务层增加消息幂等校验逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.31 23:18:16