Skip to content

S4 ​

标签
分布式/计算
字数
15863 字
阅读时间
61 分钟

S4(Simple Scalable Streaming System)是 Yahoo! 的通用分布式流计算平台,用来处理连续、无界的数据流。它的自我描述里有一个词值得停一下:部分容错(partially fault-tolerant)—— 它界定了后面一整套取舍的边界。

基本形态:带键的数据事件按亲和性路由到处理单元(Processing Element,PE),PE 消费事件后做两件事之一或两件都做 —— 发出一个或多个可能被其他 PE 消费的事件,或者发布结果。架构接近 Actors 模型,因此天然带有封装与位置透明的语义,让应用可以大规模并发、同时对开发者只暴露一个简单接口。

这篇值得读的地方不在"它做对了什么",而在于它把**"为什么不能拿批处理平台改造成流处理"讲得很清楚,并且公开列出自己接受的两条假设**。

为什么不能把 Hadoop 改造成流处理 ​

背景是搜索广告。主流搜索引擎先给自然结果,再用文本广告补位,按每次点击计费收费。为了让最相关的广告出现在最好的位置,需要算法按上下文动态估计点击概率 —— 上下文包括用户偏好、地理位置、历史查询、历史点击等。一个大型搜索引擎每秒处理数千次查询,每页又有若干条广告。为了处理用户反馈,S4 被做成一个低延迟、可扩展的流处理引擎。

研究与生产对它的要求不一样,这一点决定了设计:

场景要求
研究能很快把算法部署到线上,从而用真实流量测在线算法,开销与运维负担最小
生产可扩展(加服务器就能提吞吐,尽量不费力)与高可用(有系统故障时无人干预地持续运行)

他们考虑过扩展开源的 Hadoop 来支持无界流,但很快发现 Hadoop 是高度针对批处理优化的:

MapReduce 系统在静态数据上调度批作业。而流计算的范式是:事件以我们无法控制的速率流入。处理系统要么跟上事件速率,要么优雅降级 —— 也就是淘汰事件,通常叫 load shedding。

由此得出一个明确判断:流式范式注定需要一套与批处理很不一样的架构;试图做一个同时服务批与流的通用平台,会得到一个高度复杂、可能在两头都不最优的系统。

当时的外部状况也印证了这个空缺:MapReduce 生态随 Hadoop 一类开源项目繁荣,但通用分布式流计算软件没有同样的趋势 —— 已有的项目与商用引擎仍限于高度专门化的应用。而实时搜索、高频交易、社交网络这些新应用在把传统数据处理系统推到极限。

对"定长分段 + MapReduce"这条流式策略的批评 ​

很多真实系统实现流式策略的方式是把输入切成定长段、交给 MapReduce 处理。这篇指出的代价很具体:

  • 延迟与段长成正比,还要加上分段本身、以及启动处理作业的开销;
  • 段小 → 延迟低,但开销上升,而且跨段依赖变得更难管(某一段可能需要前段的信息);
  • 段大 → 延迟上升。

最优段长取决于应用 —— 这句话意味着没有普适参数可调。这里的判断很直白:与其削足适履,不如另找一套简单、能实时处理数据流的编程范式。

六条设计目标,与两条公开承认的局限 ​

设计目标:

  1. 为处理数据流提供简单的编程接口;
  2. 设计一个高可用、能用商用硬件扩展的集群;
  3. 靠每个处理节点的本地内存把延迟压下来,避开磁盘 I/O 瓶颈;
  4. 去中心、对称的架构 —— 所有节点功能与职责相同,没有任何承担特殊职责的中心节点(大幅简化部署与维护);
  5. 可插拔,让设计尽量通用与可定制;
  6. 对科研友好 —— 易编程、可调。

两条假设(写得很直白,也是"部分容错"的全部含义):

S4 自述的两条边界

  • 有损故障切换可以接受。 服务器故障时,进程会被自动移到备用服务器;但进程状态存在本地内存里,交接过程中就丢了,状态靠输入流重新生成。下游系统必须优雅降级。
  • 运行中的集群不会被加减节点。

这两条对他们大多数应用可以接受,不接受的情形留作后续工作。这两条在后面两处会具体现身:CTR 那节讲"PN 挂了就退回长期估计",而"缺动态负载均衡与在线 PE 迁移"被列进未来工作。

这六条目标之间有一处张力,值得点出来:第 3 条要求"靠本地内存把延迟压下来、避开磁盘 I/O",而第 2 条要求"高可用、能用商用硬件扩展"。本地内存意味着状态不在持久存储里,而高可用要求节点挂掉之后还能继续 —— 两者的缝合处就是那两条假设:用"有损"把这个缝填上。所以"部分容错"这个自我描述正落在那个缝上,它是这组目标在逻辑上唯一自洽的收尾。

另一处张力在第 4 条与第 6 条之间:去中心对称意味着"没有需要单独维护的组件",而"对科研友好、可调"意味着"有很多旋钮"。旋钮多的地方往往会长出中心 —— 配置总得有个地方统一管。它把这部分交给了外部的 ZooKeeper,用"中心在系统之外"换来了"系统内没有中心"。这也解释了为什么"对称"这个词必须配上"借助 ZooKeeper"才完整。

为什么选 Actors 模型 ​

两个目标 —— 在商用硬件上分布式运行、以及避免跨集群共享内存 —— 把他们推向 Actors 模型。

在 S4 里:

  • 计算由 Processing Element(PE) 完成;
  • 消息以 data event 的形式在 PE 之间传递;
  • 每个 PE 的状态对其他 PE 不可访问;
  • 事件的生产与消费是 PE 之间唯一的交互方式。

框架负责把事件路由到合适的 PE、以及创建新的 PE 实例 —— 这两件事带来了封装与位置透明。

与 IBM 的 Stream Processing Core(SPC) 的对照值得记一句:两者都面向大数据、都能用用户定义算子从连续流里挖信息,但 SPC 源自订阅模型,而 S4 源自 MapReduce 与 Actors 模型的组合。他们认为 S4 因为对称性而简单得多 —— 所有节点相同、没有中心控制 —— 实现办法是借助 ZooKeeper 这个可被数据中心里多个系统共享的集群管理服务。

