【多进程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→订阅者:结果推送通道,将最终结果返回给订阅者

数据流向

  1. 数据流:data_client → Topic Broker → Proc Queue → Processor → Result Queue → Broker→订阅者
  2. 配置流:config_client → Topic Broker → Config Q → Config Listener → Shared Config
  3. 存储流: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)->bool

2.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处理任务
configConfig 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 → TDengine

4.2 大数据流

Data Client → Broker → Process Queue → Processor → Result Queue → Broker → Subscriber

4.3 配置更新流

Config Client → Broker → Config Queue → Config Listener → Shared Config

5. 接口设计

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 启动流程

  1. 启动Broker和Workerspython pipe_demon_v1.py
  2. 启动数据客户端python data_client.py
  3. 启动配置客户端python config_client.py
  4. 启动结果订阅者python result_subscriber.py

9.3 监控与日志

  • 系统运行状态监控
  • 队列深度监控
  • 处理延迟统计
  • 错误日志记录

10. 总结

本系统通过分层架构设计,实现了多进程间的可靠消息通信。核心特点包括:

  1. 解耦设计:各组件通过安全队列解耦,提高系统可维护性
  2. 灵活扩展:支持水平扩展,可根据负载动态调整Processor数量
  3. 实时反馈:完整的发布/订阅机制,确保处理结果及时返回
  4. 配置热更新:支持运行时动态调整系统参数
  5. 异常健壮:完善的异常处理机制,保证系统稳定性

系统适用于需要高并发、低延迟、可扩展的数据处理场景,特别适合物联网数据采集、实时计算、分布式任务处理等应用场景。

代码下载URL

多进程Topic通信系统设计