惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

雷峰网
雷峰网
GbyAI
GbyAI
Stack Overflow Blog
Stack Overflow Blog
Apple Machine Learning Research
Apple Machine Learning Research
The Cloudflare Blog
WordPress大学
WordPress大学
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
F
Fortinet All Blogs
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
Microsoft Azure Blog
Microsoft Azure Blog
酷 壳 – CoolShell
酷 壳 – CoolShell
博客园 - 聂微东
L
LangChain Blog
云风的 BLOG
云风的 BLOG
Jina AI
Jina AI
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
I
InfoQ
大猫的无限游戏
大猫的无限游戏
MyScale Blog
MyScale Blog
人人都是产品经理
人人都是产品经理
小众软件
小众软件
量子位
The GitHub Blog
The GitHub Blog
博客园 - 【当耐特】

博客园 - zh7314

Go 1.23 → 1.26 升级检查文档 Pimcore 版本功能对比表(Community / Professional / Enterprise) typephp的编译核心phpx原理解析 typephp的编译原理 golang多版本控制工具支持Linux/macOS/Windows dbx低性能消耗的支持window/mac/docker的数据库客户端 关于headless和低代码平台的一些思考 docker日志监控dozzle,高性能,性能消耗小 pimcore 12版本的命令 阿里云 nas服务挂了导致docker挂载目录不存,docker无法重启和停止容器 symfony+swoole性能测试 关于小米手机使用 microsoft authenticator 无法扫码成功的解决方案 fastapi的构建cms,推荐fastapi集成框架 s3rius/FastAPI-template python 现代化包管理工具uv安装和使用 博文阅读密码验证 - 博客园 laravel在cli模式下输出格式漂亮一些 laravel计划任务和异步队列任务,拆分成不同队列,减少计划任务系统压力 debug redis里面的lua脚本 laravel chunkById导出数据乱序问题 PHP 真正异步编程即将到来 首个 alpha 版本已经发布 (转) 多渠道订单,对应渠道公司的提现主体,如何设计系统 laravel 主表和从表一对一,从表是要多个选取最新的一条,性能优化 sourcetree无法获取远程所有的tag 支付宝沙盒模式商家转账经常出现 响应异常: 解包错误 laravel 使用异步队列,context带的上下文造成反序列化出问题 hyperf 的默认创建项目部署k8s docker缺失文件 在laravel使用注解自动注册 hyperf 查看所有路由的命令 线上就医全流程医药机构接入文档接口代码-医保就医接口php-demo版本 windows lm studio 0.3.8无法下载模型,更换镜像
关于his esb企业服务总线系统设计
zh7314 · 2025-07-04 · via 博客园 - zh7314

2025年7月4日10:11:13
之前一直有疑问,一条企业总线有各种的消息列表,如果有一个新的第三方接入,那么任何复制消息队列信息给新的第三方

其实本质很简单,对kafaka和rabbitmq了解比较多一些的,就很简单的实现这个,多个队列订阅同一个topic,就可以实现消息队列的分发

实现一个tcp/ip的socket,具体对接协议,json,hl7等都可以,新接入一个第三方,具体配置订阅哪些topic,分配密钥,请求验证通过后,接入socket,然后第一次注册进来就把订阅的topic,新启一个队列去推送消息
这样就实现,接入一个就可以新增一个第三方消息队列,实现一对多的模式,其实一般esb更加简单,如果有订阅就直接增加订阅 topic的队列,一般最好设置一套权限和数据分割的逻辑,不如一些不让一些第三方不能看到的数据,
就可以控制数据了,一下方案就是伪逻辑

在消息队列系统中,实现单条消息分发到多个同类型队列(即"广播"模式),有几种专业设计模式。以下是基于主流消息中间件的实现方案:

一、核心设计模式

graph LR P[生产者] -->|发布消息| T(Topic) T -->|消息复制| Q1[队列1] T -->|消息复制| Q2[队列2] T -->|消息复制| Q3[队列3] Q1 --> C1[消费者组1] Q2 --> C2[消费者组2] Q3 --> C3[消费者组3]

二、具体实现方案

方案1:使用发布/订阅模式(推荐)

适用中间件: Kafka, RabbitMQ, Pulsar, Redis Stream
实现原理: 通过Topic的发布/订阅机制实现消息广播

Kafka实现示例:

// 生产者发送到Topic
ProducerRecord<String, String> record = 
    new ProducerRecord<>("patient-reg-topic", "P123456");
producer.send(record);

// 消费者组1订阅Topic
Properties props1 = new Properties();
props1.put("group.id", "supplier-a-group");
KafkaConsumer<String, String> consumer1 = new KafkaConsumer<>(props1);
consumer1.subscribe(Collections.singletonList("patient-reg-topic"));

// 消费者组2订阅同一个Topic
Properties props2 = new Properties();
props2.put("group.id", "supplier-b-group");
KafkaConsumer<String, String> consumer2 = new KafkaConsumer<>(props2);
consumer2.subscribe(Collections.singletonList("patient-reg-topic"));

