实时推荐系统技术文档

项目名称:Spark + Kafka 实时用户行为分析与推荐系统
作者:xiaowu
技术栈:Spring Boot 3.5.3 · Java 25 · Apache Kafka 3.7 · Apache Spark 4.0 · DDD 架构
文档版本:v1.0


一、为什么设计这个项目

1.1 学习动机

这个项目不是一个"拿来即用"的业务系统,而是一个为学习而设计的工程实验室。目标是通过亲手搭建一个完整的实时流处理管道,深入理解以下几个领域:

学习领域具体收获
流式计算Spark Structured Streaming 的窗口聚合、Watermark、Trigger 机制
消息中间件Kafka Producer 异步发送、分区策略、序列化方案
领域驱动设计DDD 四层架构、聚合根、值对象、领域服务的工程落地
并发编程CAS 原子操作、状态机防并发、线程生命周期管理
Spring 生态SmartLifecycle 容器感知、依赖注入、配置管理

1.2 解决什么问题

在电商场景中,用户的每一次浏览、点击、加购、购买、评分都是一条行为事件。这些事件像水流一样源源不断产生,需要:

  1. 实时采集:事件产生后毫秒级写入消息队列
  2. 实时聚合:按时间窗口统计用户行为模式
  3. 特征提取:将原始事件转化为推荐算法可用的用户偏好特征
  4. 异常检测:识别刷量机器人等异常行为

本项目实现了这条完整链路的前三个环节。


二、系统架构全景

2.1 整体数据流

┌──────────────────────────────────────────────────────────────────────┐
│                         REST API 层                                 │
│  POST /api/stream/start    启动模拟器 + Spark                        │
│  POST /api/stream/stop     停止所有流处理                             │
│  GET  /api/stream/status   查看运行状态                               │
│  POST /api/stream/event    手动注入单条事件(调试)                     │
└────────────────────────────────┬─────────────────────────────────────┘
                                 │
                                 ▼
┌──────────────────────────────────────────────────────────────────────┐
│                    StreamSimulatorService                            │
│  以 N 条/秒 的频率持续生成随机用户行为事件                               │
│  用户池: userId 1~100    商品池: itemId 1~500                         │
│  行为类型: VIEW / CLICK / ADD_TO_CART / PURCHASE / RATE              │
└────────────────────────────────┬─────────────────────────────────────┘
                                 │ BehaviorEventDTO (JSON)
                                 ▼
┌──────────────────────────────────────────────────────────────────────┐
│                    Kafka Producer                                    │
│  Topic: user-events                                                  │
│  Key: userId(保证同一用户的事件落入同一分区,维护顺序性)                 │
│  Value: JSON 字符串(ISO-8601 时间格式)                               │
│  发送模式: 异步 + CompletableFuture 回调                              │
└────────────────────────────────┬─────────────────────────────────────┘
                                 │
                        ┌────────┴────────┐
                        │  Kafka Broker   │
                        │  (localhost:9092)│
                        └────────┬────────┘
                                 │
                                 ▼
┌──────────────────────────────────────────────────────────────────────┐
│                 Spark Structured Streaming                           │
│  1. 从 Kafka 消费原始 JSON                                           │
│  2. Schema 解析 + 时间戳转换                                         │
│  3. 30 秒滚动窗口聚合(Watermark 1 分钟)                             │
│  4. 每 10 秒触发一次输出                                              │
│  输出: windowStart | windowEnd | userId | behaviorType | count | avg │
└──────────────────────────────────────────────────────────────────────┘

2.2 DDD 分层架构

xiaowu.backed/
│
├── interfaces/              ← 接口层:对外暴露 REST API
│   └── rest/
│       └── StreamController
│
├── application/             ← 应用层:编排领域对象,不含业务规则
│   ├── dto/
│   │   └── BehaviorEventDTO        ← 数据传输对象
│   └── service/
│       └── StreamSimulatorService  ← 模拟器(应用服务)
│
├── domain/                  ← 领域层:核心业务逻辑,零外部依赖
│   └── eventburial/
│       ├── aggregate/
│       │   └── UserBehaviorAggregate  ← 聚合根
│       ├── entity/
│       │   └── BehaviorEvent          ← 实体
│       ├── valueobject/
│       │   ├── BehaviorType           ← 值对象(行为类型枚举)
│       │   └── Rating                 ← 值对象(评分)
│       └── service/
│           └── BehaviorAnalysisService ← 领域服务
│
└── infrastructure/          ← 基础设施层:技术实现细节
    ├── kafka/
    │   ├── BehaviorEventProducer      ← Kafka 生产者
    │   └── KafkaProducerConfig        ← Kafka 配置
    └── spark/
        └── BehaviorStreamProcessor    ← Spark 流处理器

