Skip to content

第30章 async/await支持

Python 的 async/await 语法和 C++ 的异步编程模型有着本质区别。本章探讨如何将两者桥接,实现跨语言的异步互操作。

Python 的异步编程基于事件循环和协程机制。

import asyncio
async def fetch_data(url: str) -> dict:
"""模拟异步数据获取"""
await asyncio.sleep(1) # 模拟网络延迟
return {"url": url, "data": "sample"}
async def process_items(items: list[str]) -> list[dict]:
"""并发处理多个项目"""
tasks = [fetch_data(item) for item in items]
results = await asyncio.gather(*tasks)
return results
async def main():
"""主协程"""
print("Starting...")
data = await fetch_data("https://example.com")
print(f"Got: {data}")
# 并发执行
items = ["a", "b", "c"]
results = await process_items(items)
print(f"Processed: {len(results)} items")
asyncio.run(main())
async def async_generator(items):
for item in items:
await asyncio.sleep(0.1)
yield item
async def consume_async():
async for item in async_generator(["x", "y", "z"]):
print(f"Got: {item}")

关键洞察:Python 的 async/await 是基于生成器的语法糖,本质是协程。事件循环调度协程的挂起和恢复,而 await 等待的是一个可等待对象(Awaitable)。

C++ 端的异步操作需要与 Python 的事件循环协调。

#include <pybind11/pybind11.h>
#include <pybind11/functional.h>
#include <thread>
#include <future>
#include <chrono>
#include <mutex>
#include <queue>
#include <atomic>
namespace py = pybind11;
// 简单的异步任务包装器
class AsyncTask {
public:
AsyncTask() = default;
// 从 Python 端创建异步任务
template<typename Func>
void submit(Func&& func) {
auto future = std::async(std::launch::async, [func]() {
py::gil_scoped_acquire gil;
func();
});
futures_.push(std::move(future));
}
// 检查任务是否完成
bool is_complete() const {
if (futures_.empty()) return true;
return futures_.front().wait_for(std::chrono::seconds(0)) == std::future_status::ready;
}
// 获取结果(如果完成)
void wait() {
if (!futures_.empty()) {
futures_.front().get();
futures_.pop();
}
}
private:
std::queue<std::future<void>> futures_;
};
// C++ 后台任务执行器
class BackgroundExecutor {
public:
BackgroundExecutor() : stop_(false) {
worker_thread_ = std::thread([this]() {
while (!stop_) {
std::function<void()> task;
{
std::lock_guard<std::mutex> lock(mutex_);
if (!tasks_.empty()) {
task = std::move(tasks_.front());
tasks_.pop();
}
}
if (task) {
task();
} else {
std::this_thread::sleep_for(std::chrono::milliseconds(10));
}
}
});
}
~BackgroundExecutor() {
stop_ = true;
if (worker_thread_.joinable()) {
worker_thread_.join();
}
}
// 提交任务
void submit(py::function func) {
{
std::lock_guard<std::mutex> lock(mutex_);
tasks_.push([func]() {
py::gil_scoped_acquire gil;
func();
});
}
}
private:
std::atomic<bool> stop_;
std::thread worker_thread_;
std::mutex mutex_;
std::queue<std::function<void()>> tasks_;
};

使用 Future/Promise 在 Python 和 C++ 之间传递异步结果。

