技术栈选型与环境搭建
为了确保系统具备高并发处理能力且易于落地,本方案采用Python作为核心开发语言,利用Pandas进行高效的数据清洗与计算,Redis作为推荐结果的缓存层以降低数据库压力,Flask作为轻量级Web服务框架。这种组合既保证了开发效率,又能满足中小规模流量的实时推荐需求。
安装Python依赖库
首先确保你的系统中安装了Python 3.8及以上版本。打开终端,执行以下命令安装必要的依赖库。这些库包含了版本号,以确保环境兼容性:
pip install pandas==1.5.3 redis==4.5.1 flask==2.2.3 numpy==1.24.2
安装并启动Redis服务
推荐使用Docker方式安装Redis,这是最快捷且环境隔离最好的方式。如果你的机器没有安装Docker,请先安装Docker Desktop。执行以下命令拉取并启动Redis:
docker run -d -p 6379:6379 --name redis-recommend redis:latest
执行完毕后,可以通过 docker ps 命令确认容器是否正在运行。如果看到redis-recommend字样,说明服务已就绪。
模拟用户行为数据生成
实战第一步是拥有数据。为了让你直接跑通流程,我们编写一个脚本来模拟电商场景下的用户行为日志。这个脚本会生成包含用户ID、商品ID、行为类型(点击或购买)以及时间戳的CSV文件。
在项目根目录下创建文件 generate_data.py,写入以下完整代码:
```python
import pandas as pd
import random
import time
from datetime import datetime, timedelta
配置参数
USER_COUNT = 1000 模拟用户数量
ITEM_COUNT = 500 模拟商品数量
ACTION_COUNT = 10000 模拟行为总条数
def generate_mock_data():
data = []
base_time = datetime.now()
for i in range(ACTION_COUNT):
user_id = f"user_{random.randint(1, USER_COUNT)}"
item_id = f"item_{random.randint(1, ITEM_COUNT)}"
模拟行为类型:1为点击,2为购买(购买行为占比约20%)
action_type = 2 if random.random() > 0.8 else 1
模拟时间戳,最近7天内的数据
time_offset = random.randint(0, 7 24 3600)
timestamp = base_time - timedelta(seconds=time_offset)
data.append({
"user_id": user_id,
"item_id": item_id,
"action_type": action_type,
"timestamp": timestamp
})
df = pd.DataFrame(data)
按时间排序
df = df.sort_values(by='timestamp')
保存到CSV
df.to_csv('user_behavior.csv', index=False)
print(f"数据生成完毕,共生成 {len(df)} 条数据,已保存至 user_behavior.csv")
if __name__ == "__main__":
generate_mock_data()
```
在终端执行 python generate_data.py。该脚本会在当前目录下生成 user_behavior.csv 文件。这个文件是我们后续构建推荐算法的基础数据源。
基于协同过滤的推荐算法实现
我们将实现一个基于物品的协同过滤算法。这是提升复购率最经典且有效的手段之一。其核心逻辑是:如果用户购买了商品A和商品B,那么A和B之间就建立了一种关联。通过统计所有购买过A的用户同时也购买了哪些商品,可以计算出A的相似商品列表。当用户再次浏览A时,我们向他推荐B。
创建文件 recommender.py,首先编写核心的数据处理与共现矩阵计算逻辑:
```python
import pandas as pd
import redis
import json
class ItemBasedRecommender:
def __init__(self, redis_host='localhost', redis_port=6379):
初始化Redis连接
self.r = redis.StrictRedis(host=redis_host, port=redis_port, decode_responses=True)
self.data_path = 'user_behavior.csv'
def load_data(self):
"""加载数据并过滤出购买行为"""
try:
df = pd.read_csv(self.data_path)
只关注购买行为(action_type == 2),因为购买行为最能体现复购意图
df_buy = df[df['action_type'] == 2]
return df_buy[['user_id', 'item_id']]
except FileNotFoundError:
print("错误:未找到数据文件,请先运行 generate_data.py")
return None
def build_co_occurrence_matrix(self, df_buy):
"""
构建商品共现矩阵
逻辑:同一个用户购买过的商品两两之间关联度+1
"""
按用户分组,获取每个用户购买的所有商品列表
user_items = df_buy.groupby('user_id')['item_id'].apply(list).reset_index()
co_occurrence = {}
for items in user_items['item_id']:
对每个用户的购买列表进行两两组合
如果用户买了 [A, B, C],则生成 (A,B), (A,C), (B,C)
n = len(items)
for i in range(n):
for j in range(i + 1, n):
item_a = items[i]
item_b = items[j]
更新共现计数,无序存储,保证 key 一致
if item_a < item_b:
key = f"{item_a}:{item_b}"
else:
key = f"{item_b}:{item_a}"
if key not in co_occurrence:
co_occurrence[key] = 0
co_occurrence[key] += 1
return co_occurrence
def get_item_recommendations(self, target_item, co_occurrence, top_k=10):
"""
根据共现矩阵计算目标商品的推荐列表
"""
related_items = {}
遍历共现矩阵,找出与 target_item 相关的所有商品
for key, count in co_occurrence.items():
item_a, item_b = key.split(':')
if item_a == target_item:
related_items[item_b] = count
elif item_b == target_item:
related_items[item_a] = count
按共现次数降序排序
sorted_items = sorted(related_items.items(), key=lambda x: x[1], reverse=True)
返回商品ID列表
return [item[0] for item in sorted_items[:top_k]]
```
集成Redis缓存加速查询
直接计算推荐列表耗时较长,且用户行为数据不会每秒都在变,因此必须引入缓存。我们将计算好的推荐结果存入Redis,设置过期时间为1小时。这样,同一个商品的推荐结果在一小时内直接从内存读取,响应速度将从毫秒级降低到微秒级。
继续在 recommender.py 中添加缓存逻辑和对外服务方法:
```python
接上文 ItemBasedRecommender 类中的代码
def recommend_for_item(self, item_id):
"""
获取商品推荐(带缓存逻辑)
"""
cache_key = f"rec:item:{item_id}"
1. 尝试从Redis获取
cached_data = self.r.get(cache_key)
if cached_data:
return json.loads(cached_data)
2. 缓存未命中,重新计算
print(f"缓存未命中,正在计算 {item_id} 的推荐列表...")
df_buy = self.load_data()
if df_buy is None:
return []
注意:在实际生产中,矩阵应该离线计算好存入Redis或数据库
这里为了演示全流程,每次缓存未命中时实时计算(仅适合小数据量)
co_occurrence = self.build_co_occurrence_matrix(df_buy)
recommendations = self.get_item_recommendations(item_id, co_occurrence)
3. 将结果写入Redis,过期时间3600秒
if recommendations:
self.r.setex(cache_key, 3600, json.dumps(recommendations))
return recommendations
def recommend_for_user(self, user_id, top_n=5):
"""
基于用户最近购买的商品,生成推荐列表
逻辑:找到用户最近买的商品,取这些商品的推荐结果的并集
"""
df_buy = self.load_data()
if df_buy is None:
return []
获取该用户的购买记录
user_history = df_buy[df_buy['user_id'] == user_id]['item_id'].tolist()
if not user_history:
return []
取用户最近购买的一个商品作为种子(简化逻辑)
last_bought_item = user_history[-1]
调用单商品推荐
rec_items = self.recommend_for_item(last_bought_item)
简单去重:排除用户已经买过的商品
final_rec = [item for item in rec_items if item not in user_history]
return final_rec[:top_n]
```
构建Web服务接口