数据模型与三个抽象 ​

一个流是一串形如 (K,A) 的元素("事件"),K 与 A 分别是元组值的键与属性。

一个完整的例子:词频 top-K ​

输入事件包含一段带英文引文(quotation)的文档,任务是以最小延迟持续产出"全部文档中出现频率最高的 K 个词"的有序列表。

三处设计意图值得说清:

  • Quote 事件不带键,所以由无键 PE QuoteSplitterPE 处理全部 Quote 事件;它对文档里每个唯一词分配一个计数,发出以 word 为键的 WordEvent。
  • WordCountPE 是按 word 的每个取值实例化的 —— word="said" 的那个实例只收 word="said" 的事件;实例不存在时新建一个。
  • 为什么要多个 SortPE:每个 WordCountPE 的 sortID 是 [1,n] 里的一个随机整数(n 是期望的 SortPE 数),一旦选定就终生使用。这样把排序负载摊到多个节点或处理器上。每个 SortPE 各维护一份局部 top-K,周期性把自己的部分列表发给唯一的 MergePE(用一个约定的键),由它合并并输出最新的权威 top-K 列表。

PE 的四个身份要素 ​

每个 PE 实例由四项唯一确定:

#要素
1功能 —— 由 PE 类与关联配置定义
2它消费的事件类型
3那些事件里被键化的属性
4被键属性的取值

一条硬约束:每个 PE 恰好消费"与它被键化的那个取值相对应"的事件,并可能产出输出事件。

PE 是按被键属性的每个取值实例化的,而这个实例化由平台完成:例子里的 WordCountPE 对输入里每个词各实例化一个 —— 事件里出现一个新词,S4 就为它建一个新 PE 实例。

无键 PE 是特殊一类:没有被键化属性、也没有取值,消费与它关联的那类事件的全部。它们通常用在 S4 集群的输入层,也就是"给事件赋键"的那一步。

内置了一批标准 PE(count、aggregate、join 等),很多任务只用配置文件就能定义、不需要额外编码;自定义 PE 用 S4 的开发工具写。

PE 对象的回收是一个必须面对的问题(唯一键很多时,PE 对象会不断累积)。最简单的办法是给每个 PE 对象设一个 TTL:指定时间内没有它的事件到达,就成为可移除的;内存回收时 PE 对象被删掉,此前的状态就丢了 —— 在词频例子里,那个词的计数会丢。设计者自己承认这个策略简单但不够省:要最大化服务质量(QoS),理想做法应当依据可用系统内存与该对象对整体性能的影响来移除,并设想让 PE 对象自己报告优先级或重要度 —— 而这部分逻辑是应用相关的,得由应用开发者实现。

Processing Node:PE 的逻辑宿主 ​

Processing Node(PN) 是 PE 的逻辑宿主,负责监听事件、对传入事件执行操作、在通信层协助下派发事件、以及发出输出事件。

  ┌─ Processing Node(PN) ────────────────────────────────┐
  │                                                        │
  │   事件 ──▶ event listener ──▶ ┌─ PEC(processing       │
  │                               │   element container)  │
  │                               │                        │
  │                               │  按合适顺序调用 ──▶ PE 实例 1
  │                               │                  ──▶ PE 实例 2
  │                               └────────────────────────┘
  │                                                        │
  │   PE prototype:只有身份前三项(功能/事件类型/被键属性),
  │   取值未定。PN 遇到新的被键属性取值时,让它克隆出取值为 V 的完整 PE
  └────────────────────────────────────────────────────────┘
  • S4 按"该事件里所有已知被键属性的取值"的哈希函数,把每个事件路由到 PN;一个事件可能被路由到多个 PN。所有可能的键属性集合从集群配置里已知。
  • PN 里的 event listener 把传入事件交给 processing element container(PEC),由它按合适的顺序调用合适的 PE。
  • PE prototype 是一类特殊的 PE 对象:它具备身份的前三项(功能、事件类型、被键属性),但取值未赋值。它在初始化时被配置好,并且对任意取值 V 都能克隆出"同配置、取值 V"的完整 PE。这个操作由 PN 在遇到某个新的被键属性取值时触发一次。

由此落出一条结构性保证,它是整个系统能去中心的前提:

具有某个被键属性取值的所有事件,保证到达对应的那一个 PN,并被路由到它内部对应的 PE 实例。 每个有键 PE 恰好映射到一个 PN(按被键属性取值的哈希);而无键 PE 可以在每个 PN 上实例化。

"一个事件可能被路由到多个 PN"这句话值得展开,因为它决定了去中心架构怎么保证状态一致。路由的依据是该事件里所有已知被键属性的取值的哈希,而"所有可能的键属性集合从集群配置里已知" —— 这句是整条保证的落点:框架必须事先知道会有哪几种键属性,才能把路由算出来。所以新增一个键属性是一次集群配置变更,而它也不只是加一行代码。

两处结构性保证合起来,才让"PE 的状态只在本地、却依然一致"成立:

保证作用
某个取值的全部事件到达同一个 PN保证状态看得见全部相关输入
该 PN 内部路由到对应的那一个 PE 实例保证状态的作用域与键的取值一一对应

而"有键 PE 恰好映射到一个 PN"这条上有单点性:某个键取值的全部流量都落在同一个 PN 上,那个 PN 就是该键的单点。好处是状态一致性天然成立,代价是这个键的吞吐上限等于那个 PN 的吞吐上限。词频例子里这不是问题(每个词的流量小);换成"按国家聚合"这类键就只有少数取值、每个都承载巨大流量 —— 键的基数决定了这套架构能不能摊平。

无键 PE 的位置也因此清楚了:它"可以在每个 PN 上实例化",所以它是唯一不承担状态一致性的角色。这正是它只适合放在输入层"赋键"那一步的原因 —— 那一步要做的事(给无键事件分配一个键)天然不需要跨事件的状态,因此放在哪里都一样,也就不存在单点。

