具身机器人视觉服务的异步回调与并发控制

写这篇文章,起因挺简单:团队内部做线程调用一直没个统一模板。每个人写多线程异步都是各写各的,有人裸起 std::thread,有人塞 std::async,回调传参、标志位复位更是五花八门,review 起来费劲,出了问题也不好排查。

于是把这套「异步发起 + 回调收尾 + CAS 无锁抢占」的模式沉淀下来,当团队内部的脚手架。往后具身工程化里,凡是「发起一个异步任务、处理完回传结果、同时只允许一个在跑」的场景,直接套这个模板就行。

简述

这个例子把三件事串在一条链路上:std::async 在后台线程异步发起请求,主线程不阻塞;客户端处理完通过 std::function 回调把结果送回服务端;std::atomic<bool> 配合 compare_exchange_strong 做无锁抢占,挡住重复请求。

代码里用忙等循环故意放大竞态窗口,main 里再不 sleep 连发 500 次请求,把并发保护的效果直观跑出来。

项目结构

multi_thread_callback/
├── main.cpp               # 入口:一次单发 + 500 次快速连发
├── obj_detect_server.hpp  # 服务端:CAS 抢占 + 异步请求 + 回调收尾
└── ultinet_client.hpp     # 客户端:处理数据并调用回调

核心机制

异步发起

std::async(std::launch::async, ...) 在独立线程里跑 ProcessData,主线程只负责把请求丢出去,不等结果。

回调收尾

客户端处理完,调 std::function<void(const std::string&)> 回到服务端的 OnDetectionDone:打印结果、存 lastResult_、复位标志。

无锁抢占

std::atomic<bool> requestInProgress_ 配合 compare_exchange_strong,只有能把 false 置成 true 的线程才继续往下走,抢失败的直接进 DetectionBusy() 返回。「检查 + 置位」是原子的,不会两个线程同时挤进处理流程。

代码

UltiNetClient

模拟客户端,处理数据,完成后调回调:

class UltiNetClient {
public:
    using Callback = std::function<void(const std::string&)>;

    void ProcessData(const std::string& request, Callback callback) {
        std::string result = "processed: " + request;
        lastRequest_ = request + " (processed)";
        for (volatile int i = 0; i < 1000000; ++i) ; // 模拟耗时
        callback(result);
    }

    void getLastRequest(std::string& request) const {
        request = lastRequest_;
    }

private:
    std::string lastRequest_;
};

那段 100 万次的忙等循环,就是拿来模拟处理耗时的,真实代码别这么写。

ObjetDetectServer

服务端,把「抢占 → 异步发起 → 回调收尾」串起来:

class ObjetDetectServer {
public:
    void Detect(const std::string& request) {
        if (request.empty()) {
            std::cout << "[ObjetDetectServer] empty request" << std::endl;
            lastResult_ = "empty request";
            return;
        }

        bool expected = false;
        if (!requestInProgress_.compare_exchange_strong(expected, true)) {
            DetectionBusy();
            return;
        } else {
            for (volatile int i = 0; i < 10000000; ++i) ; // 放大竞态窗口

            future_ = std::async(std::launch::async, [this, request]() {
                client_.ProcessData(request, [this](const std::string& result) {
                    this->OnDetectionDone(result);
                });
            });
        }
    }

    void OnDetectionDone(const std::string& result) {
        std::cout << "[ObjetDetectServer] done: " << result << std::endl;
        client_.getLastRequest(lastResult_);
        std::cout << "[ObjetDetectServer] last request: " << lastResult_ << std::endl;
        requestInProgress_.store(false);
    }

    void DetectionBusy() {
        std::cout << "[ObjetDetectServer] busy, please wait..." << std::endl;
    }

private:
    UltiNetClient client_;
    std::string lastResult_;
    std::atomic<bool> requestInProgress_{false};
    std::future<void> future_;
};

两个点说一下:

  1. CAS 抢锁成功后那段 1000 万次忙等是故意留的,放大临界区,给别的线程制造穿插执行的机会,用来验证原子标志到底靠不靠谱。
  2. future_ 是成员变量,持有 std::async 返回的 future,保证异步任务在服务端生命周期内跑完。

main

int main() {
    ObjetDetectServer server;

    server.Detect("person,car,dog");   // 第一次请求,异步处理

    for (int i = 0; i < 500; ++i) {
        server.Detect("items,objects,packages");  // 快速连发
    }
}

main 先发一次 "person,car,dog",然后不 sleep 连发 500 次。第一次请求还在处理中(requestInProgress_ 已置位),后面绝大多数都会命中 busy 分支。

构建运行

g++ -std=c++11 -pthread main.cpp -o main
./main

实测单次输出 502 行,结构是:最后 2 行是 done + last request,前面 500 行还包含几次items,objects,packages

[ObjetDetectServer] busy, please wait...
[ObjetDetectServer] busy, please wait...
...(共 500 行 busy)
[ObjetDetectServer] done: processed: person,car,dog
[ObjetDetectServer] last request: person,car,dog (processed)

在这里插入图片描述

注:具体顺序随调度略有浮动。如果先抢到锁的线程在 async 启动前被调度走,done 会排到所有 busy 之后。

线程安全说明

  • 已保护:requestInProgress_ 的「检查 + 置位」是原子的,不会两个线程同时进入处理流程;工作线程结束时用 store(false) 复位。
  • 仍有竞争:lastResult_ / lastRequest_ 由后台线程写,如果 LastResult() / getLastRequest() 被别的线程读,就有数据竞争。demo 里主线程只调 Detect() 没事,但公共接口一旦多线程用,得上 std::mutex,或者把共享结果换成原子共享指针。

实验提示

  • 想复现无保护时的错误行为:把 requestInProgress_ 改回普通 bool(保留忙等窗口),就能看到竞态。
  • 想少点 busy 提示:去掉或缩短忙等窗口,或者主循环里加 sleep 拉开间隔。
  • 生产建议:用 std::mutex / std::lock_guard(或 std::atomic_flag)替代「抢锁 + 延迟复位」的模式,并对共享结果加锁保护。

参考

  1. std::async - cppreference
  2. std::atomic - cppreference
Logo

DAMO开发者矩阵,由阿里巴巴达摩院和中国互联网协会联合发起,致力于探讨最前沿的技术趋势与应用成果,搭建高质量的交流与分享平台,推动技术创新与产业应用链接,围绕“人工智能与新型计算”构建开放共享的开发者生态。

更多推荐