我们需要一个HTTP接口供前端或业务系统调用。我们将创建一个Flask应用,暴露一个 /api/recommend 接口。创建文件 app.py:
```python
from flask import Flask, jsonify, request
from recommender import ItemBasedRecommender
app = Flask(__name__)
初始化推荐器实例
recommender = ItemBasedRecommender()
@app.route('/api/recommend', methods=['GET'])
def get_recommendations():
"""
接口参数:
- user_id: 用户ID (必填)
"""
user_id = request.args.get('user_id')
if not user_id:
return jsonify({"error": "缺少必填参数 user_id"}), 400
try:
调用推荐逻辑
result = recommender.recommend_for_user(user_id)
return jsonify({
"status": "success",
"user_id": user_id,
"recommendations": result_list
})
except Exception as e:
return jsonify({"status": "error", "message": str(e)}), 500
if __name__ == '__main__':
启动服务,监听在5000端口
app.run(host='0.0.0.0', port=5000, debug=True)
``>
注意:上述代码中为了确保严谨性,修正了返回变量名为 result_list,请确保在复制时保持一致。完整的逻辑是:获取用户ID -> 传入推荐器 -> 返回JSON。
全流程实操验证
现在我们将所有模块串联起来进行测试。请严格按照以下顺序操作:
1. 生成基础数据
确保当前目录下有 generate_data.py,执行命令生成CSV数据:
python generate_data.py
看到“数据生成完毕”提示后,检查目录下是否多出了 user_behavior.csv。
2. 启动Web服务
执行以下命令启动Flask API服务:
python app.py
终端应显示 Running on http://0.0.0.0:5000,表示服务启动成功。
3. 发起请求测试
打开一个新的终端窗口(不要关闭运行Flask的窗口),使用curl命令模拟用户请求。我们可以使用数据生成脚本中定义的用户ID格式,例如 user_1:
curl "http://127.0.0.1:5000/api/recommend?user_id=user_1"
4. 验证结果与缓存机制
第一次请求时,你会注意到Flask的终端日志中打印了“缓存未命中,正在计算...”,这是因为系统正在实时计算共现矩阵。返回的JSON结果中 recommendations 字段将包含推荐的商品ID列表。
紧接着再次执行相同的curl命令。你会发现终端不再打印“缓存未命中”,且响应速度明显变快。这证明Redis缓存层已成功生效。
5. 复购逻辑验证
为了验证复购逻辑,你可以查看 user_behavior.csv,找到 user_1 购买过的最后一个商品(假设是 item_10)。然后在数据中查找其他同时也购买了 item_10 的用户,看看他们还买了什么。API返回的结果应该正是这些“买了也买”的商品。这就是利用协同过滤提升站内流量复购的核心原理。