一处实现上的取舍值得记:"一个事件可能被路由到多个 PN"是按事件里存在的键属性算出来的,而不按下游 PE 的需要来算 —— 也就是说路由这一层不知道下游的拓扑。这带来一个正面的后果:给同一条事件流加一个消费者 PE,路由不需要改(多出来的那个 PN 本来就会收到这条事件)。代价是每个 PN 都要为它可能并不消费的事件付一次接收成本,而这个成本随"事件上带的键属性种类数"增长。键属性越多,广播面越宽 —— 这是"框架事先必须知道全部键属性"这条约束的另一重代价。

通信层 ​

通信层提供集群管理与到备用节点的自动故障切换,并把物理节点映射到逻辑节点;它自动检测硬件故障并更新这个映射。

  • 发送方只指定逻辑节点 —— 它不知道物理节点,也不知道逻辑节点因故障被重映射过;
  • API 提供多种语言绑定(Java、C++);遗留系统可以用这套 API 以轮询方式把输入事件发给集群节点,再由无键 PE 处理;
  • 可插拔地选择网络协议;事件可以带保证、也可以不带保证地发送 —— 控制消息可能要求保证送达,而数据可以不保证送达以最大化吞吐;
  • 用 ZooKeeper 在节点之间做协调。

配置管理 ​

设想里由人操作来建立与拆除 S4 任务集群,而物理节点到任务集群的分配由 ZooKeeper 协调:一部分活跃节点被分给特定任务,其余空闲节点留在池子里按需使用(故障切换、动态负载均衡)。一条细节值得注意:一个空闲节点可以同时被注册为多个活跃节点的备用,而这些活跃节点可能属于不同的任务。

事件的键怎么设计 ​

这套系统里,键不是一个路由细节,它是状态的坐标。三件事同时由键决定:事件被路由到哪个 PN、会实例化出多少个 PE、以及状态的作用域有多大。因为"具有某个被键属性取值的所有事件,保证到达对应的那一个 PN",所以凡是有状态的聚合,都必须靠键来保证它看得见全部相关事件。

由此得到一条推论:改键等于改语义。 把词频例子的键从 word 换成"word + 文档 id",统计的就不再是全局词频,而是每篇文档的词频 —— 计算逻辑一行没改,结果完全不同。

CTR 那个例子把这套推理用得最干净,值得单独看一遍:click 与 serve 的载荷不对称 —— click 只带点击本身的信息加上它关联的 serve ID,而 serve 带 serve ID、查询、用户、广告。要在"查询-广告"粒度算 CTR,就得用"查询 id + 广告 id"组成的键;但 click 事件里根本没有这两个 id。于是必须先按 serve ID 做一次 join、把查询与广告信息补上,再换成"查询-广告"这个键做二次路由。两条事件流的键不一致时,先 join 再换键 —— 这是流处理里"换键"的标准形态,代价是一次额外的跳,而这次跳在网络上是实打实的。

键选错的两种方向,代价不一样:

键的选择后果
取到实体粒度(每个词、每个用户一个值)PE 实例数等于实体数;唯一键多时 PE 对象不断累积,靠 TTL 回收,而回收即丢状态
把大量实体合并到一起并行度不足,一个 PN 要扛住过多流量
无键全部事件进同一个 PE;无键 PE 只在输入层当"赋键者"时是合理的

一段判据:键的基数(不同取值的个数)应当与期望的并行度同量级,而键的实体粒度应当与状态的作用域一致。 这两件事经常冲突 —— 词频例子里 word 恰好两样都合适(要按词统计、也要按词分散),所以那个例子读起来很顺;CTR 那里就得多一步 join。

键的第三种用法在 SortPE 上:WordCountPE 的 sortID 取 [1,n] 里的一个随机整数,一旦选定终生使用。这个键与业务毫无关系,它的作用纯粹是把排序负载摊到多个节点上 —— 键在这里从"状态作用域"变成了"分片算子"。同一个机制承担两种用途,代价是两者会互相牵制:为了让分片均匀而选随机键,就同时放弃了"按业务维度聚合"的能力。

换键这件事在流处理里是常态,代价的形态却很少被写清楚,值得补一张对照:

换键发生的位置什么时候必须做代价
无键 → 有键(输入层)事件原本没有可路由的键几乎为零:这一步本来就该在入口做完
有键 → 另一个键(join 之后)两条流的键不同一,但要按同一个维度聚合一次额外的网络跳,且两条流都要重哈希
有键 → 无键需要全局汇总(例如产出唯一的最终结果)收敛到一个 PE ⇒ 重新变成单点

CTR 那个例子走的是第二行;SortPE 走的是第三行的反向 —— 把负载从"一个"摊到 n 个局部 top-K,再由 MergePE 收回来。第三行是必要的,但它的规模必须被限制:MergePE 只有一个,所以它前面那批 SortPE 的输出频率决定了它扛多少流量。例子里用"周期性发送"压这个量,这是这个拓扑能成立的前提,而不只是一处优化。

一套键设计的检查表,可以在写完 PE 之后对着问一遍:

  1. 状态的作用域与键的粒度一致吗? 想统计全局量就要全局的键,想统计分组量就要分组的键 —— 两者混用会出现"看起来对、数字偏小"的结果,因为状态被切成了多份,每份都不完整。
  2. 键的基数与期望并行度同量级吗? 基数远小于工作单元数 ⇒ 大部分节点闲着;基数远大于工作单元数 ⇒ 大量 PE 实例等着被 TTL 回收。
  3. 有没有"必须收敛到一个点"的步骤? 有,就把它放在拓扑末端,并让它前面的输出频率尽量低。
  4. 换键在哪一步发生?多付的那一跳可接受吗?
  5. 键的选择会不会让某个取值过热? 取值的流量分布呈重尾时(少数取值承担绝大部分流量),按值取模的分片会直接映射成负载不均。

最后一条最容易被忽略:这套路由按取值做,而不按取值的频率做 —— 它天然不处理倾斜,而倾斜在真实流量里是常态。这也回过头解释了为什么 SortPE 要用一个随机整数当键:随机数是为了让取值分布均匀,从而绕开倾斜;这是一个用"放弃业务聚合"换"均匀分片"的明确取舍。