为什么用 DDD?

传统三层架构(Controller → Service → DAO)在简单 CRUD 中够用,但当业务逻辑复杂化后,Service 层会变成"上帝类"——几千行代码、无法测试、无法复用。

DDD 的核心思想是:业务规则应该内聚在领域对象中,而不是散落在 Service 里。例如:

  • “评分必须在 1.0~5.0 之间” → 封装在 Rating 值对象中
  • “5 分钟内同一商品同一行为算重复” → 封装在 UserBehaviorAggregate 聚合根中
  • “1 小时超 1000 次是异常” → 封装在 BehaviorAnalysisService 领域服务中

三、技术流转详解

3.1 事件生成 → Kafka(生产者侧)

流转过程
StreamSimulatorService.sendRandomEvent()
    │
    ├─ buildRandomEvent()          ← 构建随机 BehaviorEventDTO
    │   ├─ userId:  random(1~100)
    │   ├─ itemId:  random(1~500)
    │   ├─ type:    随机 5 种行为之一
    │   ├─ rating:  仅 RATE 类型生成 1.0~5.0
    │   └─ device:  随机 5 种设备之一
    │
    ├─ producer.sendBehaviorEvent(event)
    │   ├─ ObjectMapper → JSON 序列化
    │   │   └─ Instant → "2026-03-02T10:00:01Z"  (ISO-8601,非时间戳数字)
    │   │
    │   ├─ kafkaTemplate.send(topic, key=userId, value=json)
    │   │   └─ key=userId → 同一用户的事件路由到同一 Partition
    │   │
    │   └─ future.whenComplete(callback)  ← 异步回调记录成功/失败
    │
    └─ total.incrementAndGet()     ← 原子计数器
底层原理:为什么用 userId 做 Kafka Key?

Kafka 的分区策略:partition = hash(key) % numPartitions

userId=1  → hash → Partition 0  ─→ 保证用户1的所有事件有序
userId=2  → hash → Partition 1  ─→ 保证用户2的所有事件有序
userId=3  → hash → Partition 0  ─→ 与用户1同一分区(但不影响各自的顺序)

为什么重要? 推荐系统依赖用户行为的时间顺序。如果同一用户的事件散落在不同分区,消费时无法保证顺序,可能出现"先购买后浏览"的逻辑错乱。

底层原理:Kafka Producer 的批量发送
batch.size=16384     # 16KB 批次缓冲
linger.ms=5          # 最多等 5ms 凑批
buffer.memory=33554432  # 33MB 总缓冲
acks=1               # Leader 确认即可
事件1 ──┐
事件2 ──┤── 缓冲区(等待凑批)──→ 一次网络请求发送整批 ──→ Broker
事件3 ──┤      ↑
事件4 ──┘   5ms 到 或 16KB 满
            以先到者为准

权衡acks=1 意味着只要 Leader 副本确认就算成功,不等 Follower 同步。吞吐量高,但 Leader 宕机时可能丢几条消息。对于行为分析场景,偶尔丢几条浏览事件可以接受。

3.2 Kafka → Spark Structured Streaming(消费者侧)

流转过程
Kafka Topic (user-events)
    │
    ▼
readStream().format("kafka")
    │  option: startingOffsets = "latest"   ← 只消费启动后的新数据
    │  option: failOnDataLoss = "false"     ← 数据丢失不中断流
    │
    ▼
CAST(value AS STRING)                       ← Kafka value 是 byte[],转为字符串
    │
    ▼
from_json(json_str, schema)                 ← 按 StructType 解析 JSON
    │
    ▼
to_timestamp(e.timestamp)                   ← ISO-8601 字符串 → Spark TimestampType
    │
    ▼
