基于协同过滤的电商推荐系统实战:从Django集成到Spark离线计算
在实际电商项目中,推荐系统是提升用户粘性和转化率的核心组件。一个典型的推荐系统往往涉及数据采集、算法模型、服务部署和效果评估等多个环节。本文将以一个基于协同过滤算法的电商商品推荐系统为蓝本,串联起从数据爬取、算法实现到Web服务部署的完整技术栈,涵盖Python、Django、Scrapy、Hadoop、Spark等关键技术。我们将重点探讨如何利用协同过滤算法在Django框架中构建推荐服务,并简要介绍如何整合大数据处理工具进行离线计算。无论你是希望理解推荐系统的工程化落地,还是想学习如何将算法模型集成到Web应用中,这篇文章都将提供一个清晰的实践路径。
1. 理解协同过滤算法与推荐系统架构
1.1 推荐系统与协同过滤的核心思想
推荐系统的本质是在信息过载的场景下,通过分析用户的历史行为数据,预测其可能感兴趣的物品,并进行个性化推送。协同过滤是其中应用最广泛、最经典的算法之一,其核心假设是“物以类聚,人以群分”。它主要分为两类:
- 基于用户的协同过滤:找到与目标用户兴趣相似的其他用户,将这些相似用户喜欢而目标用户未接触过的物品推荐给目标用户。
- 基于物品的协同过滤:找到与目标用户历史喜欢物品相似的其他物品,将这些相似物品推荐给目标用户。
在电商场景中,用户对商品的评分、点击、购买、收藏等行为都可以作为计算“相似度”的依据。基于物品的协同过滤在电商中更为常见,因为它更稳定(物品的属性变化通常比用户的兴趣变化慢),且可解释性强(“购买了A商品的用户也购买了B”)。
1.2 系统整体架构设计
一个完整的、可扩展的电商推荐系统通常采用分层架构,将在线服务和离线计算分离。
用户请求 --> Django Web服务 (在线) --> 推荐引擎 (加载模型/规则) --> 返回推荐列表 ↑ 模型更新 ↑ Spark/Hadoop (离线) <-- 数据清洗/特征工程 <-- 原始行为日志- 离线计算层:这是系统的“大脑”。使用Hadoop/Spark处理海量的用户行为日志,进行数据清洗、特征提取,并运行协同过滤等机器学习算法,训练出“用户-物品”评分矩阵或物品相似度矩阵。这个过程计算量大,耗时较长(如每天一次)。
- 在线服务层:这是系统的“手脚”。使用Django等Web框架提供RESTful API。当用户访问时,服务从缓存或数据库中快速加载离线层计算好的模型(如物品相似度矩阵),进行实时查询和排序,生成最终的推荐列表并返回。
- 数据采集层:使用Scrapy等爬虫框架,可以从电商网站(在拥有合法权限的前提下)或内部日志系统采集商品信息、用户行为等原始数据,作为离线计算的输入。
本实践将聚焦于在线服务层的核心实现,即如何在Django中集成一个基于内存的协同过滤推荐引擎,并简要说明离线计算的结果如何被服务层使用。
2. 环境准备与项目初始化
2.1 开发环境与工具清单
在开始编码前,请确保你的本地开发环境已就绪。
| 组件 | 推荐版本 | 用途说明 | 安装/验证命令 |
|---|---|---|---|
| Python | 3.8+ | 项目主语言 | python --version |
| pip | 最新版 | Python包管理 | pip --version |
| Django | 4.x | Web应用框架 | pip install django |
| Virtualenv | - | 创建独立Python环境 | pip install virtualenv |
| IDE/编辑器 | VSCode/PyCharm | 代码编写与调试 | - |
注意:生产环境部署时,还需要考虑Nginx、Gunicorn/uWSGI、Redis、数据库(如PostgreSQL)等组件。本文为简化演示,使用Django自带的SQLite和开发服务器。
2.2 创建Django项目与应用
我们首先创建一个标准的Django项目,并在其中创建专门处理推荐逻辑的应用。
# 1. 创建并进入项目目录 mkdir ecommerce_recommendation && cd ecommerce_recommendation # 2. 创建虚拟环境(强烈推荐,避免包冲突) python -m venv venv # 3. 激活虚拟环境 # Windows: venv\Scripts\activate # Linux/Mac: source venv/bin/activate # 4. 安装Django pip install django # 5. 创建Django项目,项目名为`recsys` django-admin startproject recsys . # 6. 创建推荐系统应用,应用名为`recommender` python manage.py startapp recommender # 7. 创建用于管理算法模型和工具函数的模块 mkdir recommender/algorithms touch recommender/algorithms/__init__.py touch recommender/algorithms/collaborative_filtering.py touch recommender/utils.py创建完成后,项目目录结构应如下所示:
ecommerce_recommendation/ ├── manage.py ├── recsys/ │ ├── __init__.py │ ├── settings.py │ ├── urls.py │ └── wsgi.py ├── recommender/ # 我们的推荐系统应用 │ ├── __init__.py │ ├── admin.py │ ├── apps.py │ ├── migrations/ │ ├── models.py │ ├── tests.py │ ├── views.py │ ├── algorithms/ # 算法模块 │ │ ├── __init__.py │ │ └── collaborative_filtering.py │ └── utils.py # 工具函数 └── venv/ # 虚拟环境目录2.3 配置Django项目基础设置
编辑recsys/settings.py文件,将新创建的应用加入INSTALLED_APPS,并配置数据库(暂用SQLite)。
# recsys/settings.py INSTALLED_APPS = [ 'django.contrib.admin', 'django.contrib.auth', 'django.contrib.contenttypes', 'django.contrib.sessions', 'django.contrib.messages', 'django.contrib.staticfiles', # 添加我们自己的应用 'recommender', ] # 数据库配置(使用SQLite进行快速演示) DATABASES = { 'default': { 'ENGINE': 'django.db.backends.sqlite3', 'NAME': BASE_DIR / 'db.sqlite3', } } # 国际化和时区设置 LANGUAGE_CODE = 'zh-hans' TIME_ZONE = 'Asia/Shanghai' USE_I18N = True USE_TZ = True3. 数据模型设计与协同过滤算法实现
3.1 定义核心数据模型
推荐系统的基础是数据。我们需要在recommender/models.py中定义最基础的模型:用户、商品和用户-商品交互行为(如评分)。
# recommender/models.py from django.db import models from django.contrib.auth.models import User class Product(models.Model): """商品模型""" product_id = models.CharField(max_length=100, unique=True, verbose_name='商品ID') title = models.CharField(max_length=200, verbose_name='商品标题') category = models.CharField(max_length=100, verbose_name='商品类别') price = models.DecimalField(max_digits=10, decimal_places=2, verbose_name='价格') description = models.TextField(blank=True, verbose_name='描述') image_url = models.URLField(blank=True, verbose_name='图片链接') created_at = models.DateTimeField(auto_now_add=True) class Meta: verbose_name = '商品' verbose_name_plural = '商品' def __str__(self): return f'{self.title} ({self.product_id})' class UserInteraction(models.Model): """用户-商品交互行为模型(如评分、点击、购买)""" # 关联Django内置用户,也可自定义用户模型 user = models.ForeignKey(User, on_delete=models.CASCADE, verbose_name='用户') product = models.ForeignKey(Product, on_delete=models.CASCADE, verbose_name='商品') # 交互类型:'view', 'click', 'cart', 'purchase', 'rating' interaction_type = models.CharField(max_length=20, verbose_name='交互类型') # 权重或评分,例如:view=1, click=2, cart=3, purchase=5, rating=1-5 value = models.FloatField(default=1.0, verbose_name='权重/评分') timestamp = models.DateTimeField(auto_now_add=True, verbose_name='交互时间') class Meta: verbose_name = '用户交互' verbose_name_plural = '用户交互' # 一个用户对同一个商品同一种行为只记录一次最新(或最重)的,可根据业务调整 unique_together = ['user', 'product', 'interaction_type'] def __str__(self): return f'{self.user.username} - {self.product.title} - {self.interaction_type}'创建并应用数据库迁移:
python manage.py makemigrations recommender python manage.py migrate3.2 实现基于物品的协同过滤算法
接下来,我们在recommender/algorithms/collaborative_filtering.py中实现一个简单的基于物品的协同过滤算法。这里我们使用余弦相似度计算物品之间的相似性。
# recommender/algorithms/collaborative_filtering.py import numpy as np from collections import defaultdict from django.contrib.auth.models import User from recommender.models import UserInteraction, Product class ItemBasedCF: """ 基于物品的协同过滤推荐器 使用余弦相似度计算物品相似度矩阵,并为目标用户生成推荐。 """ def __init__(self): self.product_similarity_matrix = None # 物品相似度矩阵 self.product_id_to_index = {} # 商品ID到矩阵下标的映射 self.index_to_product_id = {} # 矩阵下标到商品ID的映射 self.user_item_matrix = None # 用户-物品评分矩阵(可选,用于调试) def fit(self, interactions): """ 训练模型:根据用户交互数据计算物品相似度矩阵。 Args: interactions: QuerySet of UserInteraction 或 list of dict """ # 1. 数据准备:构建用户-物品评分字典 user_item_ratings = defaultdict(dict) # {user_id: {product_id: rating}} all_product_ids = set() for interaction in interactions: # 这里简单地将交互权重作为评分,实际业务中可能需要更复杂的加权计算 user_item_ratings[interaction.user_id][interaction.product_id] = interaction.value all_product_ids.add(interaction.product_id) # 2. 建立索引映射 self.product_id_to_index = {pid: idx for idx, pid in enumerate(sorted(all_product_ids))} self.index_to_product_id = {idx: pid for pid, idx in self.product_id_to_index.items()} num_products = len(all_product_ids) # 3. 构建用户-物品评分矩阵 (稀疏矩阵,这里用列表的列表简化表示) # 获取所有用户ID all_user_ids = list(user_item_ratings.keys()) num_users = len(all_user_ids) self.user_item_matrix = np.zeros((num_users, num_products)) for u_idx, user_id in enumerate(all_user_ids): for product_id, rating in user_item_ratings[user_id].items(): p_idx = self.product_id_to_index[product_id] self.user_item_matrix[u_idx, p_idx] = rating # 4. 计算物品相似度矩阵 (余弦相似度) # 矩阵转置,使行代表物品,列代表用户 item_user_matrix = self.user_item_matrix.T # 计算范数 norms = np.linalg.norm(item_user_matrix, axis=1, keepdims=True) # 避免除以零 norms[norms == 0] = 1 # 归一化 item_user_matrix_normalized = item_user_matrix / norms # 计算余弦相似度 self.product_similarity_matrix = np.dot(item_user_matrix_normalized, item_user_matrix_normalized.T) print(f"模型训练完成。共处理 {num_users} 个用户,{num_products} 个商品。") def recommend(self, target_user_id, top_k=10, interacted_items=None): """ 为目标用户生成推荐列表。 Args: target_user_id: 目标用户的ID top_k: 返回推荐商品的数量 interacted_items: 用户已经交互过的商品ID列表(可选,用于排除) Returns: list of (product_id, predicted_score) tuples """ if self.product_similarity_matrix is None: raise ValueError("模型尚未训练,请先调用 fit 方法。") # 获取目标用户的历史交互商品(如果未提供) if interacted_items is None: interactions = UserInteraction.objects.filter(user_id=target_user_id) interacted_items = [inter.product_id for inter in interactions] # 初始化预测评分字典 scores = defaultdict(float) # 遍历用户交互过的每个商品 for interacted_product_id in interacted_items: if interacted_product_id not in self.product_id_to_index: continue # 训练集中未出现的商品,跳过 interacted_idx = self.product_id_to_index[interacted_product_id] # 获取该商品与所有其他商品的相似度 similarities = self.product_similarity_matrix[interacted_idx] # 获取用户对该交互商品的评分(这里简化处理,取最后一次交互的值或固定值) # 实际应从 user_item_matrix 中获取准确评分。此处假设评分为1。 rating = 1.0 # 将相似度加权后累加到候选商品的预测分数上 for other_idx, sim in enumerate(similarities): other_product_id = self.index_to_product_id[other_idx] # 排除用户已经交互过的商品 if other_product_id in interacted_items: continue scores[other_product_id] += sim * rating # 按预测分数排序,返回top_k个推荐 sorted_items = sorted(scores.items(), key=lambda x: x[1], reverse=True) return sorted_items[:top_k] def get_similar_items(self, product_id, top_n=5): """获取与指定商品最相似的商品列表""" if product_id not in self.product_id_to_index: return [] idx = self.product_id_to_index[product_id] # 获取相似度行,并排除自身(相似度为1) sim_scores = list(enumerate(self.product_similarity_matrix[idx])) # 按相似度降序排序 sim_scores = sorted(sim_scores, key=lambda x: x[1], reverse=True) # 返回前top_n个(从1开始,跳过自身) result = [] for i, (other_idx, score) in enumerate(sim_scores[1:top_n+1], 1): result.append((self.index_to_product_id[other_idx], score)) return result这个实现是一个简化版本,适用于中小规模数据。它直接将所有数据加载到内存中计算。对于大规模数据,需要使用Spark MLlib等分布式计算框架进行离线训练,并将生成的相似度矩阵(product_similarity_matrix)和映射关系(product_id_to_index)保存到文件或数据库中,供在线服务加载。
4. 构建Django视图、API与模型管理
4.1 创建管理命令训练模型
为了定期更新推荐模型,我们创建一个Django自定义管理命令。在recommender/management/commands/目录下创建文件train_cf_model.py。
mkdir -p recommender/management/commands touch recommender/management/__init__.py touch recommender/management/commands/__init__.py touch recommender/management/commands/train_cf_model.py# recommender/management/commands/train_cf_model.py import pickle from django.core.management.base import BaseCommand from recommender.algorithms.collaborative_filtering import ItemBasedCF from recommender.models import UserInteraction class Command(BaseCommand): help = '训练基于物品的协同过滤模型并保存到文件' def handle(self, *args, **options): self.stdout.write('开始加载用户交互数据...') # 这里可以添加时间过滤,例如只使用最近90天的数据 interactions = UserInteraction.objects.all().select_related('user', 'product') # 也可以根据交互类型过滤或加权,例如只考虑购买和评分 # interactions = interactions.filter(interaction_type__in=['purchase', 'rating']) self.stdout.write(f'共加载到 {interactions.count()} 条交互记录。') if interactions.count() == 0: self.stdout.write(self.style.WARNING('没有交互数据,无法训练模型。')) return # 初始化并训练模型 cf_model = ItemBasedCF() cf_model.fit(interactions) # 将训练好的模型保存到文件 model_data = { 'similarity_matrix': cf_model.product_similarity_matrix, 'id_to_index': cf_model.product_id_to_index, 'index_to_id': cf_model.index_to_product_id, } with open('cf_model.pkl', 'wb') as f: pickle.dump(model_data, f) self.stdout.write(self.style.SUCCESS(f'模型训练完成,已保存到 cf_model.pkl')) self.stdout.write(f'模型维度: {cf_model.product_similarity_matrix.shape}')运行命令进行训练:
python manage.py train_cf_model4.2 实现推荐API视图
接下来,我们创建Django视图,提供推荐API。编辑recommender/views.py。
# recommender/views.py import pickle import os from django.http import JsonResponse from django.contrib.auth.decorators import login_required from django.views.decorators.http import require_GET from recommender.algorithms.collaborative_filtering import ItemBasedCF from recommender.models import Product # 全局变量存储加载的模型(生产环境应使用缓存如Redis,并处理模型更新) _MODEL = None _MODEL_LOADED = False def load_cf_model(): """加载协同过滤模型""" global _MODEL, _MODEL_LOADED if _MODEL_LOADED and _MODEL is not None: return _MODEL model_path = 'cf_model.pkl' if not os.path.exists(model_path): raise FileNotFoundError(f"模型文件 {model_path} 不存在,请先运行训练命令。") with open(model_path, 'rb') as f: model_data = pickle.load(f) cf_model = ItemBasedCF() cf_model.product_similarity_matrix = model_data['similarity_matrix'] cf_model.product_id_to_index = model_data['id_to_index'] cf_model.index_to_product_id = model_data['index_to_id'] _MODEL = cf_model _MODEL_LOADED = True return cf_model @require_GET @login_required def get_recommendations(request): """为当前登录用户获取商品推荐""" try: cf_model = load_cf_model() except FileNotFoundError as e: return JsonResponse({'error': str(e)}, status=503) except Exception as e: return JsonResponse({'error': '模型加载失败'}, status=500) user_id = request.user.id top_k = request.GET.get('top_k', 10) try: top_k = int(top_k) except ValueError: top_k = 10 try: # 获取推荐结果 (product_id, score) recommendations = cf_model.recommend(target_user_id=user_id, top_k=top_k) # 根据product_id查询商品详细信息 product_ids = [rec[0] for rec in recommendations] products = Product.objects.filter(product_id__in=product_ids) # 保持推荐顺序 product_map = {p.product_id: p for p in products} result = [] for pid, score in recommendations: if pid in product_map: p = product_map[pid] result.append({ 'product_id': p.product_id, 'title': p.title, 'category': p.category, 'price': str(p.price), # Decimal转字符串便于JSON序列化 'image_url': p.image_url, 'score': round(score, 4) # 保留4位小数 }) return JsonResponse({'user_id': user_id, 'recommendations': result}) except Exception as e: # 记录日志 import traceback traceback.print_exc() return JsonResponse({'error': '推荐生成失败'}, status=500) @require_GET def get_similar_items(request, product_id): """获取与指定商品相似的商品""" try: cf_model = load_cf_model() except Exception as e: return JsonResponse({'error': '模型加载失败'}, status=500) top_n = request.GET.get('top_n', 5) try: top_n = int(top_n) except ValueError: top_n = 5 try: similar_items = cf_model.get_similar_items(product_id, top_n=top_n) product_ids = [item[0] for item in similar_items] products = Product.objects.filter(product_id__in=product_ids) product_map = {p.product_id: p for p in products} result = [] for pid, score in similar_items: if pid in product_map: p = product_map[pid] result.append({ 'product_id': p.product_id, 'title': p.title, 'score': round(score, 4) }) return JsonResponse({'product_id': product_id, 'similar_items': result}) except KeyError: return JsonResponse({'error': '商品ID不在模型中'}, status=404) except Exception as e: return JsonResponse({'error': '查询失败'}, status=500)4.3 配置URL路由
在recommender应用下创建urls.py,并配置到项目主路由中。
# recommender/urls.py from django.urls import path from . import views urlpatterns = [ path('api/recommend/', views.get_recommendations, name='get_recommendations'), path('api/similar/<str:product_id>/', views.get_similar_items, name='get_similar_items'), ]# recsys/urls.py from django.contrib import admin from django.urls import path, include urlpatterns = [ path('admin/', admin.site.urls), path('recommend/', include('recommender.urls')), ]4.4 注册模型到Admin后台
为了方便在开发阶段管理商品和交互数据,我们将模型注册到Django Admin。
# recommender/admin.py from django.contrib import admin from .models import Product, UserInteraction @admin.register(Product) class ProductAdmin(admin.ModelAdmin): list_display = ('product_id', 'title', 'category', 'price') search_fields = ('product_id', 'title', 'category') list_filter = ('category',) @admin.register(UserInteraction) class UserInteractionAdmin(admin.ModelAdmin): list_display = ('user', 'product', 'interaction_type', 'value', 'timestamp') list_filter = ('interaction_type', 'timestamp') search_fields = ('user__username', 'product__title')创建超级用户并访问/admin后台:
python manage.py createsuperuser python manage.py runserver访问http://127.0.0.1:8000/admin登录后即可管理数据。
5. 系统运行验证与测试
5.1 准备测试数据
为了测试推荐功能,我们需要先注入一些模拟数据。可以编写一个简单的脚本或使用Django Shell。
python manage.py shell在Django Shell中执行以下代码(示例):
from django.contrib.auth.models import User from recommender.models import Product, UserInteraction import random # 创建一些测试用户 users = [] for i in range(1, 6): user, created = User.objects.get_or_create(username=f'test_user_{i}') users.append(user) # 创建一些测试商品 products = [] categories = ['电子产品', '图书', '家居', '服饰'] for i in range(1, 21): product, created = Product.objects.get_or_create( product_id=f'P{i:03d}', defaults={ 'title': f'测试商品 {i}', 'category': random.choice(categories), 'price': round(random.uniform(10, 1000), 2), 'description': f'这是测试商品{i}的描述。' } ) products.append(product) # 创建用户-商品交互记录(模拟评分) interaction_types = ['view', 'click', 'cart', 'purchase', 'rating'] for user in users: # 每个用户随机与5-15个商品交互 interacted_products = random.sample(products, random.randint(5, 15)) for product in interacted_products: # 随机生成一种交互类型和值 inter_type = random.choice(interaction_types) if inter_type == 'rating': value = random.randint(1, 5) # 评分1-5分 else: value = 1.0 # 其他行为权重为1 UserInteraction.objects.get_or_create( user=user, product=product, interaction_type=inter_type, defaults={'value': value} ) print("测试数据创建完成!")5.2 训练模型并调用API
训练模型:运行之前创建的管理命令。
python manage.py train_cf_model控制台应输出类似信息:“模型训练完成。共处理 5 个用户,20 个商品。模型训练完成,已保存到 cf_model.pkl”。
启动开发服务器:
python manage.py runserver测试推荐API:
- 首先,使用你创建的超级用户或测试用户登录Django Admin (
http://127.0.0.1:8000/admin),获取其用户ID。 - 然后,使用浏览器或
curl命令访问推荐API。由于视图使用了@login_required装饰器,你需要先通过会话登录。更常用的方式是使用API认证(如Token),这里为简化,我们可以先临时注释掉@login_required进行测试,或者使用Django的测试客户端。 - 简化测试:暂时修改
views.py中的get_recommendations视图,移除@login_required,并硬编码一个测试用户ID。# 临时修改,测试后请恢复 # @login_required def get_recommendations(request): # user_id = request.user.id # 注释掉 user_id = 1 # 使用测试用户ID,例如ID为1的用户 ... - 访问
http://127.0.0.1:8000/recommend/api/recommend/?top_k=5,你应该能收到一个JSON响应,包含为用户ID为1的用户推荐的5个商品列表。 - 测试相似商品API:访问
http://127.0.0.1:8000/recommend/api/similar/P001/,获取与商品P001相似的商品。
- 首先,使用你创建的超级用户或测试用户登录Django Admin (
5.3 验证结果与逻辑检查
收到API响应后,需要人工检查推荐结果是否合理。
- 推荐给用户的商品,不应该是该用户已经有过交互的商品(除非业务允许重复推荐)。
- 相似商品,应该与目标商品在类别或其他特征上具有一定关联性(由于我们使用的是随机生成的模拟数据,关联性可能不明显,这验证了数据质量对算法效果的重要性)。
6. 工程化扩展:集成Scrapy爬虫与Spark离线计算
6.1 使用Scrapy构建商品爬虫(概念与结构)
在实际项目中,商品信息库需要维护。我们可以使用Scrapy定期爬取电商网站(需遵守robots.txt及法律法规)来更新Product表。
一个简化的Scrapy项目结构可能如下:
scrapy_crawler/ ├── scrapy.cfg └── ecommerce_crawler/ ├── __init__.py ├── items.py # 定义爬取的数据结构 ├── middlewares.py ├── pipelines.py # 数据处理管道,用于保存到Django数据库 ├── settings.py └── spiders/ ├── __init__.py └── product_spider.py # 具体的爬虫关键点在于pipelines.py,它需要能够与Django项目交互,将爬取到的Product数据保存到数据库中。这通常通过设置环境变量DJANGO_SETTINGS_MODULE并调用django.setup()来实现。
# pipelines.py 示例片段 import os import sys import django # 将Django项目路径加入系统路径,并设置Django环境 sys.path.append('/path/to/your/ecommerce_recommendation') os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'recsys.settings') django.setup() from recommender.models import Product class DjangoPipeline: def process_item(self, item, spider): # 使用update_or_create避免重复创建 product, created = Product.objects.update_or_create( product_id=item['product_id'], defaults={ 'title': item['title'], 'category': item['category'], 'price': item['price'], # ... 其他字段 } ) spider.logger.info(f"{'Created' if created else 'Updated'} product: {product.product_id}") return item6.2 使用Spark进行大规模离线计算
当用户和商品数量达到百万甚至千万级时,单机的Python算法将无法满足性能和内存需求。此时需要将协同过滤算法的训练任务迁移到Spark集群上。
核心思路:
- 数据准备:将
UserInteraction表的数据定期导出到HDFS或直接通过Spark JDBC读取。 - Spark MLlib 训练:使用
pyspark.ml.recommendation.ALS(交替最小二乘法) 或自定义的分布式相似度计算。 - 模型输出:将训练得到的物品相似度矩阵(或用户/物品隐因子向量)保存为Parquet/CSV文件或写入HBase/Redis。
- Django服务加载:在线服务启动时,从存储中加载处理好的模型数据(如相似度矩阵),而不是重新计算。
一个简化的Spark ALS示例:
# spark_als_train.py from pyspark.sql import SparkSession from pyspark.ml.recommendation import ALS from pyspark.sql.functions import col # 初始化Spark会话 spark = SparkSession.builder \ .appName("EcommerceCF") \ .config("spark.executor.memory", "4g") \ .getOrCreate() # 1. 加载数据 (假设数据格式: user_id, product_id, rating) df = spark.read.parquet("hdfs://path/to/user_interactions.parquet") # 或从数据库读取 # df = spark.read.format("jdbc").options(...).load() # 2. 使用ALS算法训练模型 als = ALS( maxIter=10, regParam=0.01, userCol="user_id", itemCol="product_id", ratingCol="rating", coldStartStrategy="drop" # 处理冷启动问题 ) model = als.fit(df) # 3. 获取物品因子向量,可用于计算物品相似度 item_factors = model.itemFactors # DataFrame: [id, features] item_factors.write.parquet("hdfs://path/to/output/item_factors.parquet") # 4. 也可以直接为所有用户生成推荐 # user_recs = model.recommendForAllUsers(10) # user_recs.write.parquet("hdfs://path/to/output/user_recommendations.parquet") spark.stop()Django服务则需要修改load_cf_model函数,改为从Spark输出的文件(或经过进一步处理生成的相似度矩阵文件)中加载数据。
7. 常见问题排查与优化实践
7.1 部署与运行常见问题
| 问题现象 | 可能原因 | 检查与解决思路 |
|---|---|---|
运行train_cf_model报内存错误 | 交互数据量太大,内存不足。 | 1. 增加数据过滤(如只取最近数据)。 2. 使用Spark等分布式计算。 3. 优化算法,使用稀疏矩阵库(如 scipy.sparse)。 |
| 推荐API返回“模型文件不存在” | cf_model.pkl文件未生成或路径错误。 | 1. 确认已成功运行训练命令。 2. 检查Django应用的工作目录。 3. 在 load_cf_model中使用绝对路径。 |
| 推荐结果总是空列表或重复 | 1. 用户没有历史行为。 2. 相似度矩阵计算有误或全为零。 3. 推荐逻辑过滤了所有商品。 | 1. 为新用户实施冷启动策略(如推荐热门商品、基于用户属性推荐)。 2. 检查训练数据质量和相似度计算代码。 3. 调试 recommend方法,打印中间变量。 |
| API响应速度慢 | 1. 模型加载每次请求都进行。 2. 相似度矩阵过大,查询慢。 | 1. 使用全局变量+锁或Django缓存框架缓存模型。 2. 将相似度矩阵存入Redis等内存数据库,按需查询。 3. 对热门商品或用户的推荐结果进行缓存。 |
| 新商品/新用户无法被推荐(冷启动) | 协同过滤依赖历史行为,新实体无数据。 | 实施混合推荐策略: 1.基于内容的推荐:利用商品属性(类别、标签)计算相似度。 2.热门推荐:推荐近期最畅销或最热门的商品。 3.探索与利用:在推荐结果中混入少量随机商品。 |
7.2 性能与生产环境最佳实践
模型更新策略:
- 全量更新:每天低峰期(如凌晨)使用Spark任务重新训练全量模型。
- 增量更新:对于基于物品的CF,新用户行为对物品相似度影响较小,可定期(如每小时)用小批量数据微调或只更新受影响的部分相似度。
- 模型版本化:每次训练生成新版本模型文件,在线服务通过配置或信号平滑切换,便于回滚。
在线服务优化:
- 缓存层:使用Redis缓存热门推荐结果、用户画像、商品相似度列表。为缓存设置合理的TTL。
- 异步计算:对于实时性要求不高的“猜你喜欢”列表,可以使用Celery等异步任务队列提前计算好并缓存。
- API设计:推荐API应支持分页、AB测试参数(
algorithm_version)、场景参数(scene=homepage/cart)等。
数据与算法监控:
- 数据质量监控:监控用户行为日志的采集覆盖率、异常值(如异常高评分)。
- 算法效果监控:通过A/B测试对比不同推荐策略的CTR(点击率)、转化率等业务指标。离线评估指标如准确率、召回率、覆盖率。
- 服务健康监控:监控API响应时间、错误率、模型加载状态。
从协同过滤到深度学习: 协同过滤是入门首选,但工业级系统通常会向更复杂的模型演进:
- 特征工程:引入用户画像(年龄、地域)、商品属性(价格、品牌)、上下文(时间、地点)作为特征。
- 模型升级:尝试矩阵分解(MF)、因子分解机(FM)、深度神经网络(如YouTube DNN、Wide & Deep)等。
- 在线学习:使用Flink等流处理框架进行近实时特征计算和模型更新。
8. 总结与扩展方向
本文详细演示了如何在Django框架中,从零构建一个基于协同过滤算法的电商商品推荐系统。我们完成了从数据模型设计、核心算法实现、Django API暴露到基础测试的完整闭环。这个系统虽然简单,但清晰地勾勒出了推荐系统在线服务层的核心架构。
对于希望进一步深入和实践的开发者,可以从以下几个方向进行扩展:
- 完善数据管道:使用Scrapy或日志采集工具,构建稳定、自动化的商品与用户行为数据流入通道。
- 集成大数据计算引擎:将
ItemBasedCF.fit中的训练逻辑,用PySpark在Hadoop/YARN集群上重写,处理千万级以上的数据。 - 构建混合推荐系统:在协同过滤的基础上,加入基于内容的推荐(利用商品标题、描述做文本相似度计算)和热门推荐,有效解决冷启动问题。
- 引入评估体系:在推荐API中埋点,记录每次推荐曝光和用户点击,构建离线评估(准确率、召回率)和在线A/B测试平台,用数据驱动算法迭代。
- 服务化与部署:将推荐模块拆分为独立的微服务,使用Django REST Framework提供更规范的API,并通过Docker容器化,使用Kubernetes进行编排管理。
推荐系统的构建是一个迭代和优化的长期过程,核心在于紧密围绕业务目标,平衡算法效果、系统性能和工程复杂度。从这个可运行的Demo出发,逐步引入更真实的数据、更复杂的算法和更健壮的架构,是掌握推荐系统实践的最佳路径。
