本文目录导读:

在PHP项目中实现实时推荐系统,通常需要结合后端缓存、消息队列和前端实时交互技术,由于PHP本身是同步阻塞的脚本语言,纯粹的PHP实现“实时”推荐有局限性,因此通常采用以下架构方案:
核心架构设计(推荐方案)
用户行为 → Nginx → PHP(采集/记录)
→ Redis(实时特征/缓存)
→ 消息队列(Kafka/RabbitMQ)
→ 推荐引擎(Python/Go/Java)
→ 计算结果写回Redis
→ PHP读取Redis → 返回前端
具体实现步骤
1 用户行为采集(PHP层)
<?php
// 用户行为日志采集
function logUserAction($userId, $itemId, $actionType) {
// 1. 写入Redis(实时特征)
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
// 最近N个浏览记录
$redis->lPush("user:{$userId}:history", json_encode([
'item_id' => $itemId,
'time' => time(),
'action' => $actionType
]));
$redis->lTrim("user:{$userId}:history", 0, 99); // 保留最近100条
// 2. 异步写入消息队列(解耦推荐计算)
sendToQueue([
'user_id' => $userId,
'item_id' => $itemId,
'action' => $actionType,
'timestamp' => time()
]);
}
// 异步消息队列发送(使用Redis Stream或RabbitMQ)
function sendToQueue($data) {
// 示例:Redis Stream
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$redis->xAdd('user_actions', '*', $data);
}
2 推荐计算引擎(建议用其他语言)
推荐算法通常在独立的服务中运行,PHP只负责获取结果:
# Python推荐引擎示例(消费消息队列)
import redis
from kafka import KafkaConsumer
from collections import Counter
consumer = KafkaConsumer('user_actions', ...)
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
for msg in consumer:
action = json.loads(msg.value)
user_id = action['user_id']
item_id = action['item_id']
# 1. 基于协同过滤的简单推荐
# 获取与当前用户相似的用户(基于共同行为)
similar_users = get_similar_users(user_id)
# 2. 计算推荐候选集
candidates = set()
for other_user in similar_users:
other_items = r.lrange(f"user:{other_user}:history", 0, -1)
for item_data in other_items:
item = json.loads(item_data)
if item['item_id'] != item_id:
candidates.add(item['item_id'])
# 3. 过滤已看过的内容
seen_items = r.smembers(f"user:{user_id}:seen")
candidates = candidates - seen_items
# 4. 更新推荐结果到Redis(设置过期时间)
top_n = list(candidates)[:10] # 取top10
r.setex(f"user:{user_id}:recommendations", 3600, json.dumps(top_n))
3 PHP获取实时推荐结果
<?php
class RecommendService {
private $redis;
public function __construct() {
$this->redis = new Redis();
$this->redis->connect('127.0.0.1', 6379);
}
/**
* 获取用户实时推荐
* 策略:优先返回实时计算的结果,降级用离线推荐
*/
public function getRealtimeRecommendations($userId, $limit = 10) {
// 1. 尝试获取实时推荐缓存
$cached = $this->redis->get("user:{$userId}:recommendations");
if ($cached) {
return json_decode($cached, true);
}
// 2. 实时推荐不存在,使用实时规则推荐
return $this->fallbackRecommend($userId, $limit);
}
/**
* 降级策略:基于用户近期行为的实时规则推荐
*/
private function fallbackRecommend($userId, $limit) {
// 获取用户最近的浏览历史
$history = $this->redis->lRange("user:{$userId}:history", 0, 20);
if (empty($history)) {
return $this->getPopularItems($limit);
}
// 提取最近的3个物品类别
$recentCategories = [];
foreach (array_slice($history, 0, 3) as $itemData) {
$item = json_decode($itemData, true);
$category = $this->getItemCategory($item['item_id']);
if ($category) {
$recentCategories[] = $category;
}
}
// 同类别下未看过的物品
$recommendations = [];
foreach ($recentCategories as $categoryId) {
$items = $this->getItemsByCategory($categoryId);
$recommendations = array_merge($recommendations, $items);
}
// 去重并排除已看过的
$seen = $this->redis->sMembers("user:{$userId}:seen");
return array_diff($recommendations, $seen);
}
/**
* 获取热门物品(冷启动)
*/
private function getPopularItems($limit) {
// 从Redis热榜获取
$popular = $this->redis->zRevRange('popular_items', 0, $limit - 1, true);
return array_keys($popular);
}
}
4 前端实时推送(WebSocket + SSE)
推荐结果变化时,通过WebSocket或Server-Sent Events推送给用户:
// 使用Ratchet实现PHP WebSocket (or Swoole)
use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface;
class RecommendPush implements MessageComponentInterface {
protected $clients;
public function onOpen(ConnectionInterface $conn) {
$this->clients[$conn->resourceId] = $conn;
// 验证用户身份,建立user_id -> conn映射
}
public function onMessage(ConnectionInterface $from, $msg) {
// 前端发送"recommend_update"事件时,更新推荐
$data = json_decode($msg, true);
if ($data['type'] == 'get_recommend') {
$userId = $data['user_id'];
$recommends = (new RecommendService())->getRealtimeRecommendations($userId);
$from->send(json_encode([
'type' => 'recommend',
'data' => $recommends
]));
}
}
}
前端JS示例(使用SSE):
// EventSource方式(更简单)
const eventSource = new EventSource('/realtime-recommend.php?user_id=123');
eventSource.onmessage = function(e) {
const recommends = JSON.parse(e.data);
updateRecommendUI(recommends);
};
PHP SSE端点:
// realtime-recommend.php
header('Content-Type: text/event-stream');
header('Cache-Control: no-cache');
ini_set('max_execution_time', 0);
$userId = $_GET['user_id'];
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$lastTime = time();
while (true) {
// 轮询Redis看是否有推荐更新
$latestUpdate = $redis->get("user:{$userId}:recommend_updated_at");
if ($latestUpdate && $latestUpdate > $lastTime) {
$recommends = (new RecommendService())->getRealtimeRecommendations($userId);
echo "data: " . json_encode($recommends) . "\n\n";
ob_flush();
flush();
$lastTime = $latestUpdate;
}
sleep(1); // 每秒检查一次
}
性能优化技巧
1 Redis实时特征存储
| 数据结构 | 用途 | 示例 |
|---|---|---|
| Hash | 用户实时特征向量 | user:123:features |
| Sorted Set | 热门物品/实时榜单 | realtime_hot_items |
| List | 用户最近行为 | user:123:history |
| Set | 已看过的物品 | user:123:seen |
| String | 推荐结果缓存 | user:123:recommendations |
2 缓存策略
// 分层缓存
class RecommendCache {
// L1: 本地内存(PHP进程内,如APCu)
private function getL1Cache($key) {
return apcu_fetch($key);
}
// L2: Redis
private function getL2Cache($key) {
$redis = new Redis();
return $redis->get($key);
}
// L3: MySQL(兜底)
private function getMySQLCache($userId) {
// 直接从推荐结果表查询
}
}
3 异步处理架构
PHP请求 → 异步任务队列 → 推荐计算服务器
↓
结果写回Redis
↓
下一次请求直接读Redis
使用消息队列解耦:
// 使用Redis Stream实现简单队列
$redis->xAdd('user_actions', '*', [
'user_id' => $userId,
'item_id' => $itemId,
'action' => 'click',
'timestamp' => time()
]);
// 后台消费者进程
// php consumer.php
while (true) {
$messages = $redis->xRead(['user_actions' => '$'], 1, 0);
// 处理推荐计算...
}
完整的推荐算法示例(PHP实现)
虽然强烈建议将推荐引擎用C++/Go/Python实现,但小型项目可用PHP:
class RealtimeRecommender {
private $redis;
/**
* 协同过滤推荐(物品相似度)
*/
public function itemBasedCF($userId, $limit = 10) {
// 获取用户最近行为
$recentItems = $this->redis->lRange("user:{$userId}:history", 0, 9);
$scores = [];
foreach ($recentItems as $itemJson) {
$item = json_decode($itemJson, true);
$itemId = $item['item_id'];
// 获取与该物品相似的其他物品
$similarItems = $this->redis->zRevRange("item:{$itemId}:similar", 0, 20, true);
foreach ($similarItems as $similarId => $similarity) {
if (!isset($scores[$similarId])) {
$scores[$similarId] = 0;
}
$scores[$similarId] += $similarity;
}
}
// 过滤已看过的
$seen = $this->redis->sMembers("user:{$userId}:seen");
foreach ($seen as $itemId) {
unset($scores[$itemId]);
}
arsort($scores);
return array_slice(array_keys($scores), 0, $limit);
}
/**
* 实时热榜推荐(基于当前时间衰减)
*/
public function trendingNow($limit = 10) {
$currentTime = time();
$score = function($itemId) use ($currentTime) {
// 获取物品在最近N分钟内的行为次数
$recentActions = $this->redis->zCount("item:{$itemId}:actions",
$currentTime - 3600, $currentTime);
// 时间衰减公式:score = actions / (1 + log(time_gap))
return $recentActions / (1 + log(max(1, $currentTime - $itemCreateTime)));
};
// ... 计算并排序
}
}
最终架构建议
对于生产环境的PHP实时推荐系统:
- PHP只做数据收集和展现层,不承担计算任务
- 推荐引擎独立部署(推荐使用Go/Java/Python)
- Redis做实时缓存和特征存储
- 消息队列解耦(Kafka/RabbitMQ/Redis Stream)
- 前端使用WebSocket/SSE实现实时推送
- 离线+实时混合推荐:离线计算预推荐结果,实时修正排序
代码结构示例:
project/
├── app/
│ ├── Http/
│ │ ├── Controllers/
│ │ │ └── RecommendController.php # 接收请求,返回推荐
│ │ └── Middleware/
│ │ └── BehaviorCollector.php # 采集用户行为
│ ├── Jobs/
│ │ └── ProcessUserAction.php # 异步处理用户行为
│ └── Services/
│ ├── RealtimeRecommend.php # 获取实时推荐
│ └── Recommender/ # 推荐算法模块
│ ├── ItemBasedCF.php
│ └── TrendingNow.php
├── config/
│ ├── redis.php
│ └── queue.php
└── worker/
└── recommend_worker.php # 后台推荐计算进程
这种架构可以支撑千万级用户量的实时推荐需求,同时保持PHP开发的效率优势。