7.6

查看英文版

7.6 实时与流式数据

概述与动机

你对数据流水线的大部分认知都假定数据是静止不动的。你收集一天的记录,隔夜运行一个作业,第二天早上读取结果。实时与流式数据颠覆了这个假设。你不再处理一堆已经完成的数据,而是在事件到达时持续处理这条无穷无尽的流,并不断产生答案。这就是批处理(作用于一个有界的、完整的数据集)与流处理(作用于一个无界的、永不结束的流)之间的区别。

对大型团队而言,一旦延迟开始对业务产生影响,流式处理就会浮出水面。一个晚到一小时的反欺诈决策毫无价值。一个明天才送达的个性化信号,什么也个性化不了。一个滞后于现实一个班次的运营仪表盘,会误导正在盯着它的人。第 7.2 章(数据工程)主张你应当默认选择批处理,只有在延迟真正能带来回报的地方才求助于流式处理,本章带你走完剩下的路:实时何时值得它的成本,以及如何在不烧光你的运维预算的情况下构建它。流式处理与第 3.12 章(事件驱动架构与消息传递)中的事件驱动消息模式、第 3.4 章(数据架构与存储)中的存储选择,以及第 9.2 章(可观测性与遥测)中的遥测实践都密切相关。

企业和政府场景把风险抬得更高。一家银行要在读卡器闪烁的时间内为交易的欺诈风险打分。一家交通机构要为数百万乘客追踪车辆并预测到站时间。一家福利机构要在监测理赔异常的同时,为每一个决策保留可审计的记录。在所有这些场景中,价值来自趁数据还新鲜时采取行动,而风险来自基于错误的、不完整的或事后无法重建的数据采取行动。本章对这两点都有明确的立场。

关键原则

  • 只有在延迟具有明确的业务价值时才求助于流式处理;批处理更便宜也更简单。
  • 区分有界(有限)数据和无界(永不结束)数据,并据此进行设计。
  • 把事件时间而不是到达时间当作真相来源,并为迟到和乱序数据做好规划。
  • 窗口和水位线是你从无限的流中得到有限答案的方式。
  • 优先通过幂等的接收端获得”实质上一次”的结果,而不是依赖脆弱的”恰好一次”承诺。
  • 有状态处理需要检查点,这样它才能在不丢失或重复计数的情况下恢复。
  • 从第一天起就为背压和重新处理进行设计,而不是事后补救。
  • 让流式处理逻辑保持可观测和可审计;一条悄无声息的流比一个失败的批处理更糟糕。

建议

在构建实时能力之前先证明它的必要性

最重要的流式决策,是要不要进行流式处理。实时大致会让你的运维复杂度和成本翻倍,因为你要用一个必须每一秒都保持健康的系统,去换取一个原本运行完就可以停止的作业。在投入之前,先说清楚新鲜数据能促成什么决策,以及这个决策晚到会付出什么代价。反欺诈打分、运营告警和实时个性化通常能够跨过这道门槛。一个人一天只看两次的仪表盘,无论”实时”这个词在规划会议上听起来多么诱人,几乎永远跨不过这道门槛。把延迟需求写成一个具体的数字,以秒或分钟为单位,并对照现实进行核验。人们所说的”实时”,很多时候用每隔几分钟运行一次的微批处理就能很好地满足,而成本只是一小部分。

围绕事件时间而不是处理时间来设计

流式处理中最难的一个概念是:事件发生在某一时刻,而被处理是在另一时刻。事件时间是事情实际发生的时刻,例如一名乘客刷卡的那一刻。处理时间是你的系统实际处理它的时刻。这两者会不断地相互偏离:一部手机在隧道里失去信号,一次性上传了三分钟的刷卡记录;一次网络故障打乱了消息的顺序;一个分区滞后了。如果你依据处理时间来计算,你的数字会随着你的基础设施而摆动,而不是反映真实世界。这个迟到与乱序问题是这门学科的核心,它与事件驱动架构中的事件建模直接相关。在源头就给每个事件打上其事件时间的时间戳,让这个时间戳贯穿整条流水线,并依据它来计算你的结果。