PE 的生命周期:创建、复用与回收 ​

PE 由框架按需实例化,启动时并不建全:PN 在遇到一个新的被键属性取值时,让 prototype 克隆出取值为 V 的完整 PE。也就是说,某个键的第一次出现才带来创建开销,之后是复用。

prototype 与实例的分工值得分清:prototype 具备身份的前三项(功能、事件类型、被键属性)而取值未赋值,它在初始化时被配置好、能对任意取值克隆;触发克隆的是 PN,而不是 PE 自己。这条分工让"实例化"变成框架的职责,应用只管写好一个类。

回收是这套机制必须面对的另一半。 唯一键很多时 PE 对象会不断累积,所以必须能移除。最简单的办法是给每个 PE 对象设一个 TTL:指定时间内没有它的事件到达,就成为可移除的;内存回收时 PE 对象被删掉,此前的状态一并消失 —— 词频例子里,那个词的计数就没了。

设计者自己承认这个策略简单但不够省,并给出了理想形态:依据可用系统内存与该对象对整体性能的影响来移除,甚至设想让 PE 对象自己报告优先级或重要度 —— 而这部分逻辑是应用相关的,只能由应用开发者实现。"该丢什么"这个问题在这套系统里被推给了应用层,这是它"可插拔、通用"取向的直接代价。

一条实用判据:TTL 的长度应该由"状态重建的代价"决定,而不是由"内存够不够"决定 —— 丢失的真正损失是重建代价,不是占用的那点内存。而 TTL 是单一时间参数,它只能表达"多久没来",表达不了"丢了多贵"。这正是"让对象自己报告优先级"这个设想要解决的问题:要按价值回收,就必须有一个能表达价值的接口,而时间不是价值。

最后一处呼应前文:TTL 回收与故障丢失是同一件事的两种形态 —— 前者是有意的、由框架按时间触发;后者是被动的、由故障触发。两者都让那部分状态回到零,而下游都必须能承受这件事。这也解释了为什么"有损故障切换可以接受"这条假设能成立:系统在正常运行时就一直在做类似的丢弃,故障只是让它发生得更集中。

编程模型 ​

目标是写通用、可复用、可配置的 PE。PE 用 Java 写,用 Spring 框架组装成应用。

开发者实现两个主要钩子:

钩子何时调用做什么
processEvent()每个订阅类型的传入事件输入处理逻辑,典型是更新 PE 内部状态
output()(可选)按固定时间间隔 t,或每收到 n 个输入事件(n=1 即每个事件)输出机制,典型是把 PE 内部状态发布到外部系统

另外可定义若干状态变量。一个具体例子是 QueryCounterPE:订阅 QueryEvent,从头累计每个查询的计数,间歇性写出去。它的配置里

  • keys 属性写 QueryEvent queryString —— 即"订阅 QueryEvent、按 queryString 属性键化";
  • 绑定到 externalPersister(可以是某个数据服务系统的抽象);
  • outputFrequencyByTimeBoundary 为 600 —— 也就是每 10 分钟调一次 output()。

参数与可调项 ​

这套系统的配置面比前几篇都大,因为它把"平台代码与配置"也做成了部署期可变的东西。先看工具集,因为配置入口都在这些命令上:

命令作用
s4 newCluster定义逻辑集群(名字、分区数、初始端口)
s4 node启动节点
s4 s4r / s4 deploy打包应用、部署应用
s4 zkServer起一个测试用 ZooKeeper;加 -t 会顺带配好两个集群
s4 status查看某个 ZooKeeper 集群下所有 S4 集群的状态

配置分三层,各自的位置不同:

① 集群层 —— 配置存在 ZooKeeper 里。 启动节点之前必须先定义:集群名、分区数(约等于任务数)、以及初始端口。端口的约定很具体:每个节点开一个端口、从初始端口单调递增 —— "10 个节点、初始端口 12000"就用 12000 到 12009,这些端口用于节点间通信。所以初始端口必须预留一整段空闲,而这段范围由节点数决定。命令是 ./s4 newCluster -c=cluster1 -nbTasks=2 -flp=12000。

② 节点层 —— 默认全部可覆盖。 启动一个节点其实只需要两项:ZooKeeper 连接串(默认 localhost:2181)与逻辑集群名(./s4 node -c=cluster1 -zk=host.domain.com)。默认参数来自 classpath 上的三个文件:default.s4.base.properties、default.s4.comm.properties、default.s4.core.properties。覆盖有两种写法:命令行内联 -p=param1=value1,param2=value2,或者用 @ 引用一个配置文件(./s4 node @/path/to/config/file)。

③ 应用层 —— 参数在部署阶段给。 应用类、代码来源、用哪些模块、模块来源、以及字符串配置参数,全部在部署时指定;应用侧靠 @Inject @Named('myParam') 拿到它们。

模块化是这套配置体系的骨架:用 Guice 做依赖注入,一个节点由四个模块组成 —— base(怎么连集群管理器、怎么下代码)、comm(通信协议、事件监听器与发送器)、core(部署机制、序列化机制)、application。两处细节值得记:comm 模块有默认参数文件,而 core 模块没有默认参数;自定义模块类必须提供无参构造函数(框架要能反射地把它们实例化)。替换模块用 -emc,例如让检查点落到文件系统:-emc=org.apache.s4.core.ft.FileSystemBackendCheckpointingModule;自定义模块的 jar 用 -mu / modulesURIs 指定来源。

按六要素摊开两处:

-p 的适用范围(默认不限,但文档明确划了边界) —— 官方文档写明:它只该注入与节点相关的参数,应用参数应当在应用配置与部署阶段注入,并举了"指标日志配置"作为该放在这里的例子。所以 -p 不是万能入口。把应用参数塞进来会绕过部署配置这条路径,表现是同一个应用在不同环境里行为不一致,而排查时很容易归因到代码。

