事件驱动架构EDA与事件中心完整知识梳理

发布时间:2026/10/5 4:20:22

事件驱动架构EDA与事件中心完整知识梳理
事件驱动架构EDA与事件中心完整知识梳理一、什么是事件概念事件Event是系统中发生的一个事实的描述。它表示某件事已经发生了而不是请你做某件事。对比项命令Command事件Event语义“请你做XX”“XX已经发生了”方向发送方知道接收方发送方不关心谁订阅耦合度强耦合松耦合示例createOrder(orderDto)OrderCreatedEvent时态祈使句/将来时过去时事件的三要素{type:order.created,// 事件类型发生了什么source:order-service,// 事件来源谁产生的data:{orderId:123},// 事件数据具体内容time:2026-07-14T10:28:00Z// 事件时间什么时候发生的}注博客https://blog.csdn.net/badao_liumang_qizhi二、事件驱动架构EDA核心模型┌──────────┐ ┌──────────────┐ ┌──────────┐ │ 生产者 │ 发布事件 │ 事件中心 │ 推送事件 │ 消费者 │ │ Producer │ ────────→ │ Event Broker │ ────────→ │ Consumer │ └──────────┘ └──────────────┘ └──────────┘ │ │ 推送事件 ▼ ┌──────────┐ │ 消费者 B │ └──────────┘与传统同步调用对比同步调用点对点a服务 ──HTTP POST──→ b服务 ←──Response────问题b 挂了a 的业务也跟着失败。事件驱动发布/订阅a服务 ──发布事件──→ 事件中心 ──推送──→ b服务 ──推送──→ 监控服务 ──推送──→ 日志服务优势b挂了不影响 a将来新增消费者无需改 a 代码。适用场景系统间异步通知不需要同步等待结果一个事件需要通知多个下游系统生产者和消费者由不同团队维护需要解耦上下游系统的部署和发布周期三、CloudEvents 规范CloudEvents 是 CNCF云原生计算基金会定义的事件格式标准目的是让不同系统产生的事件有统一的格式。核心属性属性含义示例specversion规范版本1.0type事件类型aaa.xxxsource事件来源/bbb-serviceid事件唯一IDuuidtime事件产生时间2026-07-14T10:28:00Zdata事件负载业务数据JSON 对象HTTP 协议绑定CloudEvents 通过 HTTP Header 传递元信息POST /event-center/api/v3/eventTypes/aaa.xxx/send Content-Type: application/json Ce-Type: aaa.xxx Ce-Specversion: 1.0 Authorization: Bearer token { fulfillmentCode: xx.20260714.000005.xx, event: ttt, eventDate: 2026-07-14 10:28:00, systemName: uuu }四、事件中心的角色事件中心Event Center/Event Broker是事件流转的中间枢纽┌─────────────────────────────────────────────────────────┐ │ 事件中心 │ │ │ │ ┌────────────┐ ┌────────────┐ ┌────────────────┐ │ │ │ 事件接收 │ │ 事件路由 │ │ 事件投递 │ │ │ │ (Ingress) │→ │ (Routing) │→ │ (Delivery) │ │ │ └────────────┘ └────────────┘ └────────────────┘ │ │ │ │ 职责 │ │ 1. 接收生产者发布的事件 │ │ 2. 根据 eventType 路由到对应订阅者 │ │ 3. 可靠投递重试、死信队列 │ │ 4. 认证鉴权 │ └─────────────────────────────────────────────────────────┘与 MQ 的区别对比项MQRabbitMQ/Kafka事件中心范围同一系统/团队内部跨系统/跨团队/跨公司协议AMQP/Kafka协议HTTPCloudEvents订阅方式消费者主动拉取事件中心主动推送Webhook运维方自己运维平台方统一运维格式自定义标准化CloudEvents五、通用代码示例1. 事件发布工具类/** * 事件中心推送工具类. * 基于 CloudEvents 规范通过 HTTP POST 推送事件到事件中心. */ComponentSlf4jpublicclassEventPublisher{Value(${event-center.token})privateStringtoken;privatefinalRestTemplaterestTemplate;publicEventPublisher(RestTemplaterestTemplate){this.restTemplaterestTemplate;}/** * 发布事件到事件中心. * * param eventType 事件类型如 order.created * param url 事件中心接收地址 * param payload 事件数据JSON 字符串 * return 推送结果 */publicEventResultpublish(StringeventType,Stringurl,Stringpayload){try{HttpHeadersheadersnewHttpHeaders();headers.setContentType(MediaType.APPLICATION_JSON);headers.set(Authorization,token);headers.set(Ce-Type,eventType);headers.set(Ce-Specversion,1.0);HttpEntityStringrequestnewHttpEntity(payload,headers);ResponseEntityStringresponserestTemplate.postForEntity(url,request,String.class);log.info(事件发布成功, type{}, response{},eventType,response.getBody());returnEventResult.success(response.getBody());}catch(Exceptione){log.warn(事件发布失败, type{},eventType,e);returnEventResult.fail(e.getMessage());}}}2. 本地消息表 可靠发布/** * 可靠事件发布服务. * 先写本地消息表事务提交后再触发实际推送. */ServicepublicclassReliableEventPublisher{ResourceprivateEventLogRepositoryeventLogRepository;ResourceprivateEventMqSendereventMqSender;/** * 记录待发布事件在业务事务内调用. */TransactionalpublicvoidprepareEvent(StringeventType,Stringpayload,StringbizKey){// 1. 写入本地消息表EventLogeventLognewEventLog();eventLog.setEventType(eventType);eventLog.setPayload(payload);eventLog.setBizKey(bizKey);eventLog.setStatus(0);// 待发送eventLog.setRetryCount(0);eventLogRepository.save(eventLog);// 2. 注册事务提交后的回调发送 MQAfterTransactionActionCollectorcollectornewAfterTransactionActionCollector();TransactionSynchronizationManager.registerSynchronization(collector);collector.addCommitSyncAction(()-eventMqSender.send(eventLog.getId()));}}3. MQ 消费者 实际推送/** * 事件推送 MQ 消费者. * 消费本地消息表 ID执行实际推送. */ComponentSlf4jpublicclassEventPublishConsumer{ResourceprivateEventLogRepositoryeventLogRepository;ResourceprivateEventPublishereventPublisher;RabbitListener(queues${event.publish.queue})publicvoidconsume(IntegereventLogId){EventLogeventLogeventLogRepository.findById(eventLogId).orElse(null);if(eventLognull){return;}// 已成功的跳过if(eventLog.getStatus()1){return;}// 推送事件中心StringurlresolveUrl(eventLog.getEventType());EventResultresulteventPublisher.publish(eventLog.getEventType(),url,eventLog.getPayload());// 更新状态if(result.isSuccess()){eventLog.setStatus(1);eventLog.setResponse(result.getBody());}else{eventLog.setStatus(2);eventLog.setErrorMsg(result.getErrorMsg());}eventLog.setRetryCount(eventLog.getRetryCount()1);eventLogRepository.save(eventLog);}}4. 事务同步回调机制/** * Spring 事务同步回调示例. * 确保 MQ 消息在事务提交后才发送. */publicclassTransactionCallbackDemo{publicvoidbusinessMethod(){// ... 业务逻辑操作数据库 ...// 注册事务提交后的回调TransactionSynchronizationManager.registerSynchronization(newTransactionSynchronization(){OverridepublicvoidafterCommit(){// 事务提交成功后执行mqSender.send(message);}});}}5. 运维重推接口/** * 运维重推接口. * 用于手动触发失败事件的重新发布. */RestControllerRequestMapping(/api/inner/event)publicclassEventRetryController{ResourceprivateEventMqSendereventMqSender;ResourceprivateEventLogRepositoryeventLogRepository;PostMapping(/retry)publicResultBooleanretry(RequestBodyListIntegerlogIds){if(logIdsnull||logIds.isEmpty()){// 未指定ID则查询所有失败记录ListEventLogfailedLogseventLogRepository.findByStatusInAndRetryCountLessThan(Arrays.asList(0,2),6);logIdsfailedLogs.stream().map(EventLog::getId).collect(Collectors.toList());}logIds.forEach(eventMqSender::send);returnResult.success(true);}}六、完整数据流示意┌─ 你的服务内部 ─────────────────────────────┐ │ │ 业务代码 │ ① 写入本地消息表status0 │ │ │ ② 注册事务回调 │ ▼ │ │ │ 事务提交 ───────────│─────────┼────────────────────────────────────│ │ ▼ │ │ ③ 回调触发发送 MQ消息logId │ │ │ │ │ ▼ │ │ ④ MQ 消费者读取消息表 → 组装 HTTP 请求 │ │ │ │ └─────────┼────────────────────────────────────┘ │ ▼ HTTP POSTCloudEvents 格式 ┌─────────────────────┐ │ 事件中心 │ │ 接收 → 路由 → 投递 │ └─────────┬───────────┘ │ Webhook 推送 ▼ ┌─────────────────────┐ │ b 服务 │ │ 接收事件 → 处理业务 │ └─────────────────────┘七、核心概念速查表概念一句话解释事件Event已发生事实的描述过去时事件类型Event Type事件的分类标识用于路由事件中心Event Broker负责接收、路由、投递事件的中间平台生产者Producer发布事件的服务消费者Consumer/Subscriber订阅并处理事件的服务CloudEventsCNCF 定义的事件格式标准本地消息表Outbox确保事件可靠发布的数据库表事务同步Transaction Sync确保消息在事务提交后才发送最终一致性允许短暂不一致但最终状态一致幂等性同一事件重复投递不会产生副作用

相关新闻

【限时解密】ChatGPT私域社群GMV提升217%的底层逻辑:基于27个真实社群的NLP情感分析报告

【限时解密】ChatGPT私域社群GMV提升217%的底层逻辑:基于27个真实社群的NLP情感分析报告

2026/10/5 4:18:09

更多请点击: https://codechina.net 第一章:【限时解密】ChatGPT私域社群GMV提升217%的底层逻辑:基于27个真实社群的NLP情感分析报告 通过对27个垂直行业(教育、SaaS、跨境电商、知识付费等)的活跃ChatGPT私域社群进行…

Copilot团队管理功能落地全图谱(从权限配置到效能看板):微软内部PM未公开的7条黄金规则

Copilot团队管理功能落地全图谱(从权限配置到效能看板):微软内部PM未公开的7条黄金规则

2026/10/2 9:47:38

更多请点击: https://codechina.net 第一章:Copilot团队管理功能全景概览 GitHub Copilot 的团队管理功能为组织级开发协作提供了统一策略配置、成员权限控制与使用洞察分析能力,覆盖从开发者准入到资源配额分配的全生命周期管理。该功能依托…

计算机考试-n叉树节点度易混淆—东方仙盟

计算机考试-n叉树节点度易混淆—东方仙盟

2026/9/26 11:37:46

n0​:度 0,叶子节点(无孩子)n1​:度 1,只有左孩子 或 只有右孩子n2​:度 2,同时有左、右孩子公式 1(总节点)Nn0​n1​n2​公式 2(重中之重&#x…

CANN/GE ACL数据集缓冲区添加函数

CANN/GE ACL数据集缓冲区添加函数

2026/10/4 10:56:09

aclmdlAddDatasetBuffer 【免费下载链接】ge GE(Graph Engine)是面向昇腾的图编译器和执行器,提供了计算图优化、多流并行、内存复用和模型下沉等技术手段,加速模型执行效率,减少模型内存占用。 GE 提供对 PyTorch、Te…

用ffmpeg高效批量调整图片尺寸的实战指南

用ffmpeg高效批量调整图片尺寸的实战指南

2026/10/4 10:54:41

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

Transformers 音频特征提取工具库 audio_utils 全解析:从 Mel 刻度换算到对数 Mel 频谱

Transformers 音频特征提取工具库 audio_utils 全解析:从 Mel 刻度换算到对数 Mel 频谱

2026/10/4 10:55:43

Transformers 音频特征提取工具库 audio_utils 全解析:从 Mel 刻度换算到对数 Mel 频谱 【免费下载链接】transformers 🤗 Transformers: the model-definition framework for state-of-the-art machine learning models in text, vision, audio, and mu…

RustFS 多节点集群重启与滚动升级实战:Readiness、Quorum 与 Degraded 模式完全指南

RustFS 多节点集群重启与滚动升级实战:Readiness、Quorum 与 Degraded 模式完全指南

2026/10/4 17:30:40

RustFS 多节点集群重启与滚动升级实战:Readiness、Quorum 与 Degraded 模式完全指南 【免费下载链接】rustfs 🚀2.3x faster than MinIO for 4KB object payloads. RustFS is an open-source, S3-compatible high-performance object storage system sup…

Java Integer缓存揭秘:128陷阱原理、避坑与面试全解

Java Integer缓存揭秘:128陷阱原理、避坑与面试全解

2026/10/4 5:19:11

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

RustFS Scanner 数据用量发布权威性决策:配额准入如何获得可用的权威依据

RustFS Scanner 数据用量发布权威性决策:配额准入如何获得可用的权威依据

2026/10/4 10:55:30

RustFS Scanner 数据用量发布权威性决策:配额准入如何获得可用的权威依据 【免费下载链接】rustfs 🚀2.3x faster than MinIO for 4KB object payloads. RustFS is an open-source, S3-compatible high-performance object storage system supporting mi…