#include <pybind11/pybind11.h>
#include <pybind11/functional.h>
#include <thread>
#include <future>
#include <memory>
#include <string>
#include <vector>
namespace py = pybind11;
// Future 封装类,供 Python 使用
class PyFuture {
public:
PyFuture() = default;
template<typename T>
static std::shared_ptr<PyFuture> create(std::future<T> future) {
auto pf = std::make_shared<PyFuture>();
pf->future_ = std::make_any<std::future<T>>(std::move(future));
return pf;
}
bool done() const {
return check_future([](auto& f) {
return f.wait_for(std::chrono::seconds(0)) == std::future_status::ready;
});
}
void wait() {
visit_future([](auto& f) { f.get(); });
}
py::object get() {
py::gil_scoped_acquire gil;
py::object result = visit_future([](auto& f) {
using T = decltype(f.get())>;
if constexpr (std::is_same_v<T, int>) {
return py::int_(f.get());
} else if constexpr (std::is_same_v<T, std::string>) {
return py::str(f.get());
} else if constexpr (std::is_same_v<T, std::vector<int>>) {
return py::cast(f.get());
} else {
return py::none();
}
});
return result;
}
private:
template<typename Func>
bool check_future(Func&& func) const {
if (!future_.has_value()) return true;
bool result = false;
std::visit([&](auto& f) {
if constexpr (requires { func(f); }) {
result = func(f);
}
}, *future_);
return result;
}
template<typename Func>
auto visit_future(Func&& func) {
return std::visit([&](auto& f) {
return func(f);
}, *future_);
}
std::any future_;
};
// 异步计算器
class AsyncCalculator {
public:
// 异步加法
std::shared_ptr<PyFuture> add_async(int a, int b) {
auto future = std::async(std::launch::async, [a, b]() {
std::this_thread::sleep_for(std::chrono::milliseconds(100));
return a + b;
});
return PyFuture::create(std::move(future));
}
// 异步处理数据
std::shared_ptr<PyFuture> process_async(const std::vector<int>& data) {
auto future = std::async(std::launch::async, [data]() {
std::vector<int> result;
result.reserve(data.size());
for (int x : data) {
std::this_thread::sleep_for(std::chrono::milliseconds(50));
result.push_back(x * 2);
}
return result;
});
return PyFuture::create(std::move(future));
}
// 异步字符串操作
std::shared_ptr<PyFuture> string_ops_async(const std::string& s) {
auto future = std::async(std::launch::async, [s]() {
std::this_thread::sleep_for(std::chrono::milliseconds(200));
return s + "_processed";
});
return PyFuture::create(std::move(future));
}
};
// 链式异步任务
class AsyncChain {
public:
std::shared_ptr<PyFuture> then(int initial_value, py::function transform) {
auto future = std::async(std::launch::async, [initial_value, transform]() {
py::gil_scoped_acquire gil;
py::object result = transform(initial_value);
return result.cast<int>();
});
return PyFuture::create(std::move(future));
}
};
PYBIND11_MODULE(async_module, m) {
py::class_<PyFuture>(m, "PyFuture")
.def("done", &PyFuture::done)
.def("wait", &PyFuture::wait)
.def("get", &PyFuture::get);
py::class_<AsyncCalculator>(m, "AsyncCalculator")
.def(py::init<>())
.def("add_async", &AsyncCalculator::add_async)
.def("process_async", &AsyncCalculator::process_async)
.def("string_ops_async", &AsyncCalculator::string_ops_async);
py::class_<AsyncChain>(m, "AsyncChain")
.def(py::init<>())
.def("then", &AsyncChain::then);
}
import async_module as am
import time
def test_basic_async():
"""测试基本异步功能"""
calc = am.AsyncCalculator()
print("Starting async add...")
future = calc.add_async(10, 20)
print(f"Task submitted, done={future.done()}")
# 等待完成
while not future.done():
print("Waiting...")
time.sleep(0.05)
result = future.get()
print(f"Result: {result}")
def test_async_list():
"""测试异步列表处理"""
calc = am.AsyncCalculator()
print("Starting async list processing...")
future = calc.process_async([1, 2, 3, 4, 5])
# 阻塞获取
result = future.get()
print(f"Processed: {result}")
def test_async_string():
"""测试异步字符串操作"""
calc = am.AsyncCalculator()
future = calc.string_ops_async("hello")
result = future.get()
print(f"String result: {result}")

关键洞察:Future/Promise 模式允许 C++ 端的异步操作结果通过 Python 的 await 语法消费。关键是在 C++ 线程中正确获取 GIL,确保能安全地调用 Python 对象。

将 C++ 线程整合到 Python 的 asyncio 事件循环中。