用窗口和水位线得到有限的答案

一条无界的流永不结束,所以”统计事件数量”这样的问题在你给它设定边界之前是没有答案的。窗口就是用来设定这个边界的。滚动窗口(tumbling window)把时间切分成固定的、不重叠的桶,例如每分钟一个。滑动窗口(sliding window)会相互重叠,所以一个每分钟前进一次的五分钟窗口,能给你一条平滑的移动数值。会话窗口(session window)把被不活跃间隔隔开的一阵阵活动分组,这很适合用户会话。一旦有了窗口,你就需要决定一个窗口何时算完成,因为迟到的数据可能仍会到达。水位线(watermark)是系统对”它大概已经看过截至某个事件时间的所有事件”这一情况的估计。当水位线越过一个窗口的结束点时,你就发出这个结果。调整你愿意等待多久:让窗口开得更久,你就能容忍更多的迟到,代价是更高的延迟和更多的内存占用;关得更快,你就冒着丢弃迟到者的风险。要明确决定在窗口关闭之后到达的数据会怎样:是丢弃、记录,还是发出一次更正。

让接收端幂等,并优先追求”实质上一次”

投递保证听起来很简单,实则不然。至少一次(at-least-once)投递意味着每个事件都会被处理,但有些事件可能在重试后被处理不止一次,因此计数可能会虚高。恰好一次(exactly-once)听起来很理想,但代价高昂,而且如果字面地要求跨越任意外部系统实现,往往是不可能的。实用的目标是”实质上一次”(effectively-once):可观察到的结果就好像每个事件都只被处理过一次,即使底层机制曾经重试过。要做到这一点,就要让你的幂等接收端能够安全地被重复写入,使用确定性的键和 upsert(更新插入)操作,使得一个被重放的事件是覆盖而不是重复。把至少一次投递与幂等写入结合起来,你就能在不必到处支付重量级事务协调成本的情况下获得正确的结果。把真正的恰好一次机制留给那些真正需要它的狭窄场景。

为有状态处理设置检查点,使其能够恢复

许多有用的流式计算都是有状态的:运行中的计数、跨流的联接、去重、记住近期行为的反欺诈模型。这些状态存在于内存中,一旦进程重启就会消失。检查点会定期把状态和流的位置一起做快照,这样在崩溃之后,系统能从一个一致的点恢复,而不必重放所有内容或丢失记忆。要刻意地为你的状态规模做设计,因为无界状态是在生产环境中让一个流式作业耗尽内存的常见方式。对不再需要的状态使用过期时间和存活时间,并把状态规模当作一项一等指标来监控。故障之后的恢复时间是一项真实的服务水平关切,所以要在用户替你发现之前先测试它。

用变更数据捕获从运营型数据库中进行流式处理

你常常想要对一个从未被设计成能发出事件的数据库中的变更做出反应。变更数据捕获(CDC)通过读取数据库的事务日志,把每一次插入、更新和删除都变成一条变更事件流,来解决这个问题。这远远好于按定时器轮询表,因为轮询既慢,又会错过中间状态,还会对源系统造成冲击。CDC 能让你的搜索索引、缓存、分析存储或下游服务持续与系统记录(system of record)保持同步,而且不需要对应用做侵入式的改动。要把变更流当作一等的数据产品来对待:对其模式进行版本管理,为其含义编写文档,并监视其延迟,因为下游的一切都会继承这份延迟。

优先选择流式优先架构,而不是维护两套代码库

经典的 Lambda 架构让一个批处理层负责准确、完整的历史,同时让一个速度层负责新鲜的、近似的结果,然后再把二者合并。它是可行的,但会迫使你在两个系统中两次编写和维护相同的业务逻辑,并永远地去协调二者的差异。Kappa 架构把这一切收拢起来:保留一份持久的、可重放的事件日志,把所有处理都当作流处理来运行,当逻辑变化时通过重放日志来重新处理历史。整个行业已经转向这种流式优先的形态,因为单一代码库的维护和推理成本要低得多。如果你能把你的批处理需求表达为对一份留存的事件日志的重放,你就完全避免了维护两套代码库的税。使用能够留存历史的基于日志的消息代理,这样重新处理就只是回退指针的问题,而不必重新构建。

