Zulip 实时事件系统(Events System)深度解析:从事件生成、长轮询投递到初始数据原子同步

发布时间:2026/9/13 2:49:10
Zulip 实时事件系统(Events System)深度解析:从事件生成、长轮询投递到初始数据原子同步 Zulip 实时事件系统Events System深度解析从事件生成、长轮询投递到初始数据原子同步【免费下载链接】zulipZulip server and web application. Open-source team chat that helps teams stay productive and focused.项目地址: https://gitcode.com/GitHub_Trending/zu/zulipZulip 的Events System实时推送与事件系统是其服务端到客户端的推送系统承担了团队聊天产品中一个客户端修改的数据如何实时同步到其他客户端这一核心职责。本文以 docs/subsystems/events-system.md 为主线结合仓库中 zerver/tornado/event_queue.py、zerver/lib/events.py、zerver/tornado/django_api.py 等核心源码完整讲解事件生成Generation、投递Delivery、UI 更新三大部分的设计与实现并深入剖析注册register接口的原子性初始数据获取算法、apply_events的自动化测试协议与事件 Schema 变更的向后兼容策略。读完本文你将掌握 Zulip 实时同步的完整调用链并具备在 Zulip 中新增一种事件类型并配套测试的能力。为什么需要一套专门的事件系统任何单页 Web 应用都需要回答一个客户端做出的变更如何同步到其他客户端这个问题。对于 Zulip 这样的聊天工具状态时刻在变化这个问题尤为关键。这里的客户端指需要接收 Zulip 数据更新的浏览器标签页、移动端 App 或 API Bot。最简单的例子是一个客户端发送了一条新消息其他客户端必须被通知才能显示这条消息。而一个完整的应用如 Zulip有几十种需要同步的数据类型——新建频道、用户改名或换头像、设置变更等。在 Zulip 中这些需要发送给其他客户端的更新统称为事件events。设计这类系统时有一个重要原则事件需要同步给每一个持有旧数据副本的客户端否则客户端会向用户展示过期数据。因此如果一个用户开两个浏览器窗口并发送消息该用户控制的所有客户端以及消息的所有接收方包括那两个浏览器窗口都会收到事件。严格来说不需要给触发变更的那个客户端发事件但 Zulip 选择了全量下发——这样客户端无需为自己触发的变更和别人触发的变更各写一套 UI 更新代码触发方只需复用与其他客户端完全相同的代码最多再加一点操作成功的通知即可。从架构上看一个成功的实时同步系统需要三部分Generation生成数据发生变化时生成事件并确定每个事件应发给哪些用户。Delivery投递高效地把事件投递给感兴趣的客户端理想情况下做到恰好一次exactly-once。UI updates更新客户端收到事件后更新界面。React、Vue 这类响应式 JavaScript 库可以帮助简化第三部分但生成与投递没有成熟的标准系统可复用Zulip 必须自己构建。本文接下来就聚焦这两部分如何在 Zulip 中以可扩展、正确且可预测的方式实现。事件生成系统GenerationZulip 的生成系统围绕一个 Python 函数send_event_on_commit(realm, event, users)构建其真实实现位于 zerver/tornado/django_api.py。它接收三个参数realm用于分片sharding决定事件应投递到哪个 Tornado 端口/进程event事件数据结构本质就是一个 Python 字典type键始终存在其余键随具体事件类型而定users应接收事件的用户 ID 列表。在消息投递等特殊场景下users会是一组字典把用户 ID 映射到用户相关数据例如该用户是否在消息中被提及。传入send_event_on_commit的数据会被简单地序列化为 JSON放入名为notify_tornado的 RabbitMQ 队列等待投递系统消费。从源码看send_event_on_commit的核心是transaction.on_commit(lambda: send_event_rollback_unsafe(realm, event, users))——它必须在正在修改状态的那个数据库事务内部调用从而保证只有事务成功提交后才发送事件事务回滚则事件不发。send_event_rollback_unsafe则按 realm 的 Tornado 端口对用户进行分组get_realm_tornado_ports/get_user_id_tornado_port再通过queue_json_publish_rollback_unsafe发布到对应端口的notify_tornado队列。源码注释还约定这类函数只能从zerver/actions/*.py调用便于集中查找事件生成代码每个调用点都应由test_events.py中的测试覆盖并在zerver/lib/event_schema.py中校验 Schema。通常情况下users列表是以下三种之一单个用户例如用户级设置变更整个 realm 的所有人例如组织级设置变更、新增 realm 表情会收到某条消息的所有人例如消息、表情回应、消息编辑等即某频道的订阅者或某私信会话的参与者。选择正确的用户 ID 列表是调用方的责任。如果错误地把包含私信内容的事件发给整个组织会造成安全问题反之如果事件没发给足够的客户端就会出现用户可见的实时同步 Bug。事件生成过程中最困难的部分是定义一致的事件字典要清晰、可读、对各类客户端都有用并方便开发者维护。事件投递系统Delivery基于 Tornado 的长轮询架构Zulip 的事件投递实时推送系统基于Tornado——它非常适合处理大量保持打开的请求详情可参考 架构总览。整个系统约 2000 行代码集中在zerver/tornado/目录主体是 zerver/tornado/event_queue.py。投递机制采用长轮询long-polling客户端发起GET /json/events请求服务器在有事件可投递之前不响应这个请求。这种方式相当高效且兼容性好相比 WebSocketWebSocket 存在逐渐减少但并非为零的客户端兼容性问题。对每个已连接的客户端事件队列服务器event queue server维护一个事件队列event queue队列里存放该投递给这个客户端、但客户端尚未确认acknowledge的事件。忽略错误处理的细节协议非常简单客户端发起GET /json/events服务器检查队列中是否有事件有立即把事件作为响应返回没有把该队列记录为有等待中的客户端代码中常称为handler。当服务器从notify_tornadoRabbitMQ 队列拉出事件时就把它投递给目标用户关联的每个队列若队列有等待中的客户端中断长轮询向等待中的请求返回 HTTP 响应若没有等待中的客户端直接把事件压入队列。注册与轮询协议客户端启动时会先调用POST /json/register服务器为其创建一个新的事件队列返回queue_id以及初始的last_event_id可选地还可以顺带拉取初始数据以节省一个 RTT 并避免竞态详见下文初始数据获取。注册完成后客户端只需进入无限循环用这两个参数反复调用GET /json/events每次处理完事件后更新last_event_id以确认已收到Python API 绑定中的call_on_each_event是完整的示例实现。在处理每个GET /json/events请求时队列服务器可以安全地删除事件 ID 小于等于客户端last_event_id的事件事件 ID 只是该队列收到事件的计数器。last_event_id参数在无网络故障时并非必需但它是实现exactly-once 投递的关键如果没有它队列服务器只能在尝试发送事件时就删除事件一旦那次 HTTP 响应因 TCP 网络故障没有送达客户端事件就永久丢失了。心跳、垃圾回收与持久化队列服务器是超高流量系统至少为每一条投递给每个 Zulip 客户端的消息处理一次请求。此外为绕过低质量 NAT 服务器杀死空闲超过 60 秒的 HTTP 连接的问题队列服务器还会在无其他事件到达时至少每 45 秒向每个队列发送一个心跳事件。源码中的常量可佐证HEARTBEAT_MIN_FREQ_SECS 45心跳事件即dict(typeheartbeat)见create_heartbeat_eventevent_queue.py。为避免内存等资源泄漏队列在客户端默认闲置 10 分钟后会被垃圾回收GC 每分钟扫描一次见DEFAULT_EVENT_QUEUE_TIMEOUT_SECS 60 * 10与EVENT_QUEUE_GC_FREQ_MSECSevent_queue.py其假设是客户端大概率已断网或不存在。客户端重新回来时会收到 queue not found 错误其处理方式就是重启客户端/刷新浏览器像启动时一样重新拉取初始数据。由于客户端反正要实现启动流程这套方案对客户端几乎不增加复杂度。一个额外的好处是即使队列服务器队列保存在内存中崩溃丢失数据客户端也能自动恢复就像短暂断网一样仍需防范 DoS 风险。垃圾回收系统还带有钩子对 通知系统 的实现很重要。值得注意的是事件队列服务器被设计为将事件队列保存到磁盘并在重启时重新加载对应ClientDescriptor的to_dict/from_dict序列化机制见 event_queue.py并小心地捕获异常因此此类崩溃非常罕见但设计上即使发生也不会留下损坏的过期客户端。Tornado 侧的事件处理从 RabbitMQ 队列消费到事件后核心分发逻辑在process_event对每个目标用户 ID取出该用户的全部客户端描述符ClientDescriptor凡是accepts_event(event)为真的客户端都调用client.add_event(event)。process_event的调用链可通过zerver/tornado/views.py中的get_eventszerver/tornado/views.py与fetch_eventsevent_queue.py继续深入。ClientDescriptor保存了该客户端注册时的各项能力与过滤条件event_types、narrow消息范围过滤经build_narrow_predicate编译为谓词、bulk_message_deletion、stream_typing_notifications、simplified_presence_events等accepts_event正是基于这些字段决定是否把事件投给该客户端。初始数据获取与原子性保证客户端启动时通常想从服务器拿到两样东西各类数据的当前状态initial state当前设置、组织用户集合用于输入联想、频道、消息等对这些数据的更新订阅即一个事件队列。理想情况是这两者原子地获取假设其他用户此时改了名字那么要么改名发生在拉取之后初始状态里是旧名字队列里会有一个改名事件要么发生在之前初始状态是新名字队列里没有改名事件。绝不出现初始状态是旧名字、队列里也没有改名事件这种数据永远对不上的情况。实现这种原子性可以让 N 个 Zulip 客户端免于处理大量罕见且难以复现的竞态条件——只需在 Zulip 服务器端把这件事一次性做对。技术上这很有挑战拉取 Zulip 这种复杂应用的初始状态可能要执行几十次数据库、缓存查询耗时 100ms 以上几乎不可能把这些查询原子地完成。Zulip 的解决办法是用非原子的子过程组合出原子结果。逻辑位于 zerver/views/events_register.py 与 zerver/lib/events.pyregisterAPI 请求由 Django 直接处理流程如下Django 向 Tornado 发起 HTTP 请求要求创建一个新的事件队列并记下其queue_idDjango 非原子地从各个数据源执行各种数据库/缓存查询拉取数据见fetch_initial_state_dataDjango 第二次向 Tornado 发起 HTTP 请求取回自队列创建以来新增到该队列的所有事件最后 Django 把这些事件应用到拉取的初始状态上见apply_events。例如对改名事件在realm_user数据结构中找到该用户并更新为新名字。fetch_initial_state_data有大量参数见 events.py包括event_types为None时拉取支撑 Web 端page_params与/api/v1/register的核心数据指定时只拉取子集、client_gravatar、slim_presence、include_subscribers、archived_channels等并把zulip_version、zulip_feature_level等版本信息无条件放入 state。apply_events则遍历事件先按fetch_event_types过滤避免把未订阅类型的事件应用到无关状态上再逐个交给apply_event处理如 message 事件会更新state[max_message_id]。测试为什么需要apply_events协议上述设计达成了所有目标代价是必须写一个正确的apply_events函数。这个函数很难写对因为它处理的场景竞态条件在手工测试中几乎从不出现。幸运的是Zulip 的自动化后端测试有一套专门的测试协议。测试总览当你完全确信某个 action 函数在正常操作下工作正确通常意味着为对应的 POST/GET 操作写了一个全栈测试就可以在test_events.py中写测试了。一个test_events测试的实际代码可以非常简洁def test_default_streams_events(self) - None: stream get_stream(Scotland, self.user_profile.realm) events self.verify_action(lambda: do_add_default_stream(stream)) check_default_streams(events[0], events[0]) # (some details omitted)真正的技巧在于调试这些测试。上述示例做了三件事准备数据get_stream用verify_action包装一个 action 函数do_add_default_stream用 schema 检查器校验数据check_default_streams。test_events.py文件位于 zerver/tests/test_events.py。verify_action所有与apply_events相关的重活都发生在verify_action调用里它是test_events.py中BaseAction类的测试辅助方法。verify_action通过模拟可能的竞态条件来验证apply_events逻辑在某个 action 函数语境下是否正确。用上面的例子说把do_add_default_stream产生的事件通过apply_events应用到一个过期的状态副本上得到的结果应当与先执行 action 再拉取一份全新状态完全相同。具体来说verify_action依次执行调用fetch_initial_state_data获取当前状态调用 action 函数如do_add_default_stream捕获 action 函数产生的事件检查这些事件已在 OpenAPI 文档定义于zerver/openapi/zulip.yaml中登记调用apply_events(state, events)得到混合状态hybrid state再次调用fetch_initial_state_data得到正常状态normal state对比两者。如果apply_events逻辑一次写对两个状态完全一致verify_action通过并返回 action 产生的事件。通常你第一次会写错导致verify_action失败——它会打印混合状态与正常状态的 diff 帮助你调试。遇到这种 diff 可能是一场有挑战的调试建议重读本文档理解apply_events的设计动机阅读verify_action自身代码必要时在聊天中求助。verify_action只有一个必填参数即 action 函数通常用 lambda 表达以便传参events self.verify_action(lambda: do_add_default_stream(stream))它还有几个值得注意的可选参数state_change_expected如果 action 确实不引起状态变化例如输入中提示 typing notifications这类事件是临时的必须设为False否则verify_action会抱怨测试没有真正锻炼apply_events逻辑num_events告诉verify_action该 action 之后hamlet用户会收到几个事件默认 1client_gravatar、slim_presence等参数会被透传给fetch_initial_state_data对相关 action两个布尔值最好都测一遍。高级用法请直接阅读BaseAction在 zerver/tests/test_events.py 中的代码。Schema 检查test_events.py系统有两种形式的 Schema 检查。第一种是确认你已更新 GET /events API 文档 来记录新事件格式方便 Zulip 移动端、终端 App 及其他 API 客户端的开发者细节见 API 文档。第二种是test_events内部更细粒度的检查验证这个特定测试产生了预期的事件序列。看示例测试的最后一行# ... events self.verify_action(lambda: do_add_default_stream(stream)) check_default_streams(events[0], events[0])verify_action会返回 action 实际产生的事件test_events的测试纪律要求验证事件格式可预测。理想情况是测试事件与期望数据完全一致但由于数据库 id 等不可预测因素只能验证事件的 Schema——用check_default_streams这类 schema 检查器校验数据类型。如果要创建新的事件格式就得在event_schema.pyzerver/lib/event_schema.py中自己写 schema 检查器以下是与示例对应的代码default_streams_event event_dict_type( required_keys[ (type, Equals(default_streams)), (default_streams, ListType(DictType(basic_stream_fields))), ] ) check_default_streams make_checker(default_streams_event)注意basic_stream_fields未在文档中列出。理解如何编写 schema 检查器的最佳方式是阅读event_schema.py文件顶部有一大段注释然后可以快速浏览其余部分学习模式。创建一个新事件的 schema 检查器不仅让test_events测试更严格还允许其他工具复用同一检查器去校验 node 测试 fixtures 与 OpenAPI 文档中的事件格式。Node 端测试完成后端测试后还要在web/tests/lib/events.cjs添加一个示例事件在web/tests/dispatch.test.cjs中为server_events_dispatch.js的对应事件分发逻辑添加测试该文件已有约 140 处分发用例并用tools/check-schemas把示例事件与上面声明的两版 schema 做对照验证。代码覆盖率最后还需要确保apply_events始终正确即 Zulip 能生成的每一种事件类型都有相关测试。可以手动运行test-backend --coverage BaseAction然后检查所有send_event_on_commit调用点都被覆盖。未来计划用自动化手段直接通过检查覆盖率数据来验证这一点。page_params在 Zulip Web 应用中registerAPI 返回的数据通过page_params参数在页面上可用。消息初始数据协议的一个例外一个例外是真正的消息。因为 Zulip 客户端通常在站点其余部分加载完成之后用单独的 AJAX 调用拉取消息所以消息不需要包含在初始状态数据里。为正确起见客户端需要负责丢弃那些对应消息客户端尚未拉取的事件。相关机制还可参考 发送消息。Schema 变更与向后兼容当改变发送进 Tornado 的事件格式时必须正确处理向后兼容新增事件类型或给现有事件类型加字段只需在 API 文档 中仔细记录变更务必提升API_FEATURE_LEVEL并在更新的GET /eventsAPI 文档中加入**Changes**条目同时建议给移动端与终端项目开 issue 通知。如果 Web 端应在浏览器刷新时收到新事件类型的初始状态把它加进web/src/server_event_types.ts的FETCH_EVENT_TYPESWeb 端会以fetch_event_types参数传给/register。修改可能干扰现有客户端解析逻辑的字段如改变现有字段的类型/含义、删除字段需要非常谨慎因为 Zulip 支持旧客户端连接新服务器。具体政策见 发布生命周期。技术方案是为新的格式添加client_capabilities标志对未声明支持新能力的客户端继续发送旧格式数据。bulk_message_deletion就是一个很好的参考范例几年后再把该能力设为必需并移除旧路径。大多数事件类型Tornado 只是透明透传event_queue.py无需改动。但若改变了 Tornado 代码自身使用的数据格式例如重命名message事件中的presence_idle_user_ids字段必须小心因为升级 Tornado 时队列里可能还存有升级前的事件。因此必须在event_queue.py中写逻辑把旧格式翻译成新格式否则升级到相关 commit 时 Tornado 可能崩溃。这类逻辑应集中在from_dict函数用于事件队列格式变更和client_capabilities条件分支中例如process_deletion_event对旧格式客户端的逐条删除事件拆分。与client_capabilities无关的兼容代码应加# TODO/compatibility: ...注释说明何时可以安全删除主版本发布时会 grep 这些注释。源码中大量此类注释可佐证例如process_message_update_event中pm_mention_push_disabled_user_ids到dm_mention_push_disabled_user_ids的重命名翻译代码event_queue.py。Schema 变更是敏感操作与数据库 Schema 变更一样必须做认真的手工测试。例如在测试服务器上运行移动端 App 验证它正确处理新事件或安排新的 Tornado 代码真正处理一个升级前的事件并通过浏览器控制台确认输出。实战如何新增一种事件类型综合文档与源码在 Zulip 中新增一种事件类型的完整路径可归纳为生成在zerver/actions/*.py的 action 函数中于修改状态的数据库事务内部调用send_event_on_commit(realm, event, users)event字典必须包含type键确保users列表精确覆盖应收到事件的客户端。服务端 Schema在zerver/lib/event_schema.py中为该事件编写event_dict_type风格的 schema 检查器。后端测试在zerver/tests/test_events.py中为 action 编写verify_action测试并用 schema 检查器校验verify_action返回的事件用test-backend --coverage BaseAction确认send_event_on_commit调用点被覆盖。API 文档更新GET /events的 OpenAPI 文档 与zerver/openapi/zulip.yaml必要时提升API_FEATURE_LEVEL若 Web 端刷新需要初始状态更新web/src/server_event_types.ts的FETCH_EVENT_TYPES。Node 端在web/tests/lib/events.cjs加示例事件在web/tests/dispatch.test.cjs中为server_events_dispatch.js补测试并用tools/check-schemas校验。兼容性若涉及旧客户端或升级场景在event_queue.py的from_dict/client_capabilities分支中实现格式翻译并以# TODO/compatibility:注释标注清理时机。进一步阅读新应用功能教程一个完整功能如何使用该事件系统的端到端示例架构总览Tornado 在整体架构中的位置发送消息消息事件流的专门文档通知系统依赖队列垃圾回收钩子的移动/邮件通知实现API 文档GET /events等事件相关 API 的 OpenAPI 规范发布生命周期事件格式变更的兼容性政策【免费下载链接】zulipZulip server and web application. Open-source team chat that helps teams stay productive and focused.项目地址: https://gitcode.com/GitHub_Trending/zu/zulip创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询