Skip to content

Latest commit

 

History

History
334 lines (226 loc) · 15.1 KB

File metadata and controls

334 lines (226 loc) · 15.1 KB

Apache EventMesh

Apache EventMesh 是用于构建分布式事件驱动应用程序的新一代无服务器事件中间件。

EventMesh 架构

EventMesh Architecture

EventMesh 采用统一的 CloudEvents-over-MQ 架构。消息队列(MQ)仅作为持久化的预写日志(WAL)存储层——没有消费者组(consumer group)、没有 Tag,也不承担 Broker 端的订阅语义。所有投递逻辑由无状态的 EventMesh Runtime 掌控:其中的 SubscriptionManager 负责维护订阅注册表与位点(offset)跟踪,并以负载均衡(load-balance)、广播(broadcast)、多播(multicast)等语义分发事件。应用通过轻量的 HTTP + CloudEvents 1.0 SDK(publish / subscribe / unsubscribe)接入,与外部系统的集成则通过 connector SPI 在独立进程的 Connector Runtime 中运行。

特性

Apache EventMesh 提供了丰富的能力,帮助用户轻松构建事件驱动应用。以下是让 EventMesh 与众不同的核心亮点:

核心架构

  • 全链路 CloudEvents 原生 —— 完全基于 CloudEvents 1.0 规范构建,事件天然厂商中立、可移植。
  • 轻量、语言无关的 SDK —— 纯 HTTP 上仅需 publishsubscribeunsubscribe 三个操作,无重客户端、无厂商锁定。
  • 订阅与分发由 Runtime 掌控 —— 订阅状态与投递语义(负载均衡 / 广播 / 多播)由 EventMesh 自主管理,而非底层 MQ,无论后端是哪种存储,行为都保持一致。
  • MQ 仅作纯预写日志(WAL) —— 仅追加写入,无消费者组、无 Tag;Broker 退化为纯持久化存储,大幅简化运维。
  • 保证至少一次(at-least-once)投递 —— 可靠性由 EventMesh 通过自管 offset + 显式 ACK 掌控。
  • 多种下发通道 —— 订阅者可自由选择 HTTP 长轮询、Server-Sent Events(SSE)或 WebSocket 推送,并支持请求-应答(request-reply)。
  • 轻松水平扩展 —— 无状态 Runtime 实例可无缝扩容,无重平衡成本。

扩展性与生态

  • 智能体协作(A2A) —— 内置的 A2A 协议 将 EventMesh 打造成智能体协作总线,打通同步的 MCP / JSON-RPC 2.0 工具调用与异步的事件驱动发布订阅,原生支撑大模型(LLM)与多智能体(Multi-Agent)场景。
  • 可插拔存储层 —— Apache RocketMQApache KafkaApache PulsarRabbitMQRedis 等。
  • 可插拔互联层(Connector) —— connectors 作为独立进程运行,可充当 SaaS、CloudService、数据库等的 source 或 sink。
  • 可插拔元数据服务 —— ConsulNacosETCDZookeeper
  • 事件模式管理 —— 通过目录(catalog)服务实现。
  • 强大的事件编排 —— 基于 Serverless workflow 引擎。
  • 强大的事件过滤与转换能力。

子项目

快速入门

本节指南将指导您分别从本地DockerK8s部署EventMesh的步骤:

本节指南按照默认配置启动 EventMesh,如果您需要更加详细的 EventMesh 部署步骤,请访问EventMesh官方文档

部署 Event Store

EventMesh 支持多种事件存储,默认存储模式为 standalone,不依赖其他事件存储作为层。

在本地运行 EventMesh Runtime

1. 下载 EventMesh

EventMesh Download 页面下载最新版本的 Binary Distribution 发行版并解压:

wget https://dlcdn.apache.org/eventmesh/1.10.0/apache-eventmesh-1.10.0-bin.tar.gz
tar -xvzf apache-eventmesh-1.10.0-bin.tar.gz
cd apache-eventmesh-1.10.0

2. 运行 EventMesh

执行 start.sh 脚本启动 EventMesh Runtime 服务器。

bash bin/start.sh

查看输出日志:

tail -n 50 -f logs/eventmesh.out

当日志输出 server state:RUNNING,则代表 EventMesh Runtime 启动成功了。

停止:

bash bin/stop.sh

脚本打印 shutdown server ok! 时,代表 EventMesh Runtime 已停止。

在 Docker 中运行 EventMesh Runtime

1. 获取 EventMesh 镜像

使用下面的命令行下载最新版本的 EventMesh

sudo docker pull apache/eventmesh:latest

2. 运行 EventMesh

使用以下命令启动 EventMesh 容器:

sudo docker run -d --name eventmesh -p 8080:8080 -p 8081:8081 -p 8082:8082 -t apache/eventmesh:latest