filter(userId.isNotNull)                    ← 过滤解析失败的脏数据
    │
    ▼
withWatermark("eventTime", "1 minute")      ← 允许迟到 1 分钟的数据
    │
    ▼
groupBy(window("eventTime", "30s"), userId, behaviorType)
    │
    ▼
agg(count, avg(rating))                    ← 聚合计算
    │
    ▼
writeStream → console (每 10 秒触发)
底层原理:滚动窗口(Tumbling Window)
时间轴: ──────────────────────────────────────────────→

窗口1: [00:00, 00:30)    窗口2: [00:30, 01:00)    窗口3: [01:00, 01:30)
┌──────────────────┐    ┌──────────────────┐    ┌──────────────────┐
│ event1 (00:05)   │    │ event4 (00:35)   │    │ event7 (01:10)   │
│ event2 (00:12)   │    │ event5 (00:42)   │    │                  │
│ event3 (00:28)   │    │ event6 (00:55)   │    │                  │
└──────────────────┘    └──────────────────┘    └──────────────────┘
        ↓                       ↓                       ↓
   count=3, avg=4.2        count=3, avg=3.8        count=1, avg=5.0

每个事件只属于一个窗口(与滑动窗口的区别),窗口之间不重叠。30 秒的窗口大小是在实时性统计显著性之间的权衡:

  • 太短(5 秒)→ 每个窗口事件太少,聚合结果波动大
  • 太长(5 分钟)→ 延迟高,无法及时捕捉用户行为变化
底层原理:Watermark 机制
Watermark = 当前看到的最大事件时间 - 1 分钟

时间轴:
  事件到达:  event(10:05)  event(10:08)  event(10:03)  ← 迟到的!
  Watermark:  09:05         09:08         09:08
                                           ↑
                                    10:03 > 09:08? 否
                                    但 10:03 > watermark(09:08)? 是
                                    → 仍然被接受

  如果迟到事件是 event(08:55):
  08:55 < watermark(09:08)? 是 → 丢弃!太迟了

为什么需要 Watermark? 分布式系统中,事件可能因为网络延迟、设备离线等原因迟到。Watermark 定义了"我愿意等多久"。1 分钟意味着:如果一条事件迟到超过 1 分钟,就不再纳入窗口计算。

底层原理:Trigger 机制
ProcessingTime("10 seconds")

时间轴:
  ├─── 10s ───┤─── 10s ───┤─── 10s ───┤
       ↓            ↓            ↓
   触发输出      触发输出      触发输出
   (处理累积     (处理累积     (处理累积
    的微批)       的微批)       的微批)

Spark Structured Streaming 不是逐条处理,而是微批(Micro-Batch) 模式:每 10 秒收集一批数据,统一计算后输出结果。这是 Spark 在延迟吞吐量之间的经典权衡。

3.3 状态机与生命周期管理

BehaviorStreamProcessor 的状态流转
           start()                    query启动成功
STOPPED ──────────→ STARTING ──────────────→ RUNNING
   ▲                   │                        │
   │              启动失败                   stop()
   │              (finally)                     │
   │                   │                        ▼
   └───────────────────┴──────────────── STOPPING
                    (finally)
为什么用 CAS 而不是 synchronized?
// CAS 方案(当前实现)
if (!state.compareAndSet(State.STOPPED, State.STARTING)) {
    return;  // 另一个线程已经抢先了
}

// synchronized 方案(不推荐)
synchronized(this) {
    if (state != STOPPED) return;
    state = STARTING;
}

两者都能解决并发问题,但 CAS 更优:

  • 无锁:不会阻塞其他线程,失败的线程立即返回
  • 语义清晰:一行代码表达"原子性地从 A 状态转到 B 状态"
  • 适合状态机:状态转换天然是 compare-and-swap 语义
SmartLifecycle 集成
Spring 容器启动
    │
    ├─ isAutoStartup() = false → 不自动启动 Spark(等 REST API 触发)
    │
Spring 容器关闭(Ctrl+C / 部署重启)
    │
    ├─ 按 phase 从大到小关闭
    │   phase = MAX-1 → stop(callback) → query.stop() + thread.interrupt()
    │   phase = 更小值 → Kafka、Redis 等依赖随后关闭
    │
    └─ callback.run() → 通知 Spring "我关完了,继续关下一个"