将流以 SQL、物化视图和实时 OLAP 的形式暴露出来

并非每一个需要流式处理的人都应当去编写底层的流处理代码。流式 SQL 让分析师和工程师能用他们已经熟悉的语言来表达窗口、联接和聚合,并把结果作为物化视图持续保持最新。对于针对新鲜数据的低延迟分析查询,一个实时联机分析处理(OLAP)存储会摄取这条流,并在毫秒级别内回答切片钻取式查询,这正是支撑一个真正实时的运营仪表盘的能力。当目标是对功能和实验获得快速反馈时,把这些能力与第 7.4 章(产品分析与实验)中的产品分析实践搭配使用。在合适的地方选用这些更高层次的工具,把手写的流处理器留给它们无法表达的逻辑。

从一开始就为背压和重新处理做规划

一条流到达的速度可能超过你处理它的速度。背压(backpressure)是让一个较慢的消费者向上游发出信号、要求放慢速度,而不是崩溃或悄悄丢弃数据的机制。确保你流水线中的每一个阶段都遵守它,并把消费者延迟当作一项头条指标来监控,因为不断增长的延迟是你正在输掉这场竞赛的最早预警。重新处理是人们最常希望自己当初就已经内置的另一项能力。当你发现一个缺陷或改变一条规则时,你会想要用修正后的逻辑重放历史。而这只有在你的事件日志留存了足够的历史、并且你的接收端足够幂等能够吸收这次重放时才有可能。从第一天起就把这两者设计进去;在事故压力下去补建它们是很痛苦的。

权衡:优缺点

选择优点缺点最适合
批处理简单、便宜,易于测试和补数延迟高,在两次运行之间数据陈旧报表,大多数分析
微批处理(分钟级)近实时,比流式处理简单得多并非真正即时“实时”仪表盘
真正的流式处理(亚秒级)即时反应,持续产生结果复杂、昂贵、难以测试反欺诈、告警、实时个性化
至少一次 + 幂等接收端结果正确、成本可承受、有韧性需要有纪律的键设计大多数流式流水线
恰好一次机制端到端的强保证昂贵,跨系统时受限少数高风险路径
Lambda(批处理 + 速度层)准确的历史加上新鲜的视图需要维护两套代码库遗留系统迁移
Kappa(流式优先)单一代码库,可重放需要留存的、持久的日志新建的流式平台

核心张力在于延迟与复杂度之间。每一步向实时靠近,都会在运维负担、测试难度和金钱上付出代价,而回报并不是线性的:从每天一次变为每隔几分钟一次是便宜的,往往也已经够用;而从分钟级变为亚秒级则是开销真正集中的地方。解决这个张力的方法是为决策定价,而不是为技术定价。问一问新鲜度能促成什么行动、迟到又要付出什么代价,然后只购买这个行动能够证明合理的那么多延迟削减。当你确实需要流式处理时,依靠带有幂等接收端的至少一次投递,再加上一份流式优先的日志,因为这种组合能在不需要最重量级保证的情况下给你正确性和可重放性。

