Skip to content

第 21 章:广告点击事件聚合

引言

随着 Facebook、YouTube、TikTok 等平台的发展,数字广告已成为规模庞大的行业。

因此,追踪广告点击事件很重要。本章探索如何设计一个达到 Facebook 或 Google 规模的广告点击事件聚合系统。

数字广告中有一种称为**实时竞价(RTB)**的流程,用于买卖数字广告库存:

<img src="/images/chapter-21/digital-advertising-example.png" alt="数字广告示例" width="500" />

RTB 通常在一秒内完成,因此速度很重要。数据准确性同样重要,因为它会影响广告主支付的金额。

广告主可以根据点击事件聚合结果,调整目标受众和关键词等策略。


第 1 步:理解问题并确定设计范围

  • 候选人:输入数据是什么格式?
  • 面试官:每天有 10 亿次广告点击,共 200 万条广告。点击事件数每年增长 30%。
  • 候选人:系统主要要支持哪些查询?
  • 面试官:重点考虑以下查询:
    • 返回广告 X 在过去 Y 分钟内的点击次数。
    • 返回过去 1 分钟点击量最高的 100 条广告。两个参数都应可配置;每分钟聚合一次。
    • 上述查询支持按 ipuser_idcountry 过滤。
  • 候选人:需要考虑哪些边界情况?我想到:
    • 事件可能晚于预期到达。
    • 事件可能重复。
    • 系统的不同部分可能停机,因此需要考虑恢复。
  • 面试官:这些都要考虑。
  • 候选人:延迟要求是什么?
  • 面试官:广告点击聚合的端到端延迟最多几分钟;RTB 要求不到一秒。点击聚合主要用于计费和报表,因此几分钟延迟可以接受。

功能需求

  • 聚合 ad_id 在过去 Y 分钟的点击数。
  • 每分钟返回点击量最高的 100 个 ad_id
  • 支持根据不同属性过滤聚合结果。
  • 数据集规模与 Facebook 或 Google 相当。

非功能需求

  • 聚合结果用于 RTB 和广告计费,正确性很重要。
  • 妥善处理延迟到达或重复的事件。
  • 系统应有韧性,能够承受局部故障。
  • 端到端延迟最多几分钟。

粗略估算

  • 10 亿日活跃用户。
  • 假设每位用户每天点击一条广告,则每天有 10 亿次点击。
  • 广告点击平均 QPS 为 10,000。
  • 峰值 QPS 为平均值的 5 倍,即 50,000。
  • 单次广告点击占用 0.1 KB,每日所需存储约 100 GB。
  • 每月所需存储约 3 TB。

第 2 步:提出概要设计并达成共识

本节讨论查询 API、数据模型和概要设计。

查询 API 设计

API 是客户端与服务器之间的契约。这里的客户端是仪表盘用户,例如数据科学家、分析师或广告主。

功能需求包括:

  • 聚合 ad_id 在过去 Y 分钟的点击数。
  • 返回过去 M 分钟点击量最高的 N 个 ad_id
  • 支持根据不同属性过滤聚合结果。

需要两个接口;过滤条件可通过查询参数传入。

聚合某个 ad_id 在过去 M 分钟的点击数

GET /v1/ads/{:ad_id}/aggregated_count

查询参数:

  • from:起始分钟,默认当前时间前 1 分钟。
  • to:结束分钟,默认当前时间。
  • filter:过滤策略标识,例如 001 表示“非美国的点击”。

响应:

  • ad_id:广告标识。
  • count:起止分钟之间的聚合点击数。

返回过去 M 分钟点击量最高的 N 个广告 ID

GET /v1/ads/popular_ads

查询参数:

  • count:返回点击量最高的 N 条广告。
  • window:聚合窗口大小,单位为分钟。
  • filter:过滤策略标识。

响应:

  • 广告 ID 列表。

数据模型

系统保存原始数据和聚合数据。

原始数据如下:

[AdClickEvent] ad001, 2021-01-01 00:00:01, user 1, 207.148.22.22, USA

结构化示例:

ad_idclick_timestampuseripcountry
ad0012021-01-01 00:00:01user1207.148.22.22USA
ad0012021-01-01 00:00:02user1207.148.22.22USA
ad0022021-01-01 00:00:02user2209.153.56.11USA

聚合后的数据:

ad_idclick_minutefilter_idcount
ad00120210101000000122
ad00120210101000000233
ad00120210101000100121
ad00120210101000100236

filter_id 用于支持过滤需求。

filter_idregionIPuser_id
0012US**
0013*123.1.2.3*

为快速返回过去 M 分钟点击量最高的 N 条广告,还要维护以下结构:

