第30章 async/await支持
Python 的 async/await 语法和 C++ 的异步编程模型有着本质区别。本章探讨如何将两者桥接,实现跨语言的异步互操作。
30.1 Python异步编程回顾
Section titled “30.1 Python异步编程回顾”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)。
30.2 C++异步桥接
Section titled “30.2 C++异步桥接”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_;};30.3 Future/Promise模型实现
Section titled “30.3 Future/Promise模型实现”使用 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 amimport 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 对象。
30.4 事件循环集成
Section titled “30.4 事件循环集成”将 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 asyncioimport 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 asyncioimport 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++ 中运行事件循环 |
最佳实践:
- C++ 中的异步操作必须在持有 GIL 的情况下调用 Python 对象
- 使用
py::gil_scoped_acquire/py::gil_scoped_release管理 GIL - 对于长时间运行的 C++ 操作,考虑使用线程池避免阻塞事件循环
- 使用
asyncio.to_thread从 Python 端调用同步 C++ 函数 - 避免在多个 Python 协程之间共享可变状态