与团队讨论的问题

  1. 实时数据实际上为我们促成了什么决策,当那份数据晚到一分钟而不是即时到达时,代价是什么? 这是应当把关每一个流式项目的问题,因为相比批处理,流式处理大致会让你的运维成本和复杂度翻倍。一个大型团队可能花上几个季度去构建一个实时平台,结果只是为一个人一天看两次的仪表盘服务,那是把钱付之一炬。拿出数据所驱动的具体行动,无论是拦截一笔欺诈交易、呼叫一名运维人员,还是改变用户看到的内容,为每一项都标出延迟的代价数字。如果诚实的答案是一个五分钟的微批处理就能满足需求,那是一个值得庆祝、而不是该被隐藏的发现。这个答案应当直接决定你们是构建真正的流式处理、满足于微批处理,还是继续停留在批处理上。

  2. 我们如何处理迟到和乱序事件,在窗口关闭之后到达的数据会发生什么? 迟到和乱序数据是流式处理中最难的部分,跳过这个问题的团队会在生产环境中、当他们的数字拒绝对账时才发现它。相互竞争的压力是延迟与正确性:让窗口开得更久以捕捉迟到者,你就会拖延每一个结果并消耗更多内存;关得更快,你就会悄悄丢弃真实的数据。拿出证据来说明你的数据实际上会迟到多久,以你各个数据源上事件时间与处理时间之间的差距来衡量,因为一个在隧道里的移动端数据源,行为会与服务端事件截然不同。明确决定迟到数据是被丢弃、被记录,还是触发一次更正,并确保下游的每个人都知道这一点。在数字必须站得住脚的政府场景中,悄悄丢弃迟到事件可能是一个合规问题,因此这项政策需要经过深思熟虑并加以记录。

  3. 我们的接收端是否足够幂等,能让我们安全地重放历史,我们的事件日志是否留存了足够的内容来让重放成为可能? 重新处理是团队最常希望自己已经内置、却又最常没有做到的能力,它取决于两件事协同工作:能够吸收重放事件而不产生重复的幂等接收端,以及留存了足够历史以供重放的持久日志。缺了其中任何一个,修复一个逻辑缺陷就意味着你无法干净地重新计算受影响的时段,而是被迫在压力下手工修补数字。拿出你当前的留存窗口和一个具体的测试:挑选上个季度的一个真实缺陷,问一问你是否本可以用修正后的逻辑对受影响的数据进行重放。反对这样做的拉力是成本,因为留存历史和设计幂等写入需要预先投入存储和纪律。但另一种选择会在最糟糕的时刻()事故发生期间()浮现出来,所以这个答案决定了你在需要之前应当在可重放性上投入多少。

  4. 当一个流式作业崩溃时,它必须多快恢复,它被允许保存多少状态,我们是否真的在生产负载下计时测试过一次恢复? 一个挂掉的批处理作业明天可以重跑,但一个挂掉的常驻流是一场正在发生的中断,而保存着运行计数、联接或反欺诈模型的有状态作业,可能会丢失几分钟的记忆,或者在重启后需要很长时间才能重新加载状态。对大型团队而言,这正是一个不起眼的细节悄悄决定你真实可用性的地方:无界状态会不断增长,直到一个作业耗尽内存,而一次缓慢的检查点恢复会把一次十秒的小故障变成一次十分钟的中断。相互竞争的压力是新鲜度与安全性,因为更频繁的检查点会缩短恢复时间,但会增加开销,而慷慨的状态保留期能提高准确性,但会带来内存耗尽的风险。拿出一个具体的恢复时间目标、你当前的状态规模及其增长曲线、你的检查点间隔,以及一次真实故障切换演练的结果,而不是一个一厢情愿的估计。在流支撑着反欺诈打分或公共安全信息流的企业和政府场景中,一条未经测试的恢复路径是你在没有度量的情况下就默认接受的运维风险,所以要把这场演练当作一项要求,而不是锦上添花。

  5. 我们运行的是单一的流式优先代码库,还是分开的批处理层和速度层,把这两者保持一致实际上要花我们多少代价? 批处理层负责准确历史、速度层负责新鲜结果的 Lambda 模式,迫使你在两个系统中两次编写相同的业务逻辑,并永远地协调它们的答案,而流式优先(Kappa)形态则保留一份持久的、可重放的日志,把所有处理都当作流处理来运行。对大型组织而言,重复的逻辑正是漂移和有争议数字滋生的地方,因为一条规则在一层中变了、另一层没变,工程师们要花真实的时间去解释为什么两者不一致。倾向于保留两者的拉力是惯性以及一个久经验证的批处理层带来的安心感,所以要诚实地把这一点与维护税放在一起权衡。拿出你目前在两处都运行的计算清单、由两层结果不一致所引发的事故,以及对你的事件日志是否留存了足够历史、能把批处理需求表达为重放的评估。在政府和受审计的企业场景中,两个层为同一时段报告不同的数字本身就是一项合规负债,因为你必须能够说清哪个数字是权威的,以及为什么。

  6. 当这个常驻系统在凌晨三点出故障时,由谁来运维它,我们是否为它所要求的待命负荷和专业技能做了预算,还是仍然假定按批处理的方式配置人手? 流式处理把成本从构建转移到了运行:系统必须每一秒都保持健康,这意味着真实的待命值班覆盖、熟悉事件时间、水位线、状态和投递语义的工程师,以及比一个运行完就停止的作业更难做的测试。团队常常仅凭一个流式平台的能力就批准了它,却从未为维持它运转的人员拨出经费,于是平台不断退化,信任也随之流失。这里的权衡是范围与可持续性:每一条新增的实时流水线,都是又一件可能呼叫某个人的事情,所以问题在于它所购买的延迟削减,是否值得一份永久的运维承诺。拿出一份诚实的清单,说明生产环境中每条流由谁负责、你当前的待命轮值及其余量,以及事件时间方面的专业能力究竟落在哪里()是招聘、合作伙伴,还是托管服务。对公共机构或大型企业而言,还要加上采购和招聘的提前周期以及任何托管服务选项,因为一个依赖你既招不到也留不住的稀缺人才的实时平台,本质上是一个计划让人手不足、故障频发的系统运行下去的计划。