most_clicked_ads
window_sizeinteger聚合窗口大小(M 分钟)
update_time_minutetimestamp最近更新时间戳(精度为 1 分钟)
most_clicked_adsarrayJSON 格式的广告 ID 列表

保存原始数据和聚合数据各有什么优缺点?

  • 原始数据保留完整信息,支持过滤和重新计算。
  • 聚合数据的数据集更小、查询更快。
  • 原始数据占用更多存储,查询也更慢。
  • 聚合数据是派生数据,因此会损失部分原始信息。

本设计结合两种方式:

  • 保留原始数据便于调试。若聚合逻辑存在错误,可以定位问题并回填。
  • 同时存储聚合数据,以提升查询性能。
  • 原始数据可放入冷存储,降低成本。

选择数据库时,要考虑:

  • 数据是关系型、文档型还是二进制对象?
  • 负载以读取、写入还是两者为主?
  • 是否需要事务?
  • 查询是否依赖 SUMCOUNT 等 OLAP 函数?

原始数据平均 10,000 QPS、峰值 50,000 QPS,因此写入密集。原始数据主要在出问题时作为备份使用,读取量较低。

关系型数据库可以完成任务,但扩展写入能力可能困难。Cassandra 或 InfluxDB 对高写入负载有更好的原生支持。

另一选择是 Amazon S3 搭配 ORC、Parquet 或 AVRO 等列式数据格式。由于这种组合不够熟悉,这里选择 Cassandra。

聚合数据同时承受大量读写:仪表盘和告警持续查询,聚合服务也每分钟写入一次。因此,这里同样采用 Cassandra。

概要设计

系统如下:

<img src="/images/chapter-21/high-level-design-1.png" alt="概要设计一" width="500" />

输入与输出都作为无边界的数据流传输。

为了避免同步下游节点在消费者崩溃时阻塞整个系统,使用消息队列(Kafka)进行异步处理,让生产者与消费者解耦。

<img src="/images/chapter-21/high-level-design-2.png" alt="概要设计二" width="500" />

第一个消息队列存储广告点击事件:

ad_idclick_timestampuser_idipcountry

第二个消息队列存储每分钟聚合的广告点击数:

ad_idclick_minutecount

以及每分钟聚合的点击量最高的 N 条广告:

update_time_minutemost_clicked_ads

设置第二个消息队列,是为了实现端到端的“恰好一次”原子提交语义:

<img src="/images/chapter-21/atomic-commit.png" alt="原子提交" width="500" />

对于聚合服务,MapReduce 框架是一个可选方案:

<img src="/images/chapter-21/ad-count-map-reduce.png" alt="广告点击计数的 MapReduce" width="500" />

<img src="/images/chapter-21/top-100-map-reduce.png" alt="前 100 条广告的 MapReduce" width="500" />

每个节点负责单一任务,并把处理结果发送给下游节点。

Map 节点从数据源读取数据,再进行过滤和转换。

例如,Map 节点可按 ad_id 将数据分配给不同聚合节点:

<img src="/images/chapter-21/map-node.png" alt="Map 节点" width="500" />

另一种方式是将广告分布到 Kafka 分区,让聚合节点作为消费者组成员直接订阅。但独立的 Map 节点可在后续处理前清洗或转换数据。

此外,如果我们无法控制数据的生产方式,同一 ad_id 的事件可能被送入不同分区,也需要 Map 节点重新分配。

聚合节点每分钟在内存中按 ad_id 统计点击事件。

Reduce 节点汇总聚合节点的结果,生成最终结果:

<img src="/images/chapter-21/reduce-node.png" alt="Reduce 节点" width="500" />

这个有向无环图(DAG)采用 MapReduce 范式,利用分布式并行计算,将海量数据转换为较小的聚合结果。

在 DAG 中,中间数据存放在内存;不同节点通过 TCP 或共享内存通信。

下面看它如何满足各种使用场景。

场景 1:聚合点击次数

<img src="/images/chapter-21/use-case-1.png" alt="场景一" width="500" />
  • 使用 ad_id % 3 对广告分区。

场景 2:返回点击量最高的 N 条广告

<img src="/images/chapter-21/use-case-2.png" alt="场景二" width="500" />
  • 示例聚合前 3 条广告,也可轻松扩展到前 N 条。
  • 每个节点维护堆,以便快速取得前 N 条广告。

场景 3:数据过滤: 为快速过滤,可预定义过滤条件,并按条件预先聚合:

ad_idclick_minutecountrycount
ad001202101010001USA100
ad001202101010001GPB200
ad001202101010001others3000
ad002202101010001USA10
ad002202101010001GPB25
ad002202101010001others12

