【多进程Topic通信系统设计文档】
多进程Topic通信系统设计文档
- 1. 系统概述
- 1.1 项目背景
- 1.2 系统目标
- 2. 架构设计
- 2.1 整体架构
- 架构图概览
- 核心组件说明
- 1. 客户端层(Client Layer)
- 2. Topic Broker(TCP通信层)
- 3. 队列层(Queue Layer)
- 4. 处理层(Processing Layer)
- 5. 结果与配置管理层
- 数据流向
- 设计特点
- 2.2 核心组件
- 2.2.1 SafeQueue(安全队列)
- 2.2.2 TopicBroker(消息代理)
- 2.2.3 Workers(工作进程)
- 2.2.4 ProcessManager(进程管理器)
- 3. 客户端设计
- 3.1 Data Client(数据发布客户端)
- 3.2 Config Client(配置变更客户端)
- 3.3 Result Subscriber(结果订阅客户端)
- 4. 数据流设计
- 4.1 小数据流
- 4.2 大数据流
- 4.3 配置更新流
- 5. 接口设计
- 5.1 网络协议
- 5.2 消息规范
- 6. 配置管理
- 6.1 系统配置(SystemConfig)
- 6.2 动态配置(Shared Config)
- 7. 异常处理策略
- 7.1 队列异常
- 7.2 网络异常
- 7.3 进程异常
- 8. 性能考虑
- 8.1 并发模型
- 8.2 优化策略
- 9. 部署说明
- 9.1 环境要求
- 9.2 启动流程
- 9.3 监控与日志
- 10. 总结
- 代码下载URL
1. 系统概述
1.1 项目背景
本系统是一个基于TCP Socket的多进程发布/订阅消息通信框架,实现了进程间的高效数据流转和处理。系统采用生产者-消费者模式,通过安全队列机制解耦各组件,支持配置动态更新、数据处理和结果反馈。
1.2 系统目标
实现多进程间的可靠消息通信 支持多种数据类型(小数据/大数据)的分类处理 提供动态配置更新能力 保证队列操作的安全性(防溢出、防阻塞) 支持结果的实时订阅和反馈2. 架构设计
2.1 整体架构
系统采用分层架构设计,核心组件包括客户端层、Topic Broker(TCP通信层)和数据处理层,通过安全队列机制实现解耦和异步通信。
架构图概览
┌─────────────────────────────────────────────────────────────┐ │ 客户端层 │ ├──────────────┬──────────────┬──────────────────────────────┤ │ data_client │ config_client│ result_subscriber │ └──────┬───────┴──────┬───────┴──────────┬───────────────────┘ │ │ │ └──────────────┼───────────────────┘ ▼ ┌─────────────────────────────────────────────────────────────┐ │ Topic Broker (TCP) │ │ - 订阅管理(subscribers dict) │ │ - 消息路由(按topic分发) │ │ - 结果发布(独立线程) │ └──┬────────────┬────────────┬───────────────────────────────┘ │ │ │ ▼ ▼ ▼ ┌─────────┐ ┌─────────┐ ┌──────────┐ │DB Queue │ │Proc Queue│ │Config Q │ └────┬────┘ └────┬────┘ └────┬─────┘ │ │ │ ▼ ▼ ▼ ┌─────────┐ ┌─────────┐ ┌──────────┐ │DB Writer│ │Processor│ │Config │ │(TDengine)│ │(×N) │ │Listener │ └─────────┘ └─────────┘ └──────────┘ │ │ └───────────┼──────────────────┐ ▼ ▼ ┌───────────┐ ┌────────────┐ │Result Queue│ │Shared Config│ └─────┬─────┘ └────────────┘ │ ▼ ┌───────────┐ │Broker→订阅者│ └───────────┘核心组件说明
1. 客户端层(Client Layer)
- data_client:数据生产者,负责发送业务数据到系统
- config_client:配置客户端,用于动态更新系统配置
- result_subscriber:结果订阅者,接收处理结果的实时反馈
2. Topic Broker(TCP通信层)
- 订阅管理:维护订阅者字典(subscribers dict),记录各topic的订阅者
- 消息路由:根据消息topic将数据分发到对应的处理队列
- 结果发布:独立线程负责将处理结果推送回订阅者
3. 队列层(Queue Layer)
- DB Queue:数据库写入队列,存储需要持久化的数据
- Proc Queue:数据处理队列,分发到多个Processor进行并发处理
- Config Q:配置更新队列,处理动态配置变更
4. 处理层(Processing Layer)
- DB Writer:数据库写入器,将数据写入TDengine时序数据库
- Processor(×N):多进程数据处理单元,支持水平扩展
- Config Listener:配置监听器,实时响应配置变更
5. 结果与配置管理层
- Result Queue:结果队列,收集各Processor的处理结果
- Shared Config:共享配置,确保所有组件配置一致性
- Broker→订阅者:结果推送通道,将最终结果返回给订阅者
数据流向
- 数据流:data_client → Topic Broker → Proc Queue → Processor → Result Queue → Broker→订阅者
- 配置流:config_client → Topic Broker → Config Q → Config Listener → Shared Config
- 存储流:data_client → Topic Broker → DB Queue → DB Writer → TDengine
设计特点
- 解耦设计:各组件通过队列解耦,提高系统可维护性
- 水平扩展:Processor支持多实例部署,提升处理能力
- 实时反馈:结果订阅机制确保处理结果及时返回
- 配置热更新:支持运行时动态调整系统参数
- 安全队列:防溢出、防阻塞机制保证系统稳定性
2.2 核心组件
2.2.1 SafeQueue(安全队列)
基于multiprocessing.Queue的封装,提供异常处理和超时机制。
特性:
- 统一的异常捕获(Full/Empty异常)
- 超时控制(get/put方法)
- 日志记录(操作成功/失败)
- 非阻塞操作支持(get_nowait/put_nowait)
关键方法:
defput(item,timeout=None,block=True)->booldefget(timeout=1.0,default=None)->Anydefget_nowait(default=None)->Anydefput_nowait(item)->bool2.2.2 TopicBroker(消息代理)
基于TCP Socket的发布/订阅服务器。
职责:
- 管理客户端连接(accept_loop)
- 维护订阅关系(subscribers字典)
- 路由消息到内部队列或外部订阅者
- 结果发布(独立线程)
消息格式:
{"action":"subscribe"|"publish","topic":"data"|"config"|"result","payload":{...}}Topic路由规则:
| Topic | 路由目标 | 处理逻辑 |
|---|---|---|
| data (small) | DB Queue | 写入TDengine |
| data (large) | Process Queue | 处理任务 |
| config | Config Queue | 更新配置 |
| result | 订阅者 | 实时推送 |
2.2.3 Workers(工作进程)
DB Writer Worker:
- 从DB Queue读取小数据
- 写入TDengine(自动创建子表)
- 支持数据持久化
Processor Worker:
- 从Process Queue读取大数据
- 模拟处理任务(随机耗时0.5-2秒)
- 读取共享配置快照
- 发送结果到Result Queue
Config Listener Worker:
- 从Config Queue读取配置变更
- 更新Shared Config
- 支持动态参数调整
2.2.4 ProcessManager(进程管理器)
管理所有进程的生命周期。
职责:
- 创建和启动所有工作进程
- 管理队列和共享配置
- 处理信号(SIGINT/SIGTERM)
- 优雅关闭所有组件
3. 客户端设计
3.1 Data Client(数据发布客户端)
- 支持发送小数据(写入数据库)
- 支持发送大数据(处理队列)
- 支持批量发送混合数据
- 自动订阅result topic接收反馈
3.2 Config Client(配置变更客户端)
- 自定义配置键值对
- 快捷修改(batch_size/version)
- 查看可配置项列表
3.3 Result Subscriber(结果订阅客户端)
- 实时接收处理结果
- 结构化显示结果信息
- 持续监听模式
4. 数据流设计
4.1 小数据流
Data Client → Broker → DB Queue → DB Writer → TDengine4.2 大数据流
Data Client → Broker → Process Queue → Processor → Result Queue → Broker → Subscriber4.3 配置更新流
Config Client → Broker → Config Queue → Config Listener → Shared Config5. 接口设计
5.1 网络协议
- 传输协议:TCP
- 消息格式:JSON + 换行符分隔
- 端口:9999
- 编码:UTF-8
5.2 消息规范
订阅消息:
{"action":"subscribe","topic":"result"}发布消息:
{"action":"publish","topic":"data","payload":{"data_type":"small","name":"test_data","content":{"value":"test"}}}结果消息:
{"topic":"result","payload":{"worker_id":0,"task_name":"task_1","status":"success","result":"processed","process_time":"1.23秒","record_count":100,"config_snapshot":{...},"timestamp":"2026-08-02 10:30:00"}}6. 配置管理
6.1 系统配置(SystemConfig)
@dataclassclassSystemConfig:worker_num:int=3# 处理器数量broker_host:str="127.0.0.1"broker_port:int=9999td_url:str="http://127.0.0.1:6041"td_user:str="root"td_pass:str="taosdata"td_db:str="pipe_demo"queue_timeout:float=1.0# 队列超时时间max_queue_size:int=1000# 最大队列大小6.2 动态配置(Shared Config)
- app_name: 应用名称
- version: 版本号
- batch_size: 批处理大小
- 支持动态添加新配置项
7. 异常处理策略
7.1 队列异常
- Full异常:记录警告日志,丢弃数据
- Empty异常:返回默认值,继续等待
- 其他异常:记录错误日志,不影响主流程
7.2 网络异常
- 客户端断开:移除订阅者,清理资源
- 连接拒绝:友好提示用户启动Broker
- 发送失败:移除断开的客户端
7.3 进程异常
- 信号处理:优雅关闭所有进程
- 超时强制终止:5秒超时后使用terminate()
- 资源清理:关闭队列和连接
8. 性能考虑
8.1 并发模型
- Broker:多线程(每个客户端独立线程)
- Workers:多进程(CPU密集型任务)
- 队列:进程安全的Queue
8.2 优化策略
- 非阻塞I/O(select.select)
- 批量发送支持
- 配置缓存(共享字典)
- 队列大小限制(防溢出)
9. 部署说明
9.1 环境要求
- Python 3.7+
- TDengine(可选,用于数据持久化)
- 依赖:taosrest, multiprocessing
9.2 启动流程
- 启动Broker和Workers:
python pipe_demon_v1.py - 启动数据客户端:
python data_client.py - 启动配置客户端:
python config_client.py - 启动结果订阅者:
python result_subscriber.py
9.3 监控与日志
- 系统运行状态监控
- 队列深度监控
- 处理延迟统计
- 错误日志记录
10. 总结
本系统通过分层架构设计,实现了多进程间的可靠消息通信。核心特点包括:
- 解耦设计:各组件通过安全队列解耦,提高系统可维护性
- 灵活扩展:支持水平扩展,可根据负载动态调整Processor数量
- 实时反馈:完整的发布/订阅机制,确保处理结果及时返回
- 配置热更新:支持运行时动态调整系统参数
- 异常健壮:完善的异常处理机制,保证系统稳定性
系统适用于需要高并发、低延迟、可扩展的数据处理场景,特别适合物联网数据采集、实时计算、分布式任务处理等应用场景。
代码下载URL
多进程Topic通信系统设计