行业视角

初创企业。 流式处理很少是你的第一步,建立一个沉重的平台可能会拖垮一个微小的团队。挑出那唯一一个触及你核心价值的信号,把事件放到单一一个留存的、基于日志的消息代理上,运行一个使用带键、幂等接收端的轻量处理器,这样一次至少一次的重试就永远不会重复计数。保留几天的历史,以便你能通过固定的逻辑进行重放,并优先选用托管流式服务,而不是自己运维一个集群,因为你最稀缺的资源是工程注意力。

小型企业。 你很可能没有流式处理专家,也没有运维常驻基础设施的胃口,所以应当把实时能力当作你已经在用的工具内部购买的东西,而不是一个需要你配人的系统。把这个需求框定为一个带有具体数字的延迟问题,在大多数情况下,一个每隔几分钟刷新一次的微批处理,就能以一小部分的成本和风险满足它。选择那些对延迟保持透明、并且易于回退的供应商的实时功能,把定制的流式处理留给新鲜数据直接驱动收入或安全的少数场景。

企业。 这里的问题是在众多团队之间保持一致性和成本:一个共享的基于日志的平台、一套标准的事件时间和迟到数据策略,以及幂等的接收端,让各团队不再重复发明脆弱的流水线。为常驻运维和待命负荷明确做预算,标准化在一份流式优先的日志上以避免重复的批处理代码库,并把流当作有负责人、有模式版本管理、被监控延迟的治理型数据产品来管理,而不是一堆各自为政的定制作业。把延迟、恢复时间和每条流的成本作为组合层面的指标来跟踪。

政府。 可审计性和公共问责塑造着每一个选择。把每一个已处理的事件留存在一份持久的日志中,这样报告给监督机构的数字()客流量、福利理赔异常、反欺诈决策()就能被精确地重建,并把迟到数据的策略明确地文档化,而不是悄悄丢弃事件。采购应当要求数据的可移植性,并要求托管服务披露其投递和留存保证,而在规则变化之后的任何数字重述,都应当是一次通过修正后逻辑进行的、站得住脚的重放,而不是一次没人能追溯的手工修补。

示例

初创企业。 一款消费者应用想要给用户展示一条实时活动信息流,并在可疑登录发生时立即标记出来。团队抵制建立一个沉重的流式平台的冲动。他们把事件放到单一一个留存的、基于日志的消息代理上,为登录风险逻辑运行一个轻量的流处理器,并向一个支撑活动信息流的实时 OLAP 存储供数。每一个接收端都是带键且幂等的,所以一次至少一次的重试永远不会重复计数。当他们后来在风险规则中发现一个缺陷时,只需在夜间用修正后的逻辑重放日志即可,因为他们保留了一周的历史,也从未需要第二套批处理代码库。