这种技术称为星型模式,广泛应用于数据仓库。过滤字段称为维度

优点:

  • 易于理解和实现。
  • 现有聚合服务可复用,逐步增加星型模式的维度。
  • 结果已预先计算,按条件访问速度快。

局限是会产生更多分组和记录;过滤条件越多,增长越明显。


第 3 步:深入设计

下面进一步讨论几个重要主题。

流处理与批处理

前述概要架构属于流处理系统。三类系统对比如下:

服务(在线系统)批处理系统(离线系统)流处理系统(近实时系统)
响应能力快速响应客户端无需响应客户端无需响应客户端
输入用户请求有边界、有限但量大的数据无边界输入(无限数据流)
输出给客户端的响应物化视图、聚合指标等物化视图、聚合指标等
性能衡量可用性、延迟吞吐量吞吐量、延迟
示例在线购物MapReduceFlink [13]

本设计结合了流处理与批处理。

新数据到达时,通过流处理近实时生成聚合结果;历史数据备份则使用批处理。

同时存在批处理与流处理两条处理路径的架构称为 Lambda 架构。缺点是需要维护两套处理代码。

Kappa 架构将批处理与流处理合并为一条路径,核心是使用同一个流处理引擎。

Lambda 架构:

<img src="/images/chapter-21/lambda-architecture.png" alt="Lambda 架构" width="500" />

Kappa 架构:

<img src="/images/chapter-21/kappa-architecture.png" alt="Kappa 架构" width="500" />

本章概要设计采用 Kappa 架构,因为历史数据的重新处理也经过聚合服务。

例如,当聚合逻辑出现严重错误,必须重新计算聚合数据时,可以从保存的原始数据重新计算:

  • 重算服务从原始数据存储读取数据,这一步是批处理任务。

  • 将读取的数据发送到专用聚合服务,避免影响实时聚合服务。

  • 聚合结果进入第二个消息队列,再更新聚合数据库中的结果。

    重新计算示例

时间

聚合需要时间戳,可在两个时点生成:

  • 事件时间:广告点击发生时。
  • 处理时间:服务器处理事件时的系统时间。

异步消息队列和网络延迟可能使两者差异很大:

  • 使用处理时间,聚合结果可能不准确。
  • 使用事件时间,就必须处理延迟到达的事件。

没有完美方案,需要权衡:

优点缺点
事件时间聚合结果更准确客户端时钟可能有误,也可能被恶意用户伪造
处理时间服务器时间戳更可靠事件迟到时,时间戳不对应实际点击时间

由于数据准确性很重要,这里按事件时间聚合。

为缓解延迟事件带来的问题,可以使用“水位线”(watermark)技术。

下例中,事件 2 错过了本应进入的聚合窗口:

<img src="/images/chapter-21/watermark-technique.png" alt="水位线技术" width="500" />

如果有意延长聚合窗口,就能降低遗漏事件的概率。窗口延长部分由水位线控制:

<img src="/images/chapter-21/watermark-2.png" alt="水位线示例" width="500" />
  • 水位线允许的等待时间短,延迟低,但更可能漏掉事件。
  • 等待时间长,漏掉事件的概率较低,但延迟更高。

无论如何设置水位线,都可能漏掉极少数事件。与其为这些低概率事件无限延长等待时间,不如在每日结束时进行对账,修正不一致。

聚合窗口

窗口函数有四种:

  • 滚动窗口(固定窗口)。
  • 跳跃窗口。
  • 滑动窗口。
  • 会话窗口。

广告点击计数使用滚动窗口:

<img src="/images/chapter-21/tumbling-window.png" alt="滚动窗口" width="500" />

过去 M 分钟点击量最高的 N 条广告使用滑动窗口:

<img src="/images/chapter-21/sliding-window.png" alt="滑动窗口" width="500" />

投递保证

聚合数据用于计费,因此准确性优先。

需要讨论:

  • 如何避免处理重复事件。
  • 如何确保所有事件都得到处理。

可用的投递保证有至多一次、至少一次和恰好一次。

多数场景下,若能接受少量重复,“至少一次”就足够。但本系统中,即使只有几个百分点的差异,也可能造成数百万美元的偏差,因此需要“恰好一次”投递语义。

数据去重

重复数据是最常见的数据质量问题之一。

来源很多:

  • 客户端:可能多次发送同一事件。恶意重复事件最好由风控引擎处理。
  • 服务器故障:聚合服务节点在处理中停机,上游没有收到确认,于是重新发送事件。

下面是最后一跳没有确认事件而产生重复的例子:

<img src="/images/chapter-21/data-duplication-example.png" alt="数据重复示例" width="500" />

