1. 引言与业务背景在实时数据分析领域,独立访客数(Unique Visitor,UV)是衡量网站或应用用户规模的核心指标。与PV(Page View)不同,UV要求在统计周期内对同一用户的多次访问进行去重。在离线数仓中,UV通常通过COUNT(DISTINCT user_id)配合日期分区完成;但在实时流处理场景下,我们需要在无限数据流上维护一个随时间滑动的时间窗口,并计算窗口内的去重用户数。Apache Spark Structured Streaming 作为统一的流批一体引擎,提供了丰富的窗口函数和状态管理能力。然而,精确去重UV面临两大挑战:状态膨胀:需要为每个窗口保留所有已见用户ID,内存和存储压力巨大。性能开销:精确去重涉及海量数据的Shuffle和全局聚合。为了平衡资源与准确性,业界引入了近似去重算法——HyperLogLog(HLL)。它使用固定大小的内存(如1.2KB)即可估算出数十亿不同元素的基数,误差率可控在2%以内。本文将带领您从零构建一套完整的实时UV计算系统,分别实现:精确方案:基于flatMapGroupsWithState或aggregate+mapGroupsWithState的精确去重。近似方案:基于Spark内置的approx_count_distinct函数或HLL第三方库(如hyperloglog)实现低内存基数估计。所有代码均基于Spark 3.4.x / 3.5.x版本,使用Python(PySpark)编写,并附带详细的性能调优参数。目录1. 引言与业务背景2. 技术选型与架构设计2.1 数据源模拟2.2 窗口类型2.3 状态存储后端2.4 整体流程图3. 环境准备与数据生成器3.1 依赖环境3.2 模拟Kafka数据生成器4. 精确UV计算方案(基于状态的流去重)4.1 核心思路4.2 数据结构设计4.3 完整精确UV代码4.4 生产级精确UV方案(基于ForeachBatch + Redis)5. 近似UV计算方案(HyperLogLog)5.1 HyperLogLog原理简述5.2 使用内置approx_count_distinct5.3 手动实现HyperLogLog UDAF(PySpark)2. 技术选型与架构设计2.1 数据源模拟为聚焦核心逻辑,我们使用Kafka作为实时日志源。每条消息包含:user_id: 字符串,用户唯一标识(或匿名Cookie ID)page_url: 字符串event_time: 时间戳(事件产生时间)ip: 可选字段2.2 窗口类型采用滑动窗口(Sliding Window),例如:窗口长度:10分钟/