-nbTasks(分区数 ≈ 任务数) —— 它是集群级参数、存在 ZooKeeper 里,不是节点参数。这意味着改它等于改整个集群的切分方式,而不只是某个节点的行为。它与前面那条自述局限直接相关:任务总是按节点数切分,所以在早期形态里"加一个节点"就意味着"分区数变了",代价是大量数据重映射 —— 官方仓库里那条 README 把它量化成:4 个节点加 1 个,几乎全部 key 要重新映射,而理想情况只该移动 20%。一个"参数"变成一次全集群改造,是这种切分方式的直接后果。

把三层配置的位置关系再收一句:集群层的参数存在 ZooKeeper 里、对所有节点生效;节点层的参数来自本地文件或命令行、只影响那一个节点;应用层的参数在部署时下发、跟着应用走。三层边界不重叠,所以排查配置问题时先问"这个参数该在哪一层",比直接去读配置文件快得多。

一处具体到默认值的边界:ZooKeeper 连接串默认是 localhost:2181,而 s4 zkServer -t 正是为单机测试准备的(起一个 ZK 并顺带配好两个集群)。默认值服务于"能立刻跑起来",不服务于生产 —— 在生产部署里,它必须是第一个被改掉的参数。这也是三层配置里"节点层默认值"的典型形态:默认值填的是最省事的那个选择,而不是最合理的那个。

实测一:实时 CTR ​

点击是最有价值的用户行为之一 —— 它们对偏好与参与度给出即时反馈,可以用来把最热的内容放到更显眼的位置。按点击计费时,发布方、代理与广告主按点击数结算。CTR 就是点击数除以展示数;在有足够历史数据时,它是"用户会点击"概率的一个良好估计。

正因为点击对个性化与排序如此值钱,它也是欺诈的目标。 点击欺诈可以用来操纵搜索引擎里的排名,通常由跑在远端机器上的恶意软件(bot)或机器集群(botnet)实施;另一种威胁是展示刷量 —— 由机器人发起的请求,其中一些未必恶意,但会污染 CTR 估计。

实时 CTR 的做法:以"每个查询-广告组合"为单位测量,并用一组启发式规则过滤可疑的 serve 与 click。先明确两个术语:一次 serve 对应一次用户查询,被分配一个唯一 ID;该 serve 返回一个结果页,页面上产生的点击带上同一个唯一 ID。

键的设计在这里是核心,因为两个事件流的载荷不对称:

  • click 事件只含点击本身的信息,加上它所关联 serve 的 serve ID;
  • serve 事件含 serve ID、查询、用户、广告等。

所以:要在"查询-广告"粒度算 CTR,就得用"查询 id + 广告 id"组成的键来路由这两类事件。而如果 click 的载荷里不含查询与广告信息,就必须先按 serve ID 做一次 join,再用"查询-广告"当键路由。join 之后事件还要过一道 bot 过滤器,最后才把 serve 与 click 聚合出 CTR。

在线实验 ​

  • 跑在真实搜索流量的随机样本上;用户被随机分组,但按浏览器 cookie 的哈希固定,以保证体验一致;
  • 平均每天约 100 万次搜索,来自 25 万用户;跑了两周;观测到的峰值事件率是每秒 1600 个事件;
  • 集群16 台服务器,每台 4 个 32 位处理器、2 GB 内存。

窗口与延迟的取舍值得单独说:CTR 聚合在一个24 小时的滑动矩形窗口上,实现方式是把窗口切成每小时一个 slot,逐 slot 聚合点击与展示;整点时,把窗口内各 slot 的聚合相加、推给服务系统。

这个做法在内存使用上很省,代价在更新延迟上。 内存更多的话就可以维护更细的 slot —— 例如 5 分钟一个 —— 从而压低更新延迟。

系统提供的是短期 CTR 估计,再与长期估计合并使用。而这里就是"有损故障切换"落地的地方:

PN 故障时我们会丢掉那个节点的数据,被分到该节点的那部分"查询/广告"实例就没有短期估计了 —— 故障切换策略是退回到长期估计。

离线压力测试 ​

  • 8 台服务器,每台 4 个 64 位处理器、16 GB 内存;共 16 个 PN,每台跑 2 个;
  • 用搜索流量日志里的真实 click 与 serve 数据重建事件流,并把真实 CTR 当金标准来验精度;
  • 事件数据共 300 万条 serve 与 click。

结果 ​

在线部分:CTR 提升约 3%,收入没有损失 —— 主要来自很快识别并过滤掉低质量广告。

离线压力测试以逐步升高的速率灌数据,每个速率跑完都把系统估计的 CTR 与日志算出的真实 CTR 对比:

事件/秒CTR 相对误差数据率
20000.0%2.6 Mbps
36440.0%4.9 Mbps
72680.2%9.7 Mbps
104800.4%14.0 Mbps
124320.7%16.6 Mbps
149001.5%19.9 Mbps
160001.7%21.4 Mbps
200004.2%26.7 Mbps

系统在约 10 Mbps 处开始劣化,原因是S4 网格在这个速率下处理不过来事件流,从而发生事件丢失 —— 这正是开头那句"要么跟上速率、要么淘汰事件"的实际形态。

实测二:在线参数优化 ​

第二个应用是在线参数优化(OPO):用真实流量自动调优搜索广告系统的一个或多个参数。它省掉了人工调参与持续介入,同时用更少的时间搜索更大的参数空间。

闭环是:摄入目标系统发出的事件 → 度量性能 → 用自适应算法决定新参数 → 把新参数注回目标系统。这个闭环在原理上接近传统控制系统。

前提条件:目标系统(TS)的输出可表示成流,并且有可度量的性能 —— 形式是一个可配置的目标函数(OF)。系统留出两条随机分配的切片 slice1 与 slice2;要求 TS 能在两条切片上应用不同的参数值,且每条切片在输出流里可被识别。

三个功能组件:

  1. measurement —— 摄入两条切片流,度量各自的 OF 值。OF 按一个 slot 的时长来度量,slot 既可以按时间单位指定,也可以按输出流的事件计数指定;
  2. comparator —— 拿测量结果判断两条切片之间是否存在统计显著差异;判定显著就通知 optimizer。如果在指定数量的测量 slot 之后仍未发现差异,就把两条切片判为相等;
  3. optimizer —— 实现自适应策略。输入是参数影响的历史(含 comparator 的最新输出),为两条切片输出新的参数值 —— 这标志着一个新的"实验周期"开始。