此处偏移量 100 会被处理并发送给下游多次。

一种缓解办法是将最近处理的偏移量保存到 HDFS 或 S3,但这样可能出现结果从未到达下游、偏移量却已保存的情况:

<img src="/images/chapter-21/data-duplication-example-2.png" alt="数据重复示例二" width="500" />

最终,可以将保存偏移量与向下游写入作为一个原子操作;这需要分布式事务:

<img src="/images/chapter-21/data-duplication-example-3.png" alt="数据重复示例三" width="500" />

个人补充: 如果下游以幂等方式处理聚合结果,也可以不使用分布式事务。

扩展系统

下面讨论系统增长时如何扩展。

消息队列、聚合服务和数据库是三个独立组件。由于它们彼此解耦,可以分别扩展。

消息队列如何扩展:

  • 不限制生产者数量,因此生产者容易扩展。

  • 将消费者组织为消费者组,并增加消费者数量。

  • 为此还要提前创建足够多的分区。

  • 当消费者达到数千个时,再均衡可能耗时较长,建议在非高峰时段操作。

  • 也可按地域划分主题,例如 topic_natopic_eu

    扩展消费者

聚合服务如何扩展:

<img src="/images/chapter-21/aggregation-service-scaling.png" alt="聚合服务扩展" width="500" />
  • 增加 MapReduce 节点即可横向扩展。

  • 利用多线程可以提高聚合服务吞吐量。

  • 也可以借助 Apache YARN 等资源管理平台进行多进程处理。

  • 多线程较简单,但多进程方案更易扩展,因此实践中更常见。

  • 多线程示例如下:

    多线程示例

数据库如何扩展:

  • Cassandra 原生支持基于一致性哈希的横向扩展。

  • 向集群添加节点时,数据会在各个虚拟节点之间自动再均衡。

  • 无需手动重新分片。

    Cassandra 扩展

还要考虑热点:某条广告的热度远高于其他广告时怎么办?

<img src="/images/chapter-21/hotspot-issue.png" alt="热点问题" width="500" />
  • 图中聚合服务节点向资源管理器申请额外资源。
  • 资源管理器分配资源,避免原节点过载。
  • 原节点将事件拆成三组,每个聚合节点处理 100 个事件。
  • 各节点把结果写回原聚合节点。

更复杂的热点处理方式还包括:

  • 全局/局部聚合。
  • 拆分去重聚合。

容错

聚合节点在内存中处理数据。如果节点停机,已处理但尚未持久化的数据会丢失。

其他节点接管后,可以利用 Kafka 消费者偏移量从上次位置继续。不过,计算过去 M 分钟内点击量最高的 N 条广告,还需要保存额外的中间状态。

可在某个分钟边界为进行中的聚合制作快照:

<img src="/images/chapter-21/fault-tolerance-example.png" alt="容错示例" width="500" />

节点停机后,新节点读取最近提交的消费者偏移量及最新快照,然后继续任务:

<img src="/images/chapter-21/fault-tolerance-recovery-example.png" alt="容错恢复示例" width="500" />

数据监控与正确性

聚合数据用于计费,必须严格监控,保证结果正确。

可以监测:

  • 延迟:跟踪不同事件的时间戳,了解系统的端到端延迟。
  • 消息队列积压:积压突然增加时,应增加聚合节点。Kafka 基于分布式提交日志,因此应监控记录滞后量(records lag)。
  • 聚合节点资源:CPU、磁盘、JVM 等。

还需要每日结束时运行批量对账任务。它从原始数据重新计算聚合结果,并与聚合数据库中的实际数据比较:

<img src="/images/chapter-21/reconciliation-flow.png" alt="对账流程" width="500" />

替代设计

在通用系统设计面试中,不要求掌握大数据处理软件的内部细节。

解释设计思路并讨论权衡,比记住具体工具更重要,因此本章介绍了通用方案。

另一种使用现成工具的设计是:将广告点击数据存入 Hive,并在其上方建立 Elasticsearch 查询层,加快查询。

聚合通常可由 ClickHouse 或 Druid 等 OLAP 数据库完成。

<img src="/images/chapter-21/alternative-design.png" alt="替代设计" width="500" />

第 4 步:总结

本章讨论了:

  • 数据模型与 API 设计。
  • 使用 MapReduce 聚合广告点击事件。
  • 扩展消息队列、聚合服务和数据库。
  • 缓解热点问题。
  • 持续监控系统。
  • 使用对账保证正确性。
  • 容错。

广告点击事件聚合是典型的大数据处理系统。

如果事先了解以下相关技术,会更容易理解和设计它:

  • Apache Kafka。
  • Apache Spark。
  • Apache Flink。