本系统采用微服务架构,主要包含以下四个核心服务:
技术栈选择:Python(算法开发)、Redis(实时特征存储)、MySQL(元数据存储)、Django/FastAPI(服务框架)。
在Ubuntu 20.04 LTS系统上执行以下命令,安装Python环境及依赖:
``` sudo apt update sudo apt install python3.8 python3-pip redis-server mysql-server -y pip3 install numpy pandas scikit-learn django redis mysql-connector-python ```登录MySQL,创建数据库和核心表:
``` mysql -u root -p ``` ``` CREATE DATABASE IF NOT EXISTS recommend_db DEFAULT CHARSET utf8mb4; USE recommend_db; -- 用户行为日志表 CREATE TABLE user_behavior ( id BIGINT AUTO_INCREMENT PRIMARY KEY, user_id INT NOT NULL, item_id INT NOT NULL, behavior_type ENUM('view', 'cart', 'buy') NOT NULL, behavior_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP, INDEX idx_user_item (user_id, item_id), INDEX idx_time (behavior_time) ) ENGINE=InnoDB; -- 商品元信息表 CREATE TABLE item_meta ( item_id INT PRIMARY KEY, category_id INT NOT NULL, price DECIMAL(10, 2), title VARCHAR(255) ) ENGINE=InnoDB; ```创建脚本 `build_matrix.py`,从数据库读取行为数据,将行为映射为数值评分(浏览=1,加购=3,购买=5):
``` import pandas as pd from mysql.connector import connect import numpy as np def load_behavior_data(): conn = connect(host='localhost', user='root', password='your_password', database='recommend_db') query = """ SELECT user_id, item_id, behavior_type FROM user_behavior WHERE behavior_time > DATE_SUB(NOW(), INTERVAL 30 DAY) """ df = pd.read_sql(query, conn) conn.close() 行为映射 behavior_score = {'view': 1, 'cart': 3, 'buy': 5} df['score'] = df['behavior_type'].map(behavior_score) return df def create_user_item_matrix(df): 对同一用户-商品对,取最高分(如既浏览又购买,则记为购买分) df_agg = df.groupby(['user_id', 'item_id'])['score'].max().reset_index() 构建稀疏矩阵 matrix = df_agg.pivot(index='user_id', columns='item_id', values='score').fillna(0) return matrix if __name__ == '__main__': df = load_behavior_data() user_item_matrix = create_user_item_matrix(df) 保存矩阵,供后续计算使用 user_item_matrix.to_pickle('user_item_matrix.pkl') ```创建脚本 `calculate_similarity.py`,使用余弦相似度计算用户之间的相似度:
``` from scipy.spatial.distance import cosine import pickle def calculate_user_similarity(matrix): n_users = matrix.shape[0] user_ids = matrix.index.tolist() similarity_matrix = np.zeros((n_users, n_users)) for i in range(n_users): for j in range(i+1, n_users): 计算余弦相似度(1 - 余弦距离) sim = 1 - cosine(matrix.iloc[i].values, matrix.iloc[j].values) similarity_matrix[i, j] = sim similarity_matrix[j, i] = sim 自身相似度为1 similarity_matrix[i, i] = 1.0 将相似度矩阵与用户ID对应存储 sim_df = pd.DataFrame(similarity_matrix, index=user_ids, columns=user_ids) return sim_df if __name__ == '__main__': with open('user_item_matrix.pkl', 'rb') as f: matrix = pickle.load(f) sim_df = calculate_user_similarity(matrix) sim_df.to_pickle('user_similarity.pkl') ```创建FastAPI应用 `recommend_server.py`,提供实时推荐接口:
``` from fastapi import FastAPI import pickle import pandas as pd import redis import json app = FastAPI() 加载数据 with open('user_similarity.pkl', 'rb') as f: sim_matrix = pickle.load(f) with open('user_item_matrix.pkl', 'rb') as f: user_item_matrix = pickle.load(f) 连接Redis,用于缓存用户最近一次推荐结果,减轻数据库压力 r = redis.Redis(host='localhost', port=6379, db=0, decode_responses=True) def get_top_similar_users(user_id, top_k=10): """获取与目标用户最相似的前K个用户""" if user_id not in sim_matrix.index: return [] user_sim_series = sim_matrix.loc[user_id].sort_values(ascending=False) 排除自己 top_users = user_sim_series.iloc[1:top_k+1].index.tolist() return top_users def recommend_items(user_id, top_n=20): """为核心用户生成推荐商品列表""" 1. 先检查Redis缓存(缓存5分钟) cache_key = f'rec:{user_id}' cached = r.get(cache_key) if cached: return json.loads(cached) 2. 找到相似用户 similar_users = get_top_similar_users(user_id) if not similar_users: return [] 3. 聚合相似用户喜欢的商品(加权评分) similar_users_matrix = user_item_matrix.loc[similar_users] 计算加权平均分(相似度作为权重) user_sim_values = sim_matrix.loc[user_id, similar_users].values.reshape(-1, 1) weighted_scores = similar_users_matrix.values.T user_sim_values avg_scores = weighted_scores.sum(axis=1) / user_sim_values.sum() 4. 获取目标用户已交互过的商品,进行过滤 interacted_items = set(user_item_matrix.loc[user_id][user_item_matrix.loc[user_id] > 0].index) item_scores = pd.Series(avg_scores, index=similar_users_matrix.columns) 过滤掉已交互的商品,按加权分降序排序 candidate_items = item_scores[~item_scores.index.isin(interacted_items)] top_items = candidate_items.sort_values(ascending=False).head(top_n).index.tolist() 5. 将结果存入Redis缓存,过期时间300秒 r.setex(cache_key, 300, json.dumps(top_items)) return top_items @app.get("/recommend/{user_id}") async def get_recommendation(user_id: int): items = recommend_items(user_id) return {"user_id": user_id, "recommended_items": items} ```
启动服务:
``` uvicorn recommend_server:app --host 0.0.0.0 --port 8000 --reload ```在您的电商应用后端(如订单服务或用户中心服务),当需要给用户推送商品时,调用推荐接口:
``` import requests def fetch_recommendations_for_user(user_id): try: resp = requests.get(f'http://localhost:8000/recommend/{user_id}', timeout=2) if resp.status_code == 200: data = resp.json() return data['recommended_items'] else: return [] except requests.exceptions.RequestException: 降级策略:返回热门商品或空列表 return get_fallback_items() def get_fallback_items(): """备选方案:返回近期销量最高的商品""" 实现略,可从数据库直接查询 return [1001, 1002, 1003] ```安装Gunicorn,并配置多进程提高并发能力:
``` pip3 install gunicorn ```创建启动脚本 `start.sh`:
``` !/bin/bash gunicorn recommend_server:app \ --workers 4 \ --worker-class uvicorn.workers.UvicornWorker \ --bind 0.0.0.0:8000 \ --timeout 120 \ --access-logfile ./access.log \ --error-logfile ./error.log ```赋予执行权限并启动:
``` chmod +x start.sh nohup ./start.sh & ```在推送网关中,对用户进行分流(50%走推荐算法,50%走原规则),核心对比指标:
分流实现代码示例:
``` import random def should_use_recommendation(user_id): 简单哈希分桶 bucket = user_id % 100 return bucket < 50 50%的用户进入实验组 ```易频IT社区是综合性互联网IT技术门户网站,专注分享网络技术、服务器运维、网络安全、编程开发、系统架构、云计算、大数据等行业干货,实时更新IT行业资讯、零基础教程、实战案例,为IT从业者、技术爱好者提供专业的学习交流平台。
Copyright © 2021-2026 易频IT社区. All Rights Reserved. 备案号:闽ICP备2023013482号 网站地图