#include <pybind11/pybind11.h>
#include <pybind11/asyncio.h>
#include <thread>
#include <future>
#include <chrono>
#include <atomic>
namespace py = pybind11;
// 运行器类:从 asyncio 获取事件循环
class AsyncRunner {
public:
AsyncRunner() : running_(false) {}
~AsyncRunner() {
stop();
}
// 启动后台线程,运行 asyncio 事件循环
void start() {
if (running_) return;
running_ = true;
worker_ = std::thread([this]() {
py::gil_scoped_acquire gil;
// 获取或创建事件循环
py::object loop = py::module_::import("asyncio").attr("new_event_loop")();
py::module_::import("asyncio").attr("set_event_loop")(loop);
loop.attr("run_forever")();
});
}
// 停止事件循环
void stop() {
if (!running_) return;
running_ = false;
if (worker_.joinable()) {
worker_.join();
}
}
// 在后台线程运行异步任务
template<typename Func>
py::object run_async(Func&& func) {
py::gil_scoped_acquire gil;
py::object loop = py::module_::import("asyncio").attr("get_event_loop")();
// 创建协程
py::object coroutine = py::cast(func);
return loop.attr("run_until_complete")(coroutine);
}
private:
std::atomic<bool> running_;
std::thread worker_;
};
// 在 Python 线程中运行 C++ 任务
py::object run_in_thread(py::function callback, py::function task) {
std::thread([callback, task]() {
// 在新线程中执行 C++ 任务
py::gil_scoped_acquire gil;
py::object result = task();
// 通过回调返回结果
callback(result);
}).detach();
return py::none();
}
// C++ 异步工作函数
std::future<int> cpp_async_work(int value) {
return std::async(std::launch::async, [value]() {
std::this_thread::sleep_for(std::chrono::milliseconds(500));
return value * 2;
});
}
PYBIND11_MODULE(asyncio_bridge, m) {
py::class_<AsyncRunner>(m, "AsyncRunner")
.def(py::init<>())
.def("start", &AsyncRunner::start)
.def("stop", &AsyncRunner::stop);
m.def("run_in_thread", &run_in_thread);
m.def("cpp_async_work", [](int value) {
auto future = cpp_async_work(value);
return future.get();
});
}
import asyncio
import asyncio_bridge as ab
async def cpp_task_wrapper():
"""包装 C++ 异步任务的协程"""
loop = asyncio.get_event_loop()
# 方法1:使用 asyncio.to_thread 运行同步 C++ 函数
result = await asyncio.to_thread(ab.cpp_async_work, 42)
print(f"C++ async result: {result}")
return result
async def test_runner():
"""测试异步运行器"""
runner = ab.AsyncRunner()
runner.start()
# 运行协程
loop = asyncio.get_event_loop()
coro = cpp_task_wrapper()
result = loop.run_until_complete(coro)
runner.stop()
print(f"Final result: {result}")
def test_callback():
"""测试回调模式"""
results = []
def on_result(result):
results.append(result)
print(f"Callback received: {result}")
# 在线程中运行
ab.run_in_thread(on_result, lambda: 123)
# 等待回调触发
import time
time.sleep(1)
print(f"Results: {results}")
asyncio.run(test_runner())
import asyncio
import asyncio_bridge as ab
class AsyncWorker:
"""C++ 异步工作的 Python 包装器"""
def __init__(self):
self.executor = ab.AsyncRunner()
self.executor.start()
def run_cpp_async(self, func, *args):
"""在 asyncio 中运行 C++ 异步函数"""
loop = asyncio.get_event_loop()
return loop.run_in_executor(None, lambda: func(*args))
async def process_batch(self, items):
"""批量处理"""
tasks = [self.run_cpp_async(ab.cpp_async_work, item) for item in items]
results = await asyncio.gather(*tasks)
return results
def shutdown(self):
self.executor.stop()
async def main():
worker = AsyncWorker()
# 单个任务
result = await worker.run_cpp_async(ab.cpp_async_work, 10)
print(f"Single result: {result}")
# 批量任务
results = await worker.process_batch([1, 2, 3, 4, 5])
print(f"Batch results: {results}")
worker.shutdown()
asyncio.run(main())

关键洞察:Python asyncio 事件循环是单线程的,但它能调度来自其他线程的任务。关键是使用 run_in_executor 或 asyncio.to_thread 在后台线程执行 C++ 阻塞操作,结果通过回调或 Future 返回给 asyncio。

async/await 支持总结:

模式适用场景实现方式
Future/Promise需要获取异步结果C++ future → Python 可等待对象
回调模式事件驱动的场景py::function 回调
线程池CPU 密集型 C++ 工作asyncio.to_thread
事件循环桥接深度 asyncio 集成在 C++ 中运行事件循环

最佳实践:

  1. C++ 中的异步操作必须在持有 GIL 的情况下调用 Python 对象
  2. 使用 py::gil_scoped_acquire / py::gil_scoped_release 管理 GIL
  3. 对于长时间运行的 C++ 操作,考虑使用线程池避免阻塞事件循环
  4. 使用 asyncio.to_thread 从 Python 端调用同步 C++ 函数
  5. 避免在多个 Python 协程之间共享可变状态