个性化推荐系统全流程讲解附带源码地址
实时推荐系统技术文档
项目名称: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 解决什么问题
在电商场景中,用户的每一次浏览、点击、加购、购买、评分都是一条行为事件。这些事件像水流一样源源不断产生,需要:
- 实时采集:事件产生后毫秒级写入消息队列
- 实时聚合:按时间窗口统计用户行为模式
- 特征提取:将原始事件转化为推荐算法可用的用户偏好特征
- 异常检测:识别刷量机器人等异常行为
本项目实现了这条完整链路的前三个环节。
二、系统架构全景
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 │
└─────────────┴─────────────┴────────┴──────────────┴────────┴───────────┘
这个聚合结果能做什么?
- 实时用户活跃度监控:eventCount 突增 → 用户在密集操作
- 行为模式识别:某用户 30 秒内 VIEW 了 50 次 → 可能在搜索/比较商品
- 异常检测输入:结合 BehaviorAnalysisService 的规则判断是否为机器人
- 推荐时效性特征:用最近窗口的行为权重,衰减历史窗口
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 推荐算法扩展方向
- 协同过滤(Collaborative Filtering):基于 Spark MLlib ALS 矩阵分解,利用现有的 userId × itemId × rating 数据训练
- 实时特征更新:将 Spark 窗口聚合结果写入 Redis 而非 Console,供在线推荐服务查询
- 混合推荐:结合行为权重(隐式反馈)和评分(显式反馈)的混合模型
源码地址: https://github.com/dawdadsd/recommendation_system
更多推荐
所有评论(0)