在 S4 里,这三个组件就是三个 PE:

PE键
MeasurementPE按 sliceID 键化,slice1/slice2 各一个实例(要更多切片容易扩展)
ComparatorPE按 comparisonID 键化,映射到一对切片;显著性判定用配对测量的 dependent t-test,并配置了最少有效测量数
OptimizerPE同样按 slice1:slice2 键化;策略用改版的 Nelder-Mead(Amoeba)算法 —— 一种无梯度最小化算法

结果:在搜索广告系统的真实流量切片上跑,每片每天约 20 万用户,跑两周,优化一个已知影响显著的参数。目标函数是收入与用户体验的公式化表示。优化出的参数带来:收入提升 0.25%、点击产出(click yield)提升 1.4%。

底层依赖:ZooKeeper、通信与序列化 ​

"去中心、对称"这个设计标签需要补一句前提:它把中心挪到了系统之外。 S4 自己没有承担特殊职责的节点,但它依赖 ZooKeeper 这个外部协调服务。这个分工决定了它的部署形态与故障模型。

ZooKeeper 在这套系统里承担四件事:

职责具体内容
存放集群配置集群名、分区数、初始端口都在里面
节点注册与成员关系谁在线、谁属于哪个逻辑集群
物理节点到逻辑节点的映射故障后自动重映射
故障检测与切换物理节点故障时把工作交给备用节点

由此有一条结构性推论:对称的架构加中心化的协调,得到的是"每个节点都不特殊,但每个节点都依赖同一个第三方"。这也解释了它为什么强调 ZooKeeper 是"可被数据中心里多个系统共享"的服务 —— 若非共享,每个系统各起一套,对称换来的运维简化又会以另一种形式还回去。

通信层提供三样东西:集群管理、到备用节点的自动故障切换、以及物理节点到逻辑节点的映射(并自动检测硬件故障、更新这个映射)。三条性质值得逐一记住:

  • 发送方只指定逻辑节点 —— 它既不知道物理节点是谁,也不知道逻辑节点是否因为故障被重映射过。位置透明在通信层就落在这里;
  • 网络协议可插拔,而且事件可以带保证、也可以不带保证地发送 —— 官方文档点明了典型分工:控制消息可能要求保证送达,而数据可以不保证送达以最大化吞吐。这条区分让"部分容错"这个概念变得具体:它不只体现在故障时刻的取舍上,它同时体现在正常运行时的每一条消息上;
  • API 有 Java 与 C++ 绑定,而且遗留系统可以用轮询方式把输入事件发给集群节点,再由无键 PE 处理。这是它接入已有系统的入口 —— 不需要对方改成推模式。

序列化被放在 core 模块里,而且官方文档明说 core 模块没有默认参数 —— 也就是说序列化机制是必须显式选择的东西。它声明依赖的库里包括 kryo、jackson、asm 与 json,说明这一层留给用户的空间不小。

一处与相邻系统的对照:Pregel 把消息的序列化交给协议缓冲区,"每顶点一个值"这个限制的出口是值类型;S4 把序列化做成 core 模块里可替换的组件,出口是模块选择。两者回答的是同一个问题 —— 通用性与通信成本怎么同时要,只是一个靠类型系统、一个靠组件替换。

ZooKeeper 是这套系统里唯一的外部依赖,这个事实的分量值得单独称一下。 它不在 S4 的进程里,却决定了四件事(配置、成员、映射、故障检测)。由此有三个后果:

  • 它的可用性就成了 S4 集群的可用性上限。 协调服务不可用时,新的故障切换、节点加入、配置变更都无法完成 —— 已经在跑的 PE 还能继续处理事件,但系统失去了自愈能力;
  • 它的负载与 S4 的规模正相关。 每个节点的成员关系、每次故障导致的映射更新都要经过它 —— 这解释了它为什么必须"可被数据中心里多个系统共享":共享让它的运维成本被多个系统摊掉,同时也让它的容量被多个系统竞争;
  • 它的数据结构决定了 S4 的集群模型。 集群配置(名字、分区数、初始端口)直接存在 ZK 里,所以**"定义一个集群"本质上是在 ZK 里写一份数据,而不是启动一组进程**。这正是 s4 newCluster 与 s4 node 是两个命令的原因 —— 前者的产物是一份配置,后者的产物是进程。

通信层"事件可以带保证、也可以不带保证地发送"这一条,值得与"部分容错"放在一起看。 它说明这套系统的取舍不发生在故障那一刻,而是在每一条消息上就已经分好了等级:控制消息要保证送达(丢了会让状态机走错),数据可以不保证送达(丢了最多损失一点精度)。"部分容错"因此是一句关于"哪些东西丢了不要紧"的设计断言,而不只是一句关于故障概率的声明 —— 它决定了整个系统里所有的丢弃策略,包括 TTL 回收、load shedding 与故障切换。

最后一处依赖是语言绑定:Java 与 C++ 双绑定,加上"遗留系统可以用轮询方式接入",说明它把"接进已有系统"当成一等需求。代价是这套接口必须足够窄 —— 一个能用轮询对接的接口,不可能暴露太多内部语义。"PE 之间只能通过事件交互"这条规则,在接口层面同样成立。

一处容易忽略的依赖关系:序列化机制住在 core 模块里,而 core 模块没有默认参数 —— 这意味着**"用什么序列化"必须由部署者显式决定,框架不给兜底**。它与前面那条"事件可以不带保证地发送"是配套的:既然允许丢消息,序列化就不必为了可靠性加上校验与重试,可以把字节数压到最小。**两层取舍互相支撑:通信层允许丢,序列化层就只求小。**反过来,如果某个作业需要"数据也不能丢",那要改的不只是通信层的保证等级,序列化那层也要跟着重新选 —— 这两处是同一个决定的两半。