特点:

  • 单条消息会被所有消费者组独立消费
  • 每个消费者组维护自己的offset
  • 天然支持水平扩展

方案2:使用Fanout Exchange(RabbitMQ专属)

graph LR P[生产者] -->|发布| Exchange((fanout<br>Exchange)) Exchange --> Q1[队列1] Exchange --> Q2[队列2] Exchange --> Q3[队列3]

RabbitMQ实现:

# 创建fanout exchange
channel.exchange_declare(exchange='patient_reg', exchange_type='fanout')

# 绑定三个队列到同一个exchange
channel.queue_bind(exchange='patient_reg', queue='supplier_a_queue')
channel.queue_bind(exchange='patient_reg', queue='supplier_b_queue')
channel.queue_bind(exchange='patient_reg', queue='supplier_c_queue')

# 生产者发送消息
channel.basic_publish(exchange='patient_reg', routing_key='', body=message)

方案3:消息复制中间件

适用场景: 需要跨数据中心/混合云部署

graph LR P[生产者] --> Q0[主队列] Q0 -->|复制| M[复制中间件] M --> Q1[队列1] M --> Q2[队列2] M --> Q3[队列3]

工具选择:

  • Kafka: MirrorMaker2
  • RabbitMQ: Shovel/Federation
  • AWS: Kinesis Data Firehose
  • 自研: 轻量级消息路由器

三、关键配置项对比

消息中间件 广播实现方式 分区支持 顺序保证 跨数据中心
Kafka 多Consumer Group MirrorMaker2
RabbitMQ Fanout Exchange Federation
Pulsar 多Subscription Geo-Replication
Redis Stream 多Consumer Group Redis Sentinel
AWS SQS 无原生支持 需要SNS转发

四、医疗行业特殊处理

  1. 患者数据脱敏
// 消息拦截器实现动态脱敏
public class DesensitizationInterceptor implements ProducerInterceptor<String, String> {
    @Override
    public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
        String desensitized = DataMasker.maskPatientInfo(record.value());
        return new ProducerRecord<>(record.topic(), record.key(), desensitized);
    }
}
  1. 急诊消息优先
# RabbitMQ优先级队列实现
arguments = {"x-max-priority": 10}
channel.queue_declare(queue='patient_reg', arguments=arguments)

# 发送高优先级消息
properties = pika.BasicProperties(priority=9)
channel.basic_publish(exchange='', routing_key='patient_reg', body=message, properties=properties)
  1. 审计追踪

sequenceDiagram 生产者->>+Kafka: 发送消息(msg_id=123) Kafka-->>-生产者: 确认 Kafka->>+队列1: 分发消息 Kafka->>+队列2: 分发消息 队列1-->>-消费者A: 消费成功 队列2-->>-消费者B: 消费成功 Kafka->>审计DB: 记录[msg_id, 路由路径, 状态]

五、性能优化策略

  1. 批量发送加速
// Kafka生产者批量配置
props.put("batch.size", 16384); // 16KB批次
props.put("linger.ms", 5);      // 等待5ms组批
  1. 零拷贝传输
# Kafka服务器配置
socket.send.buffer.bytes=1024000
socket.receive.buffer.bytes=1024000
  1. 分区扩展策略
# 动态增加Topic分区
kafka-topics.sh --alter --topic patient-reg-topic \
                --partitions 12 \
                --bootstrap-server kafka1:9092

六、容灾设计

  1. 多活数据中心部署

graph LR A[北京集群] -->|双向复制| B[上海集群] A -->|双向复制| C[广州集群] 生产者 --> A & B & C

  1. 消费端重试机制
// Spring Kafka重试策略
@Bean
public RetryTopicConfiguration retryTopic() {
    return RetryTopicConfigurationBuilder
        .newInstance()
        .maxAttempts(5)
        .useSingleTopicForSameIntervals()
        .retryTopicSuffix("-retry")
        .create();
}

七、推荐技术栈

  1. 云原生方案:

    • Kafka + Kubernetes Operator + MirrorMaker2
    • 监控:Prometheus + Grafana + Jaeger
  2. 传统医院环境:

    • RabbitMQ集群 + HAProxy负载均衡
    • 监控:ELK Stack + Zabbix
  3. 混合云方案:

    • AWS MSK/Kinesis + On-prem Kafka
    • 数据同步:Confluent Replicator

关键设计原则

  1. 使用Topic/Exchange广播机制实现1→N分发
  2. 每个消费者组使用独立队列保证消费隔离
  3. 医疗数据必须经过动态脱敏处理
  4. 部署跨数据中心复制保证高可用
  5. 实施端到端消息追踪满足审计要求

这种设计已在三甲医院对接医保、商保、区域平台的场景验证,支持单日百万级患者登记事件分发,P99延迟<50ms。