如何整合不兼容的线程阻塞式异步等待API?
无修改权限时兼容两个异步不兼容库的高效方案
问题背景
我们面临两个无法修改代码的异步库,它们的接口互不兼容,但需要在同一个线程中同时处理:
- 异步请求的完成事件
- 其他线程发来的消息事件
两个库的核心API如下:
请求库API
req_launch(...) // 异步发起请求 req_wait() // 阻塞线程直到任意请求完成 req_test() -> bool // 检查是否有请求完成(true为有,false为全部未完成) req_pop() -> resp_t // 获取已完成请求的响应数据
消息库API
msg_send(...) // 向其他线程发送消息 msg_wait() // 阻塞线程直到收到消息 msg_test() -> bool // 检查是否有消息可用 msg_pop() -> msg_t // 获取下一条消息
我们需要在发起请求的间隙,能及时接收消息并可能发起新请求,但直接用非阻塞API忙等的方式CPU占用过高,而修改库添加唤醒机制的方案不可行。
低效的忙等实现
最直接但低效的方式是轮询两个库的非阻塞接口:
# 线程A loop do sleep 5 msg_send(THREAD_B, MSG_A) end # 线程B loop do # 处理消息 if msg_test() case msg_pop() when MSG_A then req_launch(/* A相关数据 */) when MSG_B then req_launch(/* B相关数据 */) end end # 处理请求响应 if req_test() resp = req_pop() # 处理响应数据 end end
可修改库时的理想方案
如果能修改库代码,可以添加外部信号通道,让阻塞等待可被中断:
chan_create() -> chan_t // 创建可被信号触发的通道 chan_signal(chan_t) // 触发通道信号 chan_signalled(chan_t) -> bool // 检查通道是否被触发 req_wait(chan_t) // 阻塞直到请求完成或通道被触发
对应的实现代码:
# 线程A loop do sleep 5 msg_send(THREAD_B, MSG_A) chan_signal(THREAD_B.chan) end # 线程B chan = create_channel() num_waiting = 0 loop do # 无等待请求时,阻塞等消息 msg_wait() if num_waiting == 0 # 处理消息并发起请求 if msg_test() num_waiting += 1 case msg_pop() when MSG_A then req_launch(/* A相关数据 */) when MSG_B then req_launch(/* B相关数据 */) end end # 等待请求完成或被信号唤醒 req_wait(chan) unless chan_signalled(chan) num_waiting -= 1 resp = req_pop() # 处理响应数据 end end
不可修改库时的高效解决方案
1. POSIX/Linux下使用线程信号
利用POSIX线程的信号机制,给阻塞线程发送自定义信号(如SIGUSR1),中断其阻塞的系统调用(前提是库的req_wait/msg_wait基于可中断的系统调用实现,如select/poll)。
#include <signal.h> #include <pthread.h> // 信号处理函数(仅用于中断阻塞,无实际逻辑) void sig_handler(int sig) {} // 线程A逻辑:发送消息后给线程B发信号 void* thread_a(void* arg) { pthread_t thread_b_tid = *(pthread_t*)arg; while(1) { sleep(5); msg_send(THREAD_B, MSG_A); pthread_kill(thread_b_tid, SIGUSR1); } } // 线程B逻辑 void* thread_b(void* arg) { // 注册信号处理 signal(SIGUSR1, sig_handler); int num_waiting = 0; while(1) { // 处理所有待处理消息 while(msg_test()) { msg_t msg = msg_pop(); switch(msg) { case MSG_A: req_launch(...); num_waiting++; break; case MSG_B: req_launch(...); num_waiting++; break; } } // 处理已完成的请求 while(req_test()) { resp_t resp = req_pop(); num_waiting--; // 处理响应数据 } // 阻塞等待:有请求时等请求完成,否则等消息 sigset_t mask; sigemptyset(&mask); sigaddset(&mask, SIGUSR1); if(num_waiting > 0) { // 用sigsuspend等待,会被信号中断 sigsuspend(&mask); } else { // 无请求时阻塞等消息,同样会被信号中断 msg_wait(); } } }
2. 跨平台管道唤醒方案
创建一个管道作为通用唤醒源,将阻塞的req_wait/msg_wait放到单独线程中,当事件发生时往管道写入数据,主线程通过监听管道的可读事件来统一处理所有事件。
#include <unistd.h> #include <pthread.h> int wakeup_pipe[2]; // 请求监控线程:请求完成时唤醒主线程 void* req_monitor(void* arg) { while(1) { req_wait(); char c = 0; write(wakeup_pipe[1], &c, 1); } } // 消息监控线程:消息到达时唤醒主线程 void* msg_monitor(void* arg) { while(1) { msg_wait(); char c = 0; write(wakeup_pipe[1], &c, 1); } } // 主线程逻辑 int main() { pipe(wakeup_pipe); pthread_t req_tid, msg_tid; pthread_create(&req_tid, NULL, req_monitor, NULL); pthread_create(&msg_tid, NULL, msg_monitor, NULL); char buf[1]; while(1) { // 等待管道唤醒 read(wakeup_pipe[0], buf, 1); // 处理所有待处理消息 while(msg_test()) { msg_t msg = msg_pop(); switch(msg) { case MSG_A: req_launch(...); break; case MSG_B: req_launch(...); break; } } // 处理所有已完成请求 while(req_test()) { resp_t resp = req_pop(); // 处理响应数据 } } return 0; }
3. 通用语言的线程包装方案
在Python、Java等语言中,用独立线程包装阻塞的等待函数,通过队列或条件变量通知主线程处理事件。
以Python为例:
import threading import queue # 用于传递唤醒事件的队列 wakeup_queue = queue.Queue() def req_monitor(): """监控请求完成事件""" while True: req_wait() wakeup_queue.put("req_done") def msg_monitor(): """监控消息到达事件""" while True: msg_wait() wakeup_queue.put("msg_arrived") # 启动后台监控线程 threading.Thread(target=req_monitor, daemon=True).start() threading.Thread(target=msg_monitor, daemon=True).start() # 主线程处理逻辑 while True: event = wakeup_queue.get() if event == "msg_arrived": # 处理所有待处理消息 while msg_test(): msg = msg_pop() if msg == "MSG_A": req_launch(...) elif msg == "MSG_B": req_launch(...) elif event == "req_done": # 处理所有已完成请求 while req_test(): resp = req_pop() # 处理响应数据
4. 定时器轮询(折中方案)
如果以上方案都无法实现,可采用带超时的轮询,平衡CPU占用和响应延迟:
#include <time.h> int main() { while(1) { // 处理待处理消息 while(msg_test()) { msg_t msg = msg_pop(); switch(msg) { case MSG_A: req_launch(...); break; case MSG_B: req_launch(...); break; } } // 处理已完成请求 while(req_test()) { resp_t resp = req_pop(); // 处理响应数据 } // 休眠100ms,避免忙等 struct timespec ts = {0, 100000000}; // 100毫秒 nanosleep(&ts, NULL); } return 0; }
内容的提问来源于stack exchange,提问作者Post Self
相关产品推荐
相关产品推荐