版本演进:从内部版本到退役 ​

时间线(取自 Apache 孵化器状态页):最初那篇讲的是 Yahoo! 内部的版本(成文不早于 2010);2011-09-26 进入 Apache 孵化器;2012-08-16 发布 0.5.0-incubating;2013-06-03 发布 0.6.0-incubating;2014-06-19 退役。它停在 incubating 上就到头了,从未毕业为顶级项目 —— 孵化状态页给出的全部说明只有"项目于 2014-06-19 退役"这几个字。

一次"从有损到可插拔"的变化值得单独记。 最初描述的形态是有损故障切换可以接受:状态存在进程的本地内存里、交接时即丢、靠输入流重新生成。而 0.5.0 那一版是一次全面重构,引入了 TCP 通信、可插拔的检查点机制(用于状态恢复)、发布-订阅式的集群与应用间通信、动态应用部署,以及一整套工具 —— 官方配置文档里能直接看到那个检查点模块的类名(基于文件系统的后端)。"有损"这条假设在开源版里被做成了一个可选件,而不是唯一形态。

还有一条线留在官方仓库里、没进主线:一份 README 专门讲与 Apache Helix 集成的动机,先逐条列出 S4 的局限,再逐条给出改进方向。这份材料把"运行中不加不减节点"这条假设的代价算成了数字,而那条假设最初只是一句声明:

S4 的局限(官方自述)代价
任务总是按节点数切分流若已在 S4 之外按别的因子分好片,就得重新分片 ⇒ 多一跳;两条外部已分片的流要在 S4 内 join,两条都得重哈希
加节点会改变分区数大量数据重映射 —— 4 个节点加 1 个,几乎全部 key 要重映射,而理想情况只该移动 20%
容错靠备用节点、故障时才顶上硬件利用率差(备用节点平时不干活)

Helix 给出的三条改进正好对着这三条:用一致性哈希把分区映射到节点、加节点时迁移分区而不改分区数、以及所有节点都是活跃的、节点故障后负载在剩余节点间重新分配。修这三条需要的东西(一致性哈希 + 分区迁移)在同期的键值存储里已经被用上了 —— 也就是说,这些技术在同期已经成熟,只是这套系统当时没有采用。

一条判断:"分区数 = 节点数"在小集群上是一个合理的简化,在"弹性"成为标配之后就成了硬伤。 而它留下的东西并没有消失 —— PE / PN 这套按事件键路由、按需实例化处理单元的形状,是后来所有流处理引擎的骨架;本专栏下一篇讲的那些工程化改进,做的几乎都是同一件事:把"分区与节点绑定"这条约束解开。

一处需要说清的边界:这份笔记讲的是最初那个形态,它没有覆盖 S4 后来做的事。 0.5.0 起的重构引入了检查点、发布-订阅式通信与动态部署,这些在最初那份描述里都不存在;而 Helix 那条线连主分支都没进。所以"S4 是什么"有两个答案:一份是早期形态(状态随节点消失、只做流、只关心能不能跟上速率),一份是官方仓库最后停下的地方(可插拔检查点、工具齐备、并且知道自己的切分方式是瓶颈)。两者之间的差距,就是"一个研究原型变成项目"要补的东西。

把 S4 与它同期的邻居摆在一起,"为什么是它退役"会看得更清楚:Apache 孵化器状态页上,S4 的邻居里有一批同样是"没毕业就退役"的项目,理由栏多写"缺乏活跃度"或"没能建立社区"。而 S4 那一页给出的全部说明只有"退役"两个字,没有给原因。能核到的事实只有两条:它停在了 0.6.0-incubating,以及同期流处理生态里出现了拿到更广采用的替代者(本栏下一篇讲的 Storm 就在那一批里)。一个早期标准没有被后来的生态继承,通常不在于机制错了,而在于它的约束在新场景里变成了瓶颈 —— 而"S4 加一个节点要重映射几乎全部 key"这句话,就是那个约束最具体的形态。

一条更一般的判断:评估一个系统时,"它自述的局限"比"它自述的目标"信息量大得多。 六条设计目标说的是它想成为什么,而两条假设与那份 Helix README 说的是它在什么条件下才成立。目标可以是共通的(谁都想低延迟、可扩展),而成立条件是个性化的 —— 后者才是选型时该逐条对照的东西。

还有一处对读这份材料有用的提示:最初那份描述没有版本号(讲的是一个内部系统),而进入 Apache 之后的版本号从 0.5.0-incubating 起 —— 也就是说开源发布的第一版就是 0.5.0,不是 0.1.0,因为 0.5.0 那一次是"全面重构"。这个编号本身说明:捐赠给 Apache 的是重写后的代码,而不是 Yahoo! 内部跑着的那份。 对读两边的差别时,这一点得先分清:它对应的是同一个设计的两代实现,而不只是一次版本升级 —— 所以两边的机制差异不必硬凑成"演进关系",它们是两个独立写就的实现。

它承认还缺什么 ​

这份工作留到后续的只有一句,但把边界说得干净:当前系统用的是静态路由 + 通过 ZooKeeper 的自动故障切换,缺动态负载均衡与可靠的在线 PE 迁移。

这与开头那两条假设一一对应:运行中不加不减节点(所以不需要动态负载均衡与在线迁移),有损故障切换(所以不需要把 PE 状态持久化)。S4 的优势与限制来自同一组选择。

排查:从症状到判据 ​

这类系统的症状与前面几篇都不同:它不会报错,而是悄悄丢东西。实测那条曲线就是最好的说明 —— 事件率从 2000 提到 20000(数据率 2.6 Mbps 到 26.7 Mbps),CTR 的相对误差从 0.0% 一路爬到 4.2%,而系统在约 10 Mbps 处就开始劣化,原因是网格处理不过来、发生事件丢失。

