流媒体广告请求的分布式工程:从实时决策到个性化投放

2026-08-05 47 预计阅读时间: 1 分钟
来源: infoq.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:10 分钟

在流媒体播放过程中,一次广告请求看起来只是“返回一条广告”,实际上要同时处理用户画像、广告候选、预算节奏、瀑布流优先级、实时竞价或内部决策,以及严格的播放时延约束。JioHotstar 对这类系统的实践说明,广告体验的核心并不只是更复杂的推荐模型,而是让多个分布式服务在极短时间内完成协调。

一次广告请求到底要完成什么

播放器在即将进入广告位时发起请求。请求通常携带内容、广告位、设备、地域、用户标识和时间信息。服务端需要完成一条实时决策链:

  1. 校验请求和用户隐私状态。
  2. 从候选广告池中筛选符合地域、设备、内容和定向条件的广告。
  3. 根据广告活动的预算、频控和投放节奏进行过滤。
  4. 按瀑布流层级或优先级依次尝试广告源。
  5. 在超时前返回可播放的广告素材和跟踪信息。
  6. 异步记录曝光、填充、跳过和播放完成事件。

这条链路的难点在于,某个下游服务变慢时,系统不能让播放器一直等待。广告决策服务通常需要设置总超时,并为画像、库存、预算和素材服务分别分配预算。关键路径只保留决定“能否投放”的调用,统计和分析写入则放到异步链路中。

瀑布流不是简单的优先级列表

瀑布流可以理解为一组按优先级排列的广告层级。例如,第一层尝试高价值的定向广告,第二层尝试更宽泛的广告,最后使用保底广告。每一层都可能连接不同的广告源或内部库存。

但如果系统永远从第一层开始,并在每次请求中顺序尝试所有层级,就会出现两个问题:高优先级库存可能被过快消耗,低优先级广告也可能因为等待过长而错过播放窗口。因此,实际决策通常需要结合以下信号:

  • 当前广告活动的预算消耗速度。
  • 目标曝光量与已完成曝光量之间的差距。
  • 用户频次上限和最近一次曝光时间。
  • 广告源的历史填充率、响应时间和错误率。
  • 播放器距离广告位开始的剩余时间。

所谓 pacing,即投放节奏控制,目标不是尽可能多地展示某个活动,而是在整个活动周期内平稳地消耗预算。可以用一个简化的节奏因子表达这种思想:

pacing_factor = target_spend_to_now / actual_spend_to_now

当因子大于 1,说明活动消耗偏慢,可以适度提高参与机会;当因子小于 1,说明消耗偏快,需要降低被选中的概率。生产系统还需要加入平滑、最小样本量和异常保护,避免短时间流量波动造成投放频率剧烈变化。

延迟优化要围绕播放窗口

流媒体广告的延迟不是普通接口的平均响应时间问题。用户已经在等待内容继续播放,尾部延迟,尤其是 P95 或 P99,往往比平均值更影响体验。

可以采用几种工程手段:

  • 为每个下游调用设置独立超时,而不是共享一个无限等待的请求上下文。
  • 对稳定的配置、广告元数据和候选集合使用本地缓存或区域缓存。
  • 让多个互不依赖的候选查询并行执行。
  • 在达到最低决策质量后提前返回,不等待低收益的额外候选。
  • 使用区域化部署,让请求尽量在靠近播放器的节点完成。
  • 将曝光、计费和分析事件从同步返回路径中移出。
  • 为服务调用设置熔断、限流和降级策略。

下面是一个可以直接运行和改造的简化示例。它不代表任何特定平台的实现,只演示如何在一个广告决策服务中执行“超时、瀑布流和保底广告”这三个动作。保存为 ad_decision.py 后运行 python ad_decision.py

from concurrent.futures import ThreadPoolExecutor, TimeoutError
from dataclasses import dataclass
from time import sleep


@dataclass
class Ad:
    campaign_id: str
    creative_url: str
    tier: int
    score: float


class AdSource:
    def __init__(self, name, ads, delay=0.0):
        self.name = name
        self.ads = ads
        self.delay = delay

    def query(self, user_id, region):
        sleep(self.delay)
        return [ad for ad in self.ads if region != "blocked"]


def choose_ad(user_id, region, sources, timeout_seconds=0.08):
    # 每个广告源独立查询;真实系统还应增加频控、预算和隐私校验。
    with ThreadPoolExecutor(max_workers=len(sources)) as pool:
        futures = {
            pool.submit(source.query, user_id, region): source
            for source in sources
        }
        candidates = []
        for future, source in futures.items():
            try:
                candidates.extend(future.result(timeout=timeout_seconds))
            except TimeoutError:
                print(f"skip slow source: {source.name}")

    if not candidates:
        return Ad("house-ad", "https://cdn.example/house.mp4", 99, 0.0)

    # 先按瀑布层级,再按决策分数选择广告。
    return min(candidates, key=lambda ad: (ad.tier, -ad.score))


sources = [
    AdSource("targeted", [Ad("c-101", "https://cdn.example/a.mp4", 1, 0.91)], 0.02),
    AdSource("broad", [Ad("c-202", "https://cdn.example/b.mp4", 2, 0.80)], 0.20),
]

selected = choose_ad("user-42", "IN", sources)
print(selected)

运行示例会优先选择 targeted 层广告,并跳过超过查询预算的 broad 广告;如果所有候选都超时或为空,则返回保底广告。生产实现不应把素材 URL、预算状态和用户身份直接硬编码在进程中,而应接入配置中心、库存服务和隐私控制模块。

服务协调与可观测性

广告决策链路通常会跨越多个团队维护的服务。要让它在流量增长后仍然可控,接口契约和故障边界必须清晰:

  • 决策服务对下游使用短超时和有限重试,避免重试风暴。
  • 每个请求携带 trace ID,并记录各阶段耗时,而不是只记录总耗时。
  • 指标至少覆盖请求量、填充率、P50/P95/P99 延迟、超时率、降级率和各瀑布层命中率。
  • 预算和计费相关写入需要幂等键,防止重试造成重复计费或重复曝光。
  • 用户标识、定向条件和日志字段需要遵循适用的隐私与数据保留要求。

一个实用的排障视角是把“没有返回广告”拆成几类:没有符合条件的候选、预算或频控拒绝、下游超时、素材不可用、播放器取消请求,以及服务内部错误。只有这样,运营团队和工程团队才能区分库存问题与系统问题。

落地时的检查清单

可以按以下顺序推进这类系统的建设:

  • 先定义播放器能接受的端到端时延和各服务的时间预算。
  • 把瀑布层级、保底策略、超时和降级行为写成明确契约。
  • 将 pacing、频控和预算扣减设计为可审计、可幂等的状态变更。
  • 在高峰流量下验证 P99 延迟、重试放大和区域故障场景。
  • 使用真实播放事件校准填充率、广告完成率和节奏控制,而不是只看接口成功率。

这类架构的取舍很明确:更丰富的个性化决策会带来更多服务依赖和状态协调,而更激进的降级会牺牲部分广告价值。合理的目标不是让每个请求都执行完整决策链,而是在播放窗口内稳定返回一条符合策略、可计量、可追踪的广告。


相关推荐