四、算法实现逻辑

4.1 行为权重模型

系统为每种用户行为赋予不同的隐式反馈权重,反映该行为对用户兴趣的指示强度:

行为类型        权重    理由
──────────────────────────────────────────────
VIEW (浏览)      1     最弱信号:可能只是随便看看
CLICK (点击)     1     略强于浏览,但仍属于探索行为
ADD_TO_CART      3     明确的购买意图信号
RATE (评分)      4     主动评价,表示深度参与
PURCHASE (购买)  5     最强信号:用真金白银投票

设计思路:权重不是线性增长,而是按用户投入成本分级。浏览/点击几乎零成本,加购有心理承诺,评分需要主动操作,购买是最高成本行为。

4.2 事件权重计算

EventWeight = BaseWeight × NormalizedRating

其中:
  BaseWeight = BehaviorType.weight(1 / 1 / 3 / 5 / 4)
  NormalizedRating = (rating - 1.0) / 4.0    归一化到 [0, 1]
                     如果没有评分则不乘(直接用 BaseWeight)

示例计算

用户A 浏览了商品X:     weight = 1(VIEW 基础权重)
用户A 购买了商品Y:     weight = 5(PURCHASE 基础权重)
用户A 给商品Z打了4星:  weight = 4 × (4.0 - 1.0) / 4.0 = 4 × 0.75 = 3.0
用户A 给商品W打了1星:  weight = 4 × (1.0 - 1.0) / 4.0 = 4 × 0.0 = 0.0  ← 差评=零权重

为什么评分要归一化? 原始评分 1~5 直接乘上去会让高分商品权重过大。归一化到 [0,1] 后,评分变成了一个调节因子:5 星 = 满权重,1 星 = 零权重(用户明确表示不喜欢)。

4.3 用户偏好计算(Category Preferences)

public Map<String, Double> calculateCategoryPreferences() {
    // 1. 按类别累加事件权重
    Map<String, Double> preferences = new HashMap<>();
    for (BehaviorEvent event : behaviorEvents) {
        String category = "category_" + (event.getItemId() % 10);  // 简化分类
        preferences.merge(category, event.calculateEventWeight(), Double::sum);
    }

    // 2. 归一化:各类别权重 / 总权重 → 概率分布
    double totalWeight = preferences.values().stream().mapToDouble(v -> v).sum();
    if (totalWeight > 0) {
        preferences.replaceAll((k, v) -> v / totalWeight);
    }
    return preferences;
}

算法图解

用户行为记录:
  VIEW   商品12 (category_2)  → weight = 1.0
  CLICK  商品42 (category_2)  → weight = 1.0
  购买   商品17 (category_7)  → weight = 5.0
  评分   商品42 (category_2)  → weight = 3.0  (4星)

按类别汇总:
  category_2: 1.0 + 1.0 + 3.0 = 5.0
  category_7: 5.0

总权重: 10.0

归一化:
  category_2: 5.0 / 10.0 = 0.50  → 用户 50% 的兴趣在类别2
  category_7: 5.0 / 10.0 = 0.50  → 用户 50% 的兴趣在类别7

这本质上是什么? 这是一个简化的用户画像特征向量。在完整的推荐系统中,这个向量会作为协同过滤或深度学习模型的输入特征。

4.4 Spark 窗口聚合算法

这是整个系统的实时计算核心:

输入:连续不断的 (userId, itemId, behaviorType, rating, eventTime) 流

算法:
  FOR EACH 30秒滚动窗口 [t, t+30s):
    FOR EACH (userId, behaviorType) 组合:
      eventCount = COUNT(*)            ← 该用户在此窗口内的该行为次数
      avgRating  = AVG(rating)         ← 平均评分(仅 RATE 事件有值)