症状判据先看什么常见归因
结果精度缓慢下降、但没有报错输入速率离劣化点多远事件率、数据率跟不上速率而丢事件(load shedding)
CTR 估计突然变差、退回到长期值是否有 PN 故障节点存活状态有损故障切换生效 —— 那部分"查询/广告"实例没有短期估计
某个词的计数莫名归零PE 是否被回收PE 实例数与内存占用TTL 回收把状态一起删掉了
集群加节点后行为异常分区数是否随之变化集群配置"任务总是按节点数切分" —— 改分区数会触发大量重映射
应用换环境后行为不一致参数注入在哪一层部署参数 vs -p应用参数被写进了节点命令行
自定义模块加载失败模块类有没有无参构造模块 jar自定义模块类必须提供无参构造函数
节点重启后没接管原来的分区启动时有没有给节点 id启动参数新版要求用 -id 启动节点,才能让重启的节点关联回原来的分区
事件被路由到预期之外的节点被键属性与取值是否变过PE 身份四要素键设计变了 ⇒ 路由与状态作用域都变了

三条判读原则:

  1. 先算"离处理能力上限多远",再谈优化。 这套系统的劣化是渐变的:2000 事件/秒误差 0%、7268 时 0.2%、10480 时 0.4%、16000 时 1.7%、20000 时 4.2%。判据是"你在这条曲线的哪一段",而不是"有没有报错" —— 它不报错。
  2. 把"有意丢的"与"意外丢的"分开。 三种丢法要分清:跟不上速率而丢(设计的一部分,靠降载或扩集群解决)、PE 被 TTL 回收而丢(框架按时间做的清理,靠调 TTL 或改状态设计解决)、节点故障而丢(有损故障切换,只能靠下游降级)。混在一起看,就会把"该扩集群"的问题当成"该调参数"。
  3. 状态相关的症状,先查键与 TTL。 状态的位置完全由键决定、寿命完全由 TTL 决定。任何"某部分状态不对"的问题,答案几乎都落在这两处 —— 而不在计算逻辑里。

最后一条与前文接得上:这套系统对"丢"这件事的态度是三者里最坦然的 —— 它把"有损"写进了自我描述("部分容错"),把恢复寄托在输入流重放上。代价是下游必须能承受降级:CTR 那条路退回到长期估计,就是这条约定在生产里的实际形态。

这套系统上最有效的排查动作,大多落在"看配置"上,而不在"看指标"上 —— 因为它的行为几乎全部由键、TTL 与三层配置决定,而这三样都是静态的。三条具体的动作:

  1. 把键的分布画出来。 每个键取值对应多少流量、多少 PE 实例,直接决定并行度是否有效。取值少而流量大(按国家聚合这类)会看到几个 PN 明显比别的忙;取值极多而每个都小会看到 PE 实例数一直涨、TTL 频繁回收。
  2. 核对参数注入的那一层。 同一个应用在两个环境行为不一致时,先比对部署参数与节点命令行,而不是比对代码 —— 官方文档已经明确划过边界(-p 只放节点相关参数),而"应用参数写进了节点命令行"是最常见的越界。
  3. 算一遍处理能力与当前速率的距离。 这条曲线是量过的(约 10 Mbps 处开始劣化),所以**"我们现在离它多远"是一个能算出来的数**,不必等到精度掉了才发现。

最后一条经验与这套架构的取向一致:它给的旋钮都是为了"调形状",而不是为了"调速度" —— 键决定状态怎么切、分区数决定任务怎么分、TTL 决定状态能活多久。想让作业更快,唯一有效的动作是把形状改对;改参数只能让一个形状本来就对的作业跑得更顺一点。这与它"去中心、对称"的设计是一致的:一个不区分节点的系统,也不会提供按节点加速的手段。

一条与实测条件相关的提醒:那两组实验的硬件都很小 —— 在线 CTR 是 16 台服务器、每台 4 个 32 位处理器、2 GB 内存(峰值 1600 事件/秒),离线压力测试是 8 台、每台 4 个 64 位处理器、16 GB 内存(共 16 个 PN,一路跑到 20000 事件/秒)。"约 10 Mbps 开始劣化"这个数属于这套配置,换更大的集群这个拐点会右移。但它右移得有限 —— 因为按前面那条"键的基数决定并行度"看,有些键的流量本来就落在单个 PN 上,加机器并不会把它分散开。这条约束比集群规模更硬,也是排查时要先确认的一件:当前速率是不是集中在少数几个取值上。

相关 ​

  • Pregel —— 同一出发点、不同结论:两者都是"承认 MapReduce 不适用于某类负载,于是另造一套架构",也都选择了对称/去中心的取向(Pregel 是每台机器既当 master 候选又当 worker,S4 是所有节点职责相同)。差别在状态:Pregel 允许顶点状态跨超步存活并靠 checkpoint + 部分恢复保住它,S4 把 PE 状态放在本地内存里、故障即丢
  • 参数服务器 —— 同样公开承认有损/部分容错,并把恢复寄托在"外部状态重新生成"上:S4 靠输入流重放重建 PE 状态,参数服务器靠训练数据重算。两篇都敢这么做,是因为上层任务对扰动容忍(S4 退化到长期 CTR 估计,参数服务器丢掉一点训练数据)
  • RDD —— 对"状态该不该跨任务存活"的相反回答:RDD 的全部价值在于把中间结果跨迭代留在内存里、并用血统保证可重算;S4 则明确接受"状态随节点消失"。两篇的取舍都能追到各自负载的形状 —— 一个是迭代批处理,一个是速率不受控的事件流

参考 ​

  • L. Neumeyer, B. Robbins, A. Nair, A. Kesari. S4: Distributed Stream Computing Platform. Yahoo! Labs.

  • Apache Incubator. S4 Incubation Status 与 All Incubator Projects By Status. https://incubator.apache.org/s4/ —— 用于核对进入孵化器与退役的确切日期、以及各孵化版本的发布时间

  • Apache S4. Configuration(官方配置文档)与 README(Helix 集成). https://apache.googlesource.com/incubator-retired-s4/ —— 用于核对工具集、三层配置与默认值,以及官方自述的局限

    抽取出的源文本只含标题、作者、单位与正文,没有会议名、页码或 DOI,所以出版信息未被核对;正文引用的文献最新到 2010,成文时间不早于那一年。

贡献者 ​

文件历史 ​