企业。 一家零售银行在授权窗口期内为每一笔卡交易的欺诈风险打分,把实时交易流与一个反映近期账户行为的有状态模型进行联接。检查点让打分服务能在几秒钟内从一次节点故障中恢复,而不丢失最近几分钟的记忆。与此同时,变更数据捕获把核心银行数据库中的更新流式传输到一个搜索索引和一个个性化服务中,让二者都保持新鲜而无需轮询。运营仪表盘从一个实时 OLAP 存储中读取数据,使风险和运营团队能够实时观察业务的动态,整条流水线都会发出第 9.2 章所述的延迟和吞吐量遥测数据。

政府。 一家都市交通管理机构摄取车辆位置和票务刷卡数据,实时预测到站时间并监控拥挤情况,同时为公共应用和运营中心供数。因为在隧道中的乘客会以延迟的批次上传刷卡记录,团队依据事件时间计算客流量,并把水位线调整到与观察到的迟到程度相匹配,对任何在窗口关闭之后到达的事件进行记录,而不是悄悄丢弃。每一个已处理的事件都被保留在一份可审计的日志中,这样报告给监督机构的客流数字就能被精确地重建。当票价规则变化时,他们会用修正后的逻辑对受影响的时段进行重放,并产生一份站得住脚的数字重述。

商业案例:动机、投资回报率与总拥有成本

实时数据的回报来自在行动仍然有意义的时候采取行动。在授权过程中被截获的欺诈行为,防止了一笔损失,而一次夜间批处理只会事后报告它。能在一次会话内做出响应的个性化,能以明天的推荐无法做到的方式提升转化率。反映当下情况的运营监控,让你能在一个小问题变成一次中断或一起公共事件之前进行干预。在每一种情形中,价值都是现在行动与以后行动之间的差额,而这个差额正是你在论证理由时应当量化的东西。

总拥有成本比批处理更高,对此保持诚实能保护你的信誉。你要为常驻基础设施付费,为理解事件时间、水位线、状态和投递语义的工程师付费,也要为一个必须持续保持健康、而不是运行完就停止的系统所需的更难的测试和待命负担付费。一个建立在留存日志之上的流式优先架构,通过让你免于维护重复的批处理代码库,降低了持续成本,而选择带有幂等接收端的至少一次投递,则避免了端到端恰好一次机制的开销。最昂贵的错误,是在微批处理或批处理本就足够的地方构建了实时能力,所以最有力的成本论证往往是一个不进行流式处理的决定。向领导层陈述时,应当围绕具体的、对延迟敏感的决策及其可度量的回报来组织论点,同时同样清楚地说明,在哪些地方停留在批处理上能够节省成本而不损失价值。

反模式与陷阱

  • 为了面子而构建流式处理,而每隔几分钟一次的微批处理本就能满足需求。
  • 依据处理时间进行计算,导致你的数字随着基础设施摆动,而不是反映真实世界。
  • 忽视迟到和乱序数据,直到对账在生产环境中失败。
  • 到处追求字面意义上的恰好一次,而不是带幂等接收端的至少一次。
  • 无界状态没有过期机制,悄悄增长直到一个作业耗尽内存。
  • 没有检查点,导致一次重启丢失状态或迫使进行一次完整重放。
  • 按定时器轮询运营型数据库,而不是使用变更数据捕获。
  • 维护着一个逻辑重复且不断漂移的 Lambda 批处理层和速度层。
  • 留存窗口太短,当你发现一个缺陷时无法重放历史。
  • 没有延迟、吞吐量或新鲜度指标的流,悄无声息地失败。