输出:
  ┌─────────────┬─────────────┬────────┬──────────────┬────────┬───────────┐
  │ windowStart │ windowEnd   │ userId │ behaviorType │ count  │ avgRating │
  ├─────────────┼─────────────┼────────┼──────────────┼────────┼───────────┤
  │ 10:00:00    │ 10:00:30    │ 42     │ VIEW         │ 15     │ null      │
  │ 10:00:00    │ 10:00:30    │ 42     │ PURCHASE     │ 2      │ null      │
  │ 10:00:00    │ 10:00:30    │ 7      │ RATE         │ 3      │ 4.33      │
  └─────────────┴─────────────┴────────┴──────────────┴────────┴───────────┘

这个聚合结果能做什么?

  1. 实时用户活跃度监控:eventCount 突增 → 用户在密集操作
  2. 行为模式识别:某用户 30 秒内 VIEW 了 50 次 → 可能在搜索/比较商品
  3. 异常检测输入:结合 BehaviorAnalysisService 的规则判断是否为机器人
  4. 推荐时效性特征:用最近窗口的行为权重,衰减历史窗口

4.5 异常行为检测算法

public boolean detectAnomalousPattern(UserBehaviorAggregate userBehavior) {
    List<BehaviorEvent> recentEvents = userBehavior.getRecentEvents(1);  // 最近1小时

    // 规则1: 绝对频次阈值
    if (recentEvents.size() > 1000) return true;

    // 规则2: 行为多样性检测
    Set<BehaviorType> uniqueTypes = recentEvents.stream()
            .map(BehaviorEvent::getBehaviorType)
            .collect(Collectors.toSet());
    if (recentEvents.size() > 50 && uniqueTypes.size() == 1) return true;

    return false;
}

决策树图解

                获取最近1小时事件
                       │
                事件数 > 1000?
               /            \
             是               否
             │                │
          异常!           事件数 > 50?
        (暴力刷量)        /          \
                        是            否
                        │             │
                  行为类型只有1种?   正常
                   /          \
                 是            否
                 │             │
              异常!          正常
           (单一行为刷量)

为什么用两条规则?

  • 规则 1(> 1000 次/小时):抓暴力机器人,无论行为是否多样,量就不对
  • 规则 2(> 50 次且单一行为):抓精细机器人,量不大但行为模式不像人类(正常用户不会连续点击 50 次同一类行为)

4.6 去重算法

private boolean isDuplicateEvent(BehaviorEvent newEvent) {
    Instant fiveMinutesAgo = newEvent.getTimestamp().minus(5, ChronoUnit.MINUTES);
    return behaviorEvents.stream()
            .filter(e -> e.getItemId().equals(newEvent.getItemId()))
            .filter(e -> e.getBehaviorType().equals(newEvent.getBehaviorType()))
            .anyMatch(e -> e.getTimestamp().isAfter(fiveMinutesAgo));
}

判定条件:5 分钟内,同一用户对同一商品同一行为算重复。

用户A, 商品X, VIEW, 10:00:00  → ✅ 接受
用户A, 商品X, VIEW, 10:03:00  → ❌ 重复(3分钟内,同商品同行为)
用户A, 商品X, CLICK, 10:03:00 → ✅ 接受(行为类型不同)
用户A, 商品Y, VIEW, 10:03:00  → ✅ 接受(商品不同)
用户A, 商品X, VIEW, 10:06:00  → ✅ 接受(超过5分钟窗口)

4.7 内存管理策略

public void addBehaviorEvent(BehaviorEvent event) {
    // ... 验证和去重 ...
    behaviorEvents.add(event);
    if (behaviorEvents.size() > 10000) {
        removeOldestEvents();  // 按时间排序,删除最老的1000条
    }
}
事件列表(按时间排序):
┌─────────────────────────────────────────────────────┐
│ event_1  event_2  ...  event_9999  event_10000  NEW │  ← 超过10000
└─────────────────────────────────────────────────────┘
       ↑                                          ↑
   删除最老的1000条                              保留最新的9001条

结果:
┌─────────────────────────────────────────┐
│ event_1001  event_1002  ...  event_NEW  │  ← 9001条
└─────────────────────────────────────────┘

为什么是 10000 和 1000?

  • 10000 是单用户事件上限,防止内存溢出(100 用户 × 10000 事件 × ~200B/事件 ≈ 200MB)
  • 一次删 1000 而非 1 条,是分摊删除成本:频繁排序删除的 O(n log n) 开销被分摊到每 1000 次插入

五、核心组件设计原理

