尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

C语言发布订阅模式:从轮询到事件驱动的轻量级实现

C语言发布订阅模式:从轮询到事件驱动的轻量级实现 1. 从“轮询”到“通知”为什么我们需要发布订阅模式在嵌入式或者系统级的C语言开发里我们经常要处理一个头疼的问题一个模块的状态变了怎么让其他几个关心这个状态的模块知道新手最常见的做法就是“轮询”。比如一个温度采集模块temperature_sensor一个显示模块display一个报警模块alarm。display和alarm都想知道温度值。于是代码可能就写成了这样// 伪代码示例轮询模式 float current_temperature 0.0f; void temperature_sensor_task(void) { // 读取传感器更新全局变量 current_temperature read_sensor(); } void display_task(void) { while(1) { // 不停地读取全局变量来刷新显示 update_display(current_temperature); sleep(1000); // 每秒刷新一次 } } void alarm_task(void) { while(1) { // 也不停地读取全局变量来判断是否报警 if (current_temperature THRESHOLD) { trigger_alarm(); } sleep(500); // 每500毫秒检查一次 } }这种模式的问题太明显了。首先资源浪费。display_task和alarm_task大部分时间都在空转sleep的时间间隔成了性能和实时性的矛盾体间隔短了CPU空耗严重间隔长了显示和报警的延迟又让人无法接受。其次耦合度高。display和alarm模块必须知道current_temperature这个全局变量的存在和名字一旦变量名改变或者存储方式改变比如从float改成int所有依赖它的模块都得改。最后可扩展性差。如果现在要加一个数据记录模块logger也得来轮询这个全局变量代码的修改是侵入式的。发布订阅模式就是为了解决这些问题而生的。它的核心思想是解耦和事件驱动。发布者Publisher不关心谁需要数据它只负责在数据变化时“发布”一个事件或消息。订阅者Subscriber不关心数据从哪里来它只声明自己关心某类事件并在事件发生时被“通知”到。中间由一个“经纪人”Broker或“事件通道”Event Channel来负责管理订阅关系和消息路由。在C语言这种没有原生语言级事件机制的环境里实现发布订阅模式本质上就是在构建一个轻量级、高效的消息中间件。这对于嵌入式系统、游戏引擎、网络服务器等对性能和资源有严格要求的场景是至关重要的基础设施。2. 模式核心三要素与两种实现模型发布订阅模式听起来高级拆开来看就是三个核心角色理解了它们就理解了模式的骨架。发布者Publisher/Producer它是事件的源头。当某个状态发生改变、一个操作完成、或者一个特定条件被触发时发布者就创建一个“消息”或称为“事件”、“主题”并将这个消息发送出去。在C语言里这个消息通常是一个结构体struct里面包含了事件类型标识和相关的数据载荷。发布者的责任很单一构造消息然后调用“发送”接口。它完全不知道也不应该知道这条消息最终会被谁处理。订阅者Subscriber/Consumer它是事件的处理端。订阅者会向系统注册告诉系统“我对‘XXX类型’的事件感兴趣”。当这类事件发生时系统会回调订阅者预先注册好的一个函数回调函数并把事件消息传递给它。这个回调函数就是事件的处理逻辑。订阅者只关心自己订阅的事件不关心是谁发布的。消息代理Broker/Event Bus/Dispatcher这是模式的大脑也是最需要我们动手实现的部分。它维护着一个“订阅表”通常是一个映射关系事件类型 - 回调函数列表。它的工作流程是1接收订阅者的注册将回调函数添加到对应事件类型的列表中2接收发布者的消息3根据消息的类型查找订阅表找到所有注册了的回调函数4依次调用这些回调函数并将消息传递给它们。在C语言中实现这个代理主要有两种模型选择哪一种取决于你的应用场景和对实时性、复杂度的要求。2.1 同步调用模型这是最简单、最直接的实现。当发布者调用publish(event)时这个函数会直接在当前线程的上下文中同步地、依次调用所有订阅了该事件的回调函数。就像函数调用一样。typedef void (*event_handler_t)(void* event_data); typedef struct { int event_type; event_handler_t handlers[MAX_HANDLERS]; int handler_count; } subscription_t; // 发布函数同步 void publish_event_sync(int event_type, void* event_data) { subscription_t* sub find_subscription(event_type); if (sub) { for (int i 0; i sub-handler_count; i) { // 同步调用回调函数 sub-handlers[i](event_data); } } }优点实现极其简单没有线程切换开销执行顺序确定就是注册的顺序。缺点最大的问题是阻塞。如果某个订阅者的回调函数执行时间很长或者发生阻塞比如调用了sleep、等待一个锁那么发布者线程会被卡住后续的订阅者也会被延迟通知整个系统的响应性会变差。它适用于回调函数都非常短小、快速、无阻塞的场景。2.2 异步队列模型这是更健壮、更符合“解耦”精神的实现。发布者不直接调用回调函数而是将事件消息放入一个队列通常是线程安全的环形缓冲区或链表。系统有一个或多个独立的“消费者线程”或“事件循环”专门从这个队列里取出消息然后查找订阅表并执行回调。// 伪代码异步模型核心 typedef struct { int event_type; void* event_data; } event_msg_t; event_msg_t event_queue[MAX_QUEUE_SIZE]; // 环形队列 pthread_mutex_t queue_mutex; sem_t queue_sem; // 用于消费者线程等待 // 发布函数异步 void publish_event_async(int event_type, void* event_data) { event_msg_t msg {event_type, event_data}; pthread_mutex_lock(queue_mutex); enqueue(event_queue, msg); // 消息入队 pthread_mutex_unlock(queue_mutex); sem_post(queue_sem); // 通知消费者线程 } // 独立的消费者线程函数 void* event_consumer_thread(void* arg) { while(1) { sem_wait(queue_sem); // 等待有消息 pthread_mutex_lock(queue_mutex); event_msg_t msg dequeue(event_queue); // 出队 pthread_mutex_unlock(queue_mutex); // 找到订阅者并同步调用回调在消费者线程上下文中 subscription_t* sub find_subscription(msg.event_type); if (sub) { for (int i 0; i sub-handler_count; i) { sub-handlers[i](msg.event_data); } } // 注意这里需要处理event_data的内存释放如果它是动态分配的 } return NULL; }优点彻底解耦了发布者和订阅者的执行流。发布者只需快速投递消息然后立刻返回不会因为订阅者的慢处理而阻塞。提高了系统的整体吞吐量和响应性。缺点实现复杂引入了线程、锁、队列等并发编程要素调试难度增加。事件的处理变成了异步的顺序可能因线程调度而变得不确定如果需要严格顺序需要在队列或消费者端做额外控制。同时动态分配的消息内存需要在合适的时机释放否则会导致内存泄漏。注意在资源极其受限的裸机无操作系统嵌入式环境中可能没有线程和队列库。此时异步模型可以通过“主循环状态机”来模拟。在主循环中检查事件标志位如果有事件 pending则调用对应的处理函数。这要求发布者只是“置位”一个标志而不是直接调用函数。3. 手把手实现一个轻量级、线程安全的C语言事件总线理论说再多不如一行代码。下面我们来设计并实现一个功能完整、可用于实际项目的轻量级事件总线EventBus。我们将采用异步队列模型并考虑线程安全、内存管理和易用性。3.1 数据结构设计首先定义核心的数据结构。我们需要一个结构来表示事件本身一个结构来管理一种事件类型的所有订阅者以及一个全局的事件总线上下文。// event_bus.h #ifndef EVENT_BUS_H #define EVENT_BUS_H #include stdbool.h #include stdint.h // 事件基类型所有具体事件结构体的第一个成员必须是它 typedef struct { uint32_t id; // 事件唯一ID可用于区分同一类型下的不同子事件 int type; // 事件类型如 TEMPERATURE_CHANGED, BUTTON_PRESSED void* sender; // 事件发布者可选 } event_t; // 事件处理回调函数类型 typedef void (*event_handler_fp)(const event_t* evt); // 事件总线句柄不透明类型隐藏内部实现细节 typedef struct event_bus_ctx_t event_bus_t; // 创建和销毁事件总线 event_bus_t* event_bus_create(void); void event_bus_destroy(event_bus_t* bus); // 订阅与取消订阅 bool event_bus_subscribe(event_bus_t* bus, int event_type, event_handler_fp handler); bool event_bus_unsubscribe(event_bus_t* bus, int event_type, event_handler_fp handler); // 发布事件同步和异步两个版本 void event_bus_publish(event_bus_t* bus, const event_t* evt); // 异步入队即返回 void event_bus_publish_sync(event_bus_t* bus, const event_t* evt); // 同步阻塞直到所有handler执行完 // 处理事件队列需要在某个线程中循环调用或由内部线程自动处理 void event_bus_process(event_bus_t* bus, int timeout_ms); // 处理单个事件可设置超时 bool event_bus_start_background_thread(event_bus_t* bus); // 启动内部消费者线程高级功能 #endif // EVENT_BUS_H接下来是具体的实现文件event_bus.c。我们会用到POSIX线程pthread和信号量semaphore来实现跨平台的线程同步。如果你的平台不支持如某些RTOS需要替换为相应的原语如互斥量、消息队列。// event_bus.c #include event_bus.h #include stdlib.h #include string.h #include pthread.h #include semaphore.h #define MAX_EVENT_TYPES 64 // 支持的最大事件类型数 #define MAX_SUBSCRIBERS_PER_TYPE 10 // 每种事件最多订阅者数 #define EVENT_QUEUE_SIZE 256 // 事件队列大小 // 订阅者节点 typedef struct subscriber_node_t { event_handler_fp handler; struct subscriber_node_t* next; } subscriber_node_t; // 事件类型订阅表 typedef struct { int type; subscriber_node_t* head; // 订阅者链表头 pthread_mutex_t mutex; // 保护此类型订阅链表的锁 } subscription_table_t; // 事件总线上下文实现细节 struct event_bus_ctx_t { // 1. 订阅表 subscription_table_t subscriptions[MAX_EVENT_TYPES]; int sub_table_count; // 2. 事件队列环形缓冲区 event_t* event_queue[EVENT_QUEUE_SIZE]; int queue_head; int queue_tail; pthread_mutex_t queue_mutex; sem_t queue_sem; // 计数信号量表示队列中可消费的事件数 // 3. 后台线程控制 pthread_t consumer_thread; bool thread_running; };3.2 核心函数实现初始化与销毁创建事件总线时我们需要初始化所有内部资源订阅表、队列、互斥锁和信号量。event_bus_t* event_bus_create(void) { event_bus_t* bus (event_bus_t*)calloc(1, sizeof(event_bus_t)); if (!bus) return NULL; // 初始化订阅表 bus-sub_table_count 0; for (int i 0; i MAX_EVENT_TYPES; i) { bus-subscriptions[i].head NULL; pthread_mutex_init(bus-subscriptions[i].mutex, NULL); } // 初始化队列 bus-queue_head 0; bus-queue_tail 0; pthread_mutex_init(bus-queue_mutex, NULL); sem_init(bus-queue_sem, 0, 0); // 初始值为0表示队列为空 // 初始化后台线程标志 bus-thread_running false; return bus; } void event_bus_destroy(event_bus_t* bus) { if (!bus) return; // 停止后台线程如果启动了 bus-thread_running false; sem_post(bus-queue_sem); // 唤醒可能阻塞的线程使其退出 pthread_join(bus-consumer_thread, NULL); // 释放所有订阅者节点内存 for (int i 0; i bus-sub_table_count; i) { subscriber_node_t* node bus-subscriptions[i].head; while (node) { subscriber_node_t* next node-next; free(node); node next; } pthread_mutex_destroy(bus-subscriptions[i].mutex); } // 释放队列中可能残留的事件内存注意这要求事件是动态分配的 pthread_mutex_lock(bus-queue_mutex); while (bus-queue_head ! bus-queue_tail) { free(bus-event_queue[bus-queue_head]); bus-queue_head (bus-queue_head 1) % EVENT_QUEUE_SIZE; } pthread_mutex_unlock(bus-queue_mutex); // 销毁同步原语 pthread_mutex_destroy(bus-queue_mutex); sem_destroy(bus-queue_sem); // 释放总线自身 free(bus); }3.3 核心函数实现订阅与发布订阅操作就是将回调函数添加到对应事件类型的链表中。这里需要注意线程安全因为订阅/取消订阅可能发生在任何线程。bool event_bus_subscribe(event_bus_t* bus, int event_type, event_handler_fp handler) { if (!bus || !handler) return false; // 查找是否已存在该事件类型的订阅表项 subscription_table_t* sub_table NULL; int index -1; for (int i 0; i bus-sub_table_count; i) { if (bus-subscriptions[i].type event_type) { sub_table bus-subscriptions[i]; index i; break; } } // 如果不存在创建一个新的表项需考虑数组上限 if (!sub_table) { if (bus-sub_table_count MAX_EVENT_TYPES) { return false; // 事件类型数已达上限 } index bus-sub_table_count; sub_table bus-subscriptions[index]; sub_table-type event_type; sub_table-head NULL; } // 检查是否已重复订阅可选但建议做 pthread_mutex_lock(sub_table-mutex); subscriber_node_t* curr sub_table-head; while (curr) { if (curr-handler handler) { pthread_mutex_unlock(sub_table-mutex); return false; // 重复订阅 } curr curr-next; } // 创建新的订阅者节点并插入链表头部头部插入效率高 subscriber_node_t* new_node (subscriber_node_t*)malloc(sizeof(subscriber_node_t)); if (!new_node) { pthread_mutex_unlock(sub_table-mutex); return false; } new_node-handler handler; new_node-next sub_table-head; sub_table-head new_node; pthread_mutex_unlock(sub_table-mutex); return true; }发布异步事件核心是将事件拷贝一份深拷贝放入队列。这里有一个关键设计决策事件数据的生命周期管理。我们选择在发布时动态分配内存来拷贝事件在处理完毕后由消费者线程释放。这确保了发布者可以传递栈上的临时变量也避免了发布者等待订阅者处理完再释放数据的耦合。static bool enqueue_event(event_bus_t* bus, const event_t* evt) { pthread_mutex_lock(bus-queue_mutex); int next_tail (bus-queue_tail 1) % EVENT_QUEUE_SIZE; if (next_tail bus-queue_head) { // 队列满 pthread_mutex_unlock(bus-queue_mutex); return false; } // 深拷贝事件。注意这里假设event_t结构体是“平坦”的没有内部指针需要二次拷贝。 // 如果事件数据包含指针需要更复杂的拷贝函数如dup_event。 event_t* evt_copy (event_t*)malloc(sizeof(event_t)); if (!evt_copy) { pthread_mutex_unlock(bus-queue_mutex); return false; } memcpy(evt_copy, evt, sizeof(event_t)); bus-event_queue[bus-queue_tail] evt_copy; bus-queue_tail next_tail; pthread_mutex_unlock(bus-queue_mutex); sem_post(bus-queue_sem); // 增加信号量计数通知消费者 return true; } void event_bus_publish(event_bus_t* bus, const event_t* evt) { if (!bus || !evt) return; if (!enqueue_event(bus, evt)) { // 处理队列满的情况可以打印日志、丢弃事件、或者阻塞等待不推荐在发布者线程阻塞 // 这里简单丢弃实际项目应根据需求处理。 } }3.4 核心函数实现事件处理与消费者线程事件处理函数event_bus_process可以从主线程中周期性调用适合无操作系统或主循环架构也可以由一个独立的后台线程调用。void event_bus_process(event_bus_t* bus, int timeout_ms) { if (!bus) return; // 等待队列中有事件带超时。如果timeout_ms为0则非阻塞检查。 struct timespec ts; if (timeout_ms 0) { clock_gettime(CLOCK_REALTIME, ts); ts.tv_nsec (timeout_ms % 1000) * 1000000; ts.tv_sec timeout_ms / 1000 ts.tv_nsec / 1000000000; ts.tv_nsec % 1000000000; if (sem_timedwait(bus-queue_sem, ts) ! 0) { return; // 超时或无事件 } } else { if (sem_trywait(bus-queue_sem) ! 0) { return; // 无事件 } } // 从队列中取出事件 pthread_mutex_lock(bus-queue_mutex); if (bus-queue_head bus-queue_tail) { // 理论上不会发生因为信号量计数保证了有事件但出于安全考虑检查 pthread_mutex_unlock(bus-queue_mutex); sem_post(bus-queue_sem); // 补偿信号量 return; } event_t* evt bus-event_queue[bus-queue_head]; bus-queue_head (bus-queue_head 1) % EVENT_QUEUE_SIZE; pthread_mutex_unlock(bus-queue_mutex); // 查找订阅者并调用回调 subscription_table_t* sub_table NULL; for (int i 0; i bus-sub_table_count; i) { if (bus-subscriptions[i].type evt-type) { sub_table bus-subscriptions[i]; break; } } if (sub_table) { pthread_mutex_lock(sub_table-mutex); subscriber_node_t* node sub_table-head; // 在调用回调前先复制链表头因为回调函数内部可能会执行取消订阅操作修改链表。 // 更稳健的做法是在锁内将链表复制到本地数组然后解锁再调用回调。 subscriber_node_t* handlers_to_call[10]; int handler_count 0; while (node handler_count 10) { handlers_to_call[handler_count] node; node node-next; } pthread_mutex_unlock(sub_table-mutex); // 在锁外执行回调避免死锁回调函数里可能又发布新事件 for (int i 0; i handler_count; i) { handlers_to_call[i]-handler(evt); } } // 处理完毕释放事件内存 free(evt); }启动一个独立的后台消费者线程可以让事件处理完全自动化static void* event_consumer_thread_func(void* arg) { event_bus_t* bus (event_bus_t*)arg; while (bus-thread_running) { event_bus_process(bus, -1); // -1 表示无限等待直到有事件 } return NULL; } bool event_bus_start_background_thread(event_bus_t* bus) { if (!bus || bus-thread_running) return false; bus-thread_running true; if (pthread_create(bus-consumer_thread, NULL, event_consumer_thread_func, bus) ! 0) { bus-thread_running false; return false; } return true; }4. 实战演练用事件总线重构温度监控系统现在让我们用上面实现的事件总线来重构文章开头那个笨拙的“全局变量轮询”温度监控系统。我们将看到代码如何变得更清晰、更解耦、更易于扩展。首先定义我们关心的事件类型和具体的事件数据结构。这是良好设计的第一步用枚举和结构体明确通信协议。// system_events.h #ifndef SYSTEM_EVENTS_H #define SYSTEM_EVENTS_H #include event_bus.h // 包含我们的事件总线头文件 // 定义系统内所有的事件类型 typedef enum { EVENT_TEMPERATURE_CHANGED 1, // 温度变化事件 EVENT_BUTTON_PRESSED, // 按键按下事件 EVENT_SYSTEM_ERROR, // 系统错误事件 // ... 其他事件 } system_event_type_t; // 温度变化事件的具体数据结构 // 注意第一个成员必须是 event_t typedef struct { event_t base; // 基类必须放在第一个 float temperature_c; // 摄氏度 uint32_t timestamp_ms; // 时间戳 } temperature_event_t; // 按键事件数据结构 typedef struct { event_t base; int button_id; bool is_long_press; } button_event_t; // 系统错误事件数据结构 typedef struct { event_t base; int error_code; const char* error_msg; } system_error_event_t; // 便捷函数创建并初始化一个温度事件动态分配 temperature_event_t* temperature_event_create(float temp_c, uint32_t ts); #endif // SYSTEM_EVENTS_H// system_events.c #include system_events.h #include stdlib.h temperature_event_t* temperature_event_create(float temp_c, uint32_t ts) { temperature_event_t* evt (temperature_event_t*)malloc(sizeof(temperature_event_t)); if (evt) { evt-base.type EVENT_TEMPERATURE_CHANGED; evt-base.id 0; // 可以按需生成唯一ID evt-base.sender NULL; evt-temperature_c temp_c; evt-timestamp_ms ts; } return evt; }接下来实现各个模块。首先是温度传感器模块它作为发布者职责变得非常单一。// temperature_sensor.c #include system_events.h #include event_bus.h // 假设有一个全局的事件总线实例可通过依赖注入获得这里简化 extern event_bus_t* g_event_bus; void temperature_sensor_polling_task(void) { static float last_temperature 0.0f; float current_temperature read_hardware_sensor(); // 假设的硬件读取函数 // 只有温度变化超过阈值时才发布事件避免频繁发布 if (fabs(current_temperature - last_temperature) 0.1f) { last_temperature current_temperature; // 创建事件对象动态分配 temperature_event_t* temp_evt temperature_event_create( current_temperature, get_system_tick() // 假设的系统滴答时钟 ); if (temp_evt) { // 发布事件传感器模块完全不知道谁会处理它。 event_bus_publish(g_event_bus, (event_t*)temp_evt); // 注意事件内存由事件总线的消费者线程负责释放这里不需要free。 } } }然后是显示模块和报警模块它们作为订阅者只需要注册自己对温度事件的兴趣。// display_module.c #include system_events.h #include event_bus.h extern event_bus_t* g_event_bus; // 显示模块的事件处理回调函数 static void on_temperature_changed(const event_t* evt) { // 首先进行安全类型转换 const temperature_event_t* temp_evt (const temperature_event_t*)evt; // 更新显示 update_display(temp_evt-temperature_c); } // 显示模块初始化函数 void display_module_init(void) { // 订阅温度变化事件 event_bus_subscribe(g_event_bus, EVENT_TEMPERATURE_CHANGED, on_temperature_changed); }// alarm_module.c #include system_events.h #include event_bus.h extern event_bus_t* g_event_bus; #define TEMP_THRESHOLD 85.0f // 报警模块的事件处理回调函数 static void on_temperature_changed_alarm(const event_t* evt) { const temperature_event_t* temp_evt (const temperature_event_t*)evt; if (temp_evt-temperature_c TEMP_THRESHOLD) { trigger_alarm(); } } void alarm_module_init(void) { event_bus_subscribe(g_event_bus, EVENT_TEMPERATURE_CHANGED, on_temperature_changed_alarm); }最后在系统主函数中我们初始化所有组件并启动事件处理循环。// main.c #include event_bus.h #include system_events.h #include temperature_sensor.h #include display_module.h #include alarm_module.h #include unistd.h // for sleep event_bus_t* g_event_bus NULL; int main() { // 1. 创建事件总线 g_event_bus event_bus_create(); if (!g_event_bus) { // 处理创建失败 return -1; } // 2. 启动后台消费者线程异步模式 if (!event_bus_start_background_thread(g_event_bus)) { // 处理线程启动失败可以回退到在主循环中调用process event_bus_destroy(g_event_bus); return -1; } // 3. 初始化各功能模块内部会完成事件订阅 display_module_init(); alarm_module_init(); // ... 初始化其他模块如按键模块等 // 4. 主循环模拟传感器轮询 while (1) { temperature_sensor_polling_task(); // 这里会发布事件 // 其他任务... usleep(100 * 1000); // 休眠100ms } // 5. 清理实际中应有退出机制 event_bus_destroy(g_event_bus); return 0; }对比与优势解耦temperature_sensor不再需要知道display和alarm的存在。它只和event_bus交互。高效display和alarm只在温度真正变化时被调用没有轮询开销。易扩展现在要增加一个数据记录模块logger只需要新建一个logger.c在其中实现一个回调函数并订阅EVENT_TEMPERATURE_CHANGED事件即可。无需修改temperature_sensor、display或alarm的任何一行代码。灵活性如果将来显示模块需要同时响应温度变化和按键事件它只需再订阅一个EVENT_BUTTON_PRESSED即可。5. 进阶话题与避坑指南实现一个能用的发布订阅框架不难但要让它稳定、高效、安全地运行在真实项目中还需要考虑很多边界情况和进阶特性。5.1 内存管理谁分配谁释放这是C语言实现事件总线最棘手的问题之一。在我们的实现中采用了“发布者分配消费者释放”的策略。发布者调用temperature_event_create内部malloc创建事件通过event_bus_publish传递指针。总线负责将指针入队最终由event_bus_process函数在调用完所有回调后free掉。潜在坑点回调中再次发布事件如果某个订阅者的回调函数里又发布了新事件并且可能也动态分配了内存这没有问题只要新事件走同样的流程即可。事件数据包含指针如果temperature_event_t里包含一个char* error_msg那么在malloc事件结构体后还需要为这个字符串额外分配内存并拷贝。同样在free事件结构体前需要先free这个字符串。这需要更复杂的拷贝和释放函数如dup_event和free_event最好通过函数指针让用户注册。替代方案对于极度追求性能或禁止动态内存的系统可以使用静态事件池。预先分配一个固定大小的数组作为事件对象池发布时从池中取一个空闲事件填充数据使用后标记为空闲。这避免了malloc/free的开销和碎片但增加了管理的复杂性并且有池大小的限制。5.2 线程安全与死锁预防我们的实现已经对共享资源订阅表、队列加了锁但并发编程的陷阱无处不在。锁的粒度我们为每种事件类型的订阅链表单独设置了锁sub_table-mutex而不是用一个全局锁锁住整个订阅表。这样在发布一个事件时其他类型事件的订阅/取消订阅操作不会被阻塞提高了并发性。在锁内调用用户代码这是死锁的经典根源。绝对不能在持有sub_table-mutex的时候调用用户注册的回调函数handler(evt)。因为用户回调里可能做任何事情包括尝试订阅或取消订阅其他事件这需要获取锁。我们的解决方法是在锁内将需要调用的处理器指针复制到一个本地数组然后释放锁再遍历本地数组调用回调。递归发布如果回调函数A发布了事件B而事件B的回调函数又直接或间接发布了事件A可能会导致无限递归和栈溢出。总线本身很难防止这一点这需要模块设计者遵循良好的实践避免产生事件循环。5.3 事件优先级与顺序保证默认情况下事件按照入队的顺序FIFO被处理。但在某些实时系统中可能需要优先级。例如系统错误事件应该比温度更新事件更优先处理。实现方案可以为事件结构体增加一个priority字段。事件队列不再是一个简单的FIFO队列而是一个优先级队列例如用堆实现。event_bus_process总是从队列中取出优先级最高的事件进行处理。这增加了队列操作的复杂度O(log n)。顺序问题异步模型下即使事件A先于事件B发布由于线程调度B的回调也有可能先于A执行。如果两个事件有严格的先后依赖关系比如“配置加载完成”事件必须在“开始工作”事件之前那么这种模型就不适用。要么改用同步发布publish_sync要么在订阅者逻辑内部通过状态机来处理顺序依赖。5.4 性能优化考量无锁队列在超高并发场景下互斥锁可能成为瓶颈。可以考虑使用无锁lock-free环形缓冲区来实现事件队列。这能极大提升入队/出队的性能但实现难度很高且通常需要处理ABA问题。事件合并对于一些高频但非关键的状态更新事件如每秒100次的传感器读数可以设计一个合并机制。例如在事件总线内部对于同一类型的事件如果队列中已经有一个未处理的新的到来时可以合并比如只保留最新的温度值而不是无限制地入队。这能防止队列被快速产生的事件撑爆也减少了不必要的处理开销。静态订阅对于在编译期就确定的、永不改变的订阅关系可以不使用动态的链表注册而是在编译期就生成一个静态的“事件-处理器”映射表。这完全消除了运行时注册的开销和内存分配适合资源极度受限的嵌入式系统。
返回列表