成熟度模型

  • 第 1 级,启动(Initiate): 一切都是批处理,或者少数几个手工搭建的流式作业在没有监控的情况下被动运行。数字依据处理时间计算,迟到数据被忽视,一次重启就会丢失状态。没有人能够重放历史来修复缺陷,问题只有在下游数字拒绝对账时才会被发现。
  • 第 2 级,发展(Develop): 一些团队在一个带检查点的基于日志的消息代理上运行核心流式流水线,他们能区分事件时间和处理时间,并使用基本的窗口。团队之间的实践并不一致:投递是至少一次,但并非所有接收端都是幂等的,迟到数据的处理是临时拼凑的,延迟被非正式地关注,而不是被设置告警。
  • 第 3 级,标准化(Standardize): 事件时间、水位线和一套明确的迟到数据策略在整个组织范围内被文档化并加以应用。接收端为了实现”实质上一次”的结果而幂等,状态设有过期机制,变更数据捕获按约定为下游系统供数。一份留存的日志支持重放,延迟、吞吐量和新鲜度作为一项全组织标准被监控并设有告警,而不是各团队的各自习惯。
  • 第 4 级,管理(Manage): 流式处理体系被度量并对照基线进行控制。每条流水线都带有端到端延迟、消费者延迟、恢复时间、事件时间偏差、迟到事件比率、状态规模,以及每百万事件成本的服务水平目标,所有这些都对照约定的目标进行跟踪,并对回归发出告警。恢复被演练并计时,而不是被想当然地假定,背压余量和状态增长被当作容量信号来关注,一条新的流必须先通过这些指标才能上线到生产环境。
  • 第 5 级,编排(Orchestrate): 一个流式优先的架构从一份可重放的日志中同时满足新鲜和历史需求,流式 SQL、物化视图和实时 OLAP 让新鲜数据被广泛地访问。重新处理是常规且经过测试的操作,平台会依据度量得到的负载和成本自动伸缩和再平衡,各条流也会依据证据被淘汰、重新界定范围或替换。流式处理与业务和风险规划相整合,随着负载和成本图景的变化,每条流都端到端地保持可观测和可审计。

讨论话题

  1. 在你的技术栈中,“实时”真正值回票价的地方在哪里,又在哪里只是一个未经检验的愿望?
  2. 在你各个数据源上,事件时间与处理时间之间的差距有多大,你测量过吗?
  3. 你能否把一套 Lambda 式的批处理加速度层设置收拢为单一一个流式优先的代码库,什么会阻碍你这样做?
  4. 你的哪些接收端是真正幂等的,你今天能否安全地用修正后的逻辑重放上个季度的数据?
  5. 对于在窗口关闭之后到达的数据,你的策略是什么,下游的每个人都知道吗?
  6. 变更数据捕获会如何改变你保持搜索、缓存和分析同步的方式?

关键要点

  • 只有在一个对延迟敏感的决策能够为其买单时才求助于流式处理;批处理和微批处理是更便宜的默认选择。
  • 依据事件时间进行计算,把迟到和乱序数据当作核心问题,用窗口和水位线来处理。
  • 优先选择带幂等接收端的至少一次投递以获得”实质上一次”的结果,而不是到处追求字面意义上的恰好一次。
  • 为有状态处理设置检查点,限定你的状态规模,并把消费者延迟当作一项头条指标来监控。
  • 用变更数据捕获从运营型数据库中进行流式处理,而不是轮询。
  • 优先选择建立在留存的、可重放日志之上的流式优先架构,而不是维护两套代码库。
  • 通过流式 SQL、物化视图和实时 OLAP 暴露流,并让每一条流都保持可观测和可审计。

参考资料与延伸阅读

  • Tyler Akidau, Slava Chernyak, and Reuven Lax, “Streaming Systems.”
  • Martin Kleppmann, “Designing Data-Intensive Applications.”
  • Nathan Marz and James Warren, “Big Data” (Lambda architecture).
  • Jay Kreps, “Questioning the Lambda Architecture” (O’Reilly Radar).
  • Fabian Hueske and Vasiliki Kalavri, “Stream Processing with Apache Flink.”
  • Ben Stopford, “Designing Event-Driven Systems.”
  • Tyler Akidau and colleagues, “The Dataflow Model” (VLDB paper on windowing and watermarks).