端口说明:8080 = 流量 HTTP(/events/* CloudEvents API),8081 = 管理 HTTP(/admin/*),8082 = WebSocket 推送(按需开启)。Connector 运行时镜像另暴露 8083 作为其可选的管理端口。

进入容器:

sudo docker exec -it eventmesh /bin/bash

查看日志:

cd logs
tail -n 50 -f eventmesh.out

在 Kubernetes 中运行 EventMesh Runtime

1. 部署 Operator

运行以下命令部署(删除部署, 只需将 deploy 替换为 undeploy 即可):

$ cd eventmesh-operator && make deploy

运行 kubectl get podskubectl get crd | grep eventmesh-operator.eventmesh 查看部署的 EventMesh-Operator 状态以及 CRD 信息.

$ kubectl get pods
NAME                                  READY   STATUS    RESTARTS   AGE
eventmesh-operator-59c59f4f7b-nmmlm   1/1     Running   0          20s

$ kubectl get crd | grep eventmesh-operator.eventmesh
connectors.eventmesh-operator.eventmesh   2024-01-10T02:40:27Z
runtimes.eventmesh-operator.eventmesh     2024-01-10T02:40:27Z

2. 部署 EventMesh Runtime

运行以下命令部署 runtime、connector-rocketmq (删除部署, 只需将 create 替换为 delete 即可):

$ make create

运行 kubectl get pods 查看部署是否成功.

NAME                                  READY   STATUS    RESTARTS   AGE
connector-rocketmq-0                  1/1     Running   0          9s
eventmesh-operator-59c59f4f7b-nmmlm   1/1     Running   0          3m12s
eventmesh-runtime-0-a-0               1/1     Running   0          15s

使用 EventMesh API

应用通过 /events/* HTTP 端点以标准 CloudEvents 1.0 与 EventMesh 交互,仅需 publish / subscribe / unsubscribe 等操作。订阅与分发由 EventMesh 自主管理,无需消费者组或 Tag。

1. 发布 CloudEvent

以 CloudEvents 1.0 格式通过 HTTP POST 发布事件(返回 202 Accepted 表示事件已写入 WAL):

POST /events/publish HTTP/1.1
Host: localhost:8080
Content-Type: application/cloudevents+json

{
  "specversion": "1.0",
  "id": "89010a5a-3c6f-4a1e-9b2d-0f7c1f2e3a4b",
  "source": "/example/producer",
  "type": "com.example.order.created",
  "subject": "orders",
  "datacontenttype": "application/json",
  "data": {
    "content": "Hello, EventMesh!"
  }
}

2. 订阅 Topic

向 EventMesh 注册订阅(无消费者组、无 Tag)。提供 clientIdtopic 以及分发模式 modeLOAD_BALANCEBROADCASTMULTICASTLOAD_BALANCE_STICKY),响应返回 subscriptionId

POST /events/subscribe HTTP/1.1
Host: localhost:8080
Content-Type: application/json

{
  "clientId": "order-svc",
  "topic": "orders",
  "mode": "LOAD_BALANCE"
}

3. 接收事件

订阅者可通过三种传输方式接收下发的事件,按业务场景任选其一——它们下发相同的 CloudEvents,并遵循 EventMesh 的 ACK / 至少一次语义。

a) HTTP 长轮询 —— 主动拉取,请求会阻塞直到有事件或超时:

GET /events/poll?clientId=order-svc&topics=orders&timeout=30000 HTTP/1.1
Host: localhost:8080

处理完事件后,进行 ACK 以推进 offset(至少一次投递):

POST /events/ack HTTP/1.1
Host: localhost:8080
Content-Type: application/json

{
  "subId": "sub-123",
  "clientId": "order-svc",
  "topic": "orders",
  "partition": 0,
  "offset": 42
}

b) Server-Sent Events(SSE) —— 服务端经长连接主动推送,客户端无需轮询:

GET /events/stream?clientId=order-svc&topics=orders HTTP/1.1
Host: localhost:8080
Accept: text/event-stream

c) WebSocket —— 经独立 WebSocket 端口的全双工服务端推送:

GET /events/stream HTTP/1.1
Host: localhost:8082
Upgrade: websocket
Connection: Upgrade

长轮询、SSE、WebSocket 是可互换的下发传输——订阅者任选其一。SSE 与 WebSocket 由服务端推送(无需轮询循环),长轮询则由客户端驱动。此外还支持请求-应答(POST /events/request + POST /events/reply)。CloudEventsClient Java SDK 封装了以上全部方式(subscribe / subscribeSse / subscribeWs),用法见CloudEvents 客户端使用指引

4. 取消订阅

当不再需要接收某 Topic 的事件时,按 clientId 取消订阅(可选带 subscriptionId):

POST /events/unsubscribe HTTP/1.1
Host: localhost:8080
Content-Type: application/json

{
  "clientId": "order-svc",
  "topic": "orders"
}

贡献

每个贡献者在推动 Apache EventMesh 的健康发展中都发挥了重要作用。我们真诚感谢所有为代码和文档作出贡献的贡献者。

这里是贡献者列表,感谢大家! :)

CNCF Landscape

Apache EventMesh enriches the CNCF Cloud Native Landscape.

License

Apache EventMesh 的开源协议遵循 Apache License, Version 2.0.

Community

微信小助手 微信公众号 Slack
加入 Slack

双周会议 : #Tencent meeting : 346-6926-0133

双周会议记录 : bilibili

邮件名单

名称 描述 订阅 取消订阅 邮件列表存档
用户 用户支持与用户问题 订阅 取消订阅 邮件存档
开发 开发相关 (设计文档, Issues等等.) 订阅 取消订阅 邮件存档
Commits 所有与仓库相关的 commits 信息通知 订阅 取消订阅 邮件存档
Issues Issues 或者 PR 提交和代码Review 订阅 取消订阅 邮件存档