5.1 Rating 值对象——为什么不直接用 double?

public final class Rating {
    private static final double MIN_RATING = 1.0;
    private static final double MAX_RATING = 5.0;
    private final double value;

    public static Rating of(double value) {
        // 边界校验 + NaN/Infinity 防护
    }

    public double normalize() {
        return (value - MIN_RATING) / (MAX_RATING - MIN_RATING);  // → [0, 1]
    }
}

如果直接用 double rating

  • 调用方可以传 -1.0、100.0、NaN → 运行时才爆炸
  • 每个用到 rating 的地方都要写校验逻辑 → 重复代码
  • (rating - 1) / 4 这个归一化公式散落在多处 → 逻辑泄漏

用值对象后:

  • 构造时就校验,不合法的值根本无法创建 → 编译期安全
  • 归一化逻辑封装在 normalize() 里 → 单一职责
  • Rating.empty() 优雅处理无评分场景 → 消除 null 检查

5.2 BehaviorEvent 实体——Builder 模式

BehaviorEvent event = new BehaviorEvent.Builder()
    .userId(42L)
    .itemId(100L)
    .behaviorType(BehaviorType.PURCHASE)
    .rating(Rating.of(4.5))
    .build();

为什么用 Builder 而非构造函数?

构造函数有 8 个参数,其中部分可选(rating、sessionId、deviceInfo)。如果用构造函数:

// 哪个参数是什么?可读性极差
new BehaviorEvent("id", 42L, 100L, PURCHASE, Rating.of(4.5), now(), "sess", "iOS");

Builder 让每个参数都有名字,且可选参数有默认值,同时在 build() 时做完整性校验。

5.3 Kafka Producer——为什么用 ISO-8601 而非时间戳?

objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS);
// 输出: "2026-03-02T10:00:01.234Z"  而非  1740909601234
维度ISO-8601 字符串Unix 时间戳
可读性人类可直接阅读需要转换
Spark 解析to_timestamp() 直接解析需要手动 from_unixtime()
时区安全自带时区信息 (Z = UTC)无时区,易出错
体积~24 字节~13 字节

对于学习项目,可读性和调试便利性远比那 11 字节的体积差异重要。


六、配置参数说明

# ─── Kafka ───
spring.kafka.bootstrap-servers=localhost:9092          # Kafka 集群地址
kafka.topic.user-events=user-events                    # 行为事件 Topic

# ─── Spark ───
spark.master=local[*]                                  # 本地模式,使用所有 CPU 核心
spark.sql.streaming.checkpointLocation=./spark-checkpoint  # 容错检查点目录

# ─── 推荐引擎(预留) ───
recommendation.default.count=10                        # 默认推荐数量
recommendation.user.min-ratings=5                      # 用户最少评分数(冷启动阈值)
recommendation.similarity.threshold=0.5                # 相似度阈值

七、未来演进方向

7.1 当前状态 vs 完整推荐系统

                  当前已实现                         未来可扩展
              ┌────────────────┐              ┌────────────────────┐
  数据采集    │ ✅ Kafka 生产者  │              │                    │
  实时聚合    │ ✅ Spark 窗口    │              │                    │
  特征工程    │ ✅ 偏好计算      │              │                    │
  异常检测    │ ✅ 规则引擎      │              │                    │
              └────────────────┘              │                    │
  推荐算法    │ ⬜ 待实现        │  ──────→    │ 协同过滤 / ALS      │
  模型训练    │ ⬜ 待实现        │              │ Spark MLlib         │
  在线服务    │ ⬜ 待实现        │              │ Redis 缓存 + API    │
  A/B 测试    │ ⬜ 待实现        │              │ 多策略对比          │
              └────────────────┘              └────────────────────┘

7.2 推荐算法扩展方向

  1. 协同过滤(Collaborative Filtering):基于 Spark MLlib ALS 矩阵分解,利用现有的 userId × itemId × rating 数据训练
  2. 实时特征更新:将 Spark 窗口聚合结果写入 Redis 而非 Console,供在线推荐服务查询
  3. 混合推荐:结合行为权重(隐式反馈)和评分(显式反馈)的混合模型

源码地址: https://github.com/dawdadsd/recommendation_system

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