5.1 智能推荐系统


5.1 智能推荐系统 — Milvus 推荐引擎实战

本节导读:掌握基于 Milvus 的智能推荐系统设计,实现协同过滤、内容推荐和混合推荐算法,打造精准的个性化推荐体验。

学习目标

  • 掌握用户行为向量化技术
  • 实现协同过滤推荐算法
  • 构建内容推荐系统
  • 设计混合推荐策略

核心概念

1. 推荐系统架构

推荐系统是 Milvus 的核心应用场景之一:

推荐流程

  • 用户画像构建:基于用户行为构建用户向量
  • 物品特征提取:提取物品的语义特征向量
  • 相似度计算:计算用户与物品的相似度
  • 排序优化:结合多种因素进行最终排序

推荐策略

  • 协同过滤:基于用户行为相似性
  • 内容推荐:基于物品内容相似性
  • 混合推荐:结合多种推荐策略

分步实战

步骤 1:用户行为向量化

import torch import numpy as np from typing import List, Dict, Any from pymilvus import Collection, connections import json from datetime import datetime class UserBehaviorVectorizer: """用户行为向量化器""" def __init__(self, milvus_host="localhost", milvus_port="19530"): self.milvus_host = milvus_host self.milvus_port = milvus_port # 连接 Milvus connections.connect("default", host=milvus_host, port=milvus_port) # 初始化集合 self.create_user_collection() self.create_item_collection() def create_user_collection(self): """创建用户集合""" from pymilvus import CollectionSchema, FieldSchema, DataType fields = [ FieldSchema("user_id", DataType.INT64, is_primary=True), FieldSchema("user_vector", DataType.FLOAT_VECTOR, dim=128), FieldSchema("user_profile", DataType.JSON), FieldSchema("behavior_history", DataType.JSON), FieldSchema("last_updated", DataType.INT64) ] schema = CollectionSchema(fields, "user_collection") self.user_collection = Collection("user_collection", schema) # 创建索引 index_params = { "index_type": "HNSW", "params": {"M": 32, "ef": 256}, "metric_type": "IP" } self.user_collection.create_index("user_vector", index_params) def create_item_collection(self): """创建物品集合""" from pymilvus import CollectionSchema, FieldSchema, DataType fields = [ FieldSchema("item_id", DataType.INT64, is_primary=True), FieldSchema("item_vector", DataType.FLOAT_VECTOR, dim=128), FieldSchema("item_features", DataType.JSON), FieldSchema("category", DataType.VARCHAR, max_length=100), FieldSchema("popularity", DataType.FLOAT) ] schema = CollectionSchema(fields, "item_collection") self.item_collection = Collection("item_collection", schema) # 创建索引 index_params = { "index_type": "HNSW", "params": {"M": 32, "ef": 256}, "metric_type": "IP" } self.item_collection.create_index("item_vector", index_params) def extract_user_behavior_features(self, user_behaviors: List[Dict]) -> np.ndarray: """提取用户行为特征""" # 行为权重配置 behavior_weights = { "view": 0.1, "like": 0.3, "comment": 0.5, "share": 0.7, "purchase": 1.0 } # 初始化特征向量 feature_vector = np.zeros(128) # 统计各类行为 behavior_counts = {} for behavior in user_behaviors: behavior_type = behavior["type"] weight = behavior_weights.get(behavior_type, 0.1) if behavior_type not in behavior_counts: behavior_counts[behavior_type] = 0 behavior_counts[behavior_type] += weight # 构建特征向量 for i, (behavior_type, count) in enumerate(behavior_counts.items()): if i < 128: feature_vector[i] = count return feature_vector def update_user_vector(self, user_id: int, user_behaviors: List[Dict]): """更新用户向量""" # 提取用户行为特征 user_vector = self.extract_user_behavior_features(user_behaviors) # 构建用户画像 user_profile = { "user_id": user_id, "behavior_count": len(user_behaviors), "last_activity": datetime.now().isoformat(), "preferences": self.analyze_user_preferences(user_behaviors) } # 构建行为历史 behavior_history = { "behaviors": user_behaviors[-10:], # 保留最近10条行为 "updated_at": datetime.now().isoformat() } # 插入或更新用户数据 user_data = [ user_id, user_vector, user_profile, behavior_history, int(datetime.now().timestamp()) ] # 检查用户是否存在 existing_users = self.user_collection.query( expr=f"user_id == {user_id}", output_fields=["user_id"] ) if existing_users: # 更新用户数据 self.user_collection.delete(expr=f"user_id == {user_id}") self.user_collection.insert([user_data]) else: # 插入新用户数据 self.user_collection.insert([user_data]) def analyze_user_preferences(self, user_behaviors: List[Dict]) -> Dict[str, float]: """分析用户偏好""" category_scores = {} for behavior in user_behaviors: category = behavior.get("category", "unknown") behavior_type = behavior.get("type", "view") # 行为权重 weights = { "view": 1.0, "like": 2.0, "comment": 3.0, "share": 4.0, "purchase": 5.0 } weight = weights.get(behavior_type, 1.0) if category not in category_scores: category_scores[category] = 0 category_scores[category] += weight # 归一化权重 total_score = sum(category_scores.values()) if total_score > 0: category_scores = { category: score / total_score for category, score in category_scores.items() } return category_scores

步骤 2:推荐算法实现

class RecommendationEngine: """推荐引擎""" def __init__(self, milvus_host="localhost", milvus_port="19530"): self.milvus_host = milvus_host self.milvus_port = milvus_port # 初始化向量化器 self.vectorizer = UserBehaviorVectorizer(milvus_host, milvus_port) # 推荐策略配置 self.strategies = { "collaborative_filtering": self.collaborative_filtering, "content_based": self.content_based_recommendation, "hybrid": self.hybrid_recommendation } def collaborative_filtering(self, user_id: int, top_k: int = 10) -> List[Dict]: """协同过滤推荐""" # 获取用户向量 user_results = self.vectorizer.user_collection.query( expr=f"user_id == {user_id}", output_fields=["user_vector"] ) if not user_results: return [] user_vector = user_results[0].get("user_vector") # 搜索相似用户 similar_users = self.vectorizer.user_collection.search( data=[user_vector], anns_field="user_vector", param={"metric_type": "IP", "ef": 50}, limit=10, output_fields=["user_id", "user_profile"] ) # 获取相似用户喜欢的物品 recommended_items = [] for hit in similar_users[0]: similar_user_id = hit.id # 获取相似用户的行为 user_behaviors = self.get_user_behaviors(similar_user_id) # 添加推荐物品 for behavior in user_behaviors: if behavior["type"] in ["like", "comment", "share", "purchase"]: recommended_items.append({ "item_id": behavior["item_id"], "score": hit.score, "reason": f"与用户 {similar_user_id} 行为相似" }) # 去重和排序 unique_items = {} for item in recommended_items: item_id = item["item_id"] if item_id not in unique_items or item["score"] > unique_items[item_id]["score"]: unique_items[item_id] = item # 按分数排序 sorted_items = sorted(unique_items.values(), key=lambda x: x["score"], reverse=True) return sorted_items[:top_k] def content_based_recommendation(self, user_id: int, top_k: int = 10) -> List[Dict]: """基于内容的推荐""" # 获取用户偏好 user_results = self.vectorizer.user_collection.query( expr=f"user_id == {user_id}", output_fields=["user_profile"] ) if not user_results: return [] user_profile = user_results[0].get("user_profile", {}) preferences = user_profile.get("preferences", {}) # 获取用户向量 user_results = self.vectorizer.user_collection.query( expr=f"user_id == {user_id}", output_fields=["user_vector"] ) if not user_results: return [] user_vector = user_results[0].get("user_vector") # 搜索相似物品 similar_items = self.vectorizer.item_collection.search( data=[user_vector], anns_field="item_vector", param={"metric_type": "IP", "ef": 50}, limit=top_k, output_fields=["item_id", "category", "popularity"] ) # 处理推荐结果 recommended_items = [] for hit in similar_items[0]: item_id = hit.id # 根据类别权重调整分数 item_results = self.vectorizer.item_collection.query( expr=f"item_id == {item_id}", output_fields=["category"] ) if item_results: item_category = item_results[0].get("category") category_weight = preferences.get(item_category, 0.5) adjusted_score = hit.score * category_weight recommended_items.append({ "item_id": item_id, "score": adjusted_score, "original_score": hit.score, "category": item_category, "reason": f"类别 {item_category} 匹配用户偏好" }) # 按调整后的分数排序 recommended_items.sort(key=lambda x: x["score"], reverse=True) return recommended_items[:top_k] def hybrid_recommendation(self, user_id: int, top_k: int = 10) -> List[Dict]: """混合推荐""" # 获取不同策略的推荐结果 cf_results = self.collaborative_filtering(user_id, top_k * 2) cb_results = self.content_based_recommendation(user_id, top_k * 2) # 合并推荐结果 combined_results = {} # 协同过滤结果 for item in cf_results: item_id = item["item_id"] combined_results[item_id] = { "item_id": item_id, "cf_score": item["score"], "cb_score": 0, "reason": item["reason"] } # 基于内容的结果 for item in cb_results: item_id = item["item_id"] if item_id in combined_results: combined_results[item_id]["cb_score"] = item["score"] combined_results[item_id]["category"] = item.get("category") else: combined_results[item_id] = { "item_id": item_id, "cf_score": 0, "cb_score": item["score"], "reason": item["reason"] } # 计算混合分数 hybrid_results = [] for item_id, scores in combined_results.items(): cf_score = scores.get("cf_score", 0) cb_score = scores.get("cb_score", 0) # 混合权重 hybrid_score = 0.6 * cf_score + 0.4 * cb_score hybrid_results.append({ "item_id": item_id, "hybrid_score": hybrid_score, "cf_score": cf_score, "cb_score": cb_score, "reason": scores.get("reason", "混合推荐") }) # 按混合分数排序 hybrid_results.sort(key=lambda x: x["hybrid_score"], reverse=True) return hybrid_results[:top_k] def generate_recommendations(self, user_id: int, strategy="hybrid", top_k=10) -> List[Dict]: """生成推荐""" if strategy not in self.strategies: strategy = "hybrid" recommendation_func = self.strategies[strategy] return recommendation_func(user_id, top_k)

步骤 3:推荐系统优化

class RecommendationOptimizer: """推荐系统优化器""" def __init__(self, engine: RecommendationEngine): self.engine = engine def optimize_recommendation_quality(self, user_id: int, top_k: int = 10) -> List[Dict]: """优化推荐质量""" # 获取基础推荐 recommendations = self.engine.generate_recommendations(user_id, "hybrid", top_k) # 添加多样性优化 diversified_recommendations = self.add_diversity(recommendations) # 添加新颖性优化 final_recommendations = self.add_novelty(diversified_recommendations) return final_recommendations def add_diversity(self, recommendations: List[Dict]) -> List[Dict]: """添加多样性""" category_diversity = {} diversified = [] for rec in recommendations: category = rec.get("category", "unknown") if category not in category_diversity or len(category_diversity[category]) < 3: if category not in category_diversity: category_diversity[category] = [] category_diversity[category].append(rec) diversified.append(rec) return diversified def add_novelty(self, recommendations: List[Dict]) -> List[Dict]: """添加新颖性""" # 为推荐添加新颖性分数 for rec in recommendations: # 根据物品的新颖度调整分数 novelty_factor = self.calculate_novelty(rec["item_id"]) rec["novelty_score"] = novelty_factor rec["final_score"] = rec.get("hybrid_score", 0) * (1 + novelty_factor * 0.1) # 按最终分数排序 recommendations.sort(key=lambda x: x["final_score"], reverse=True) return recommendations def calculate_novelty(self, item_id: int) -> float: """计算物品新颖度""" # 这里可以实现复杂的新颖度计算逻辑 # 简化版:基于物品的流行度 item_results = self.engine.vectorizer.item_collection.query( expr=f"item_id == {item_id}", output_fields=["popularity"] ) if item_results: popularity = item_results[0].get("popularity", 0.5) # 流行度越低,新颖度越高 return 1.0 - popularity return 0.5

常见问题 FAQ

Q1:如何处理新用户的冷启动问题?

A

  • 基于内容的推荐:为新用户提供基于内容的推荐
  • 热门推荐:推荐当前热门的物品
  • 随机推荐:随机推荐一些物品收集用户反馈
  • 引导式推荐:通过问卷收集用户偏好

Q2:如何优化推荐系统的实时性?

A

  • 缓存热门结果:缓存热门用户的推荐结果
  • 增量更新:只更新变化的部分数据
  • 预计算:提前计算可能的推荐结果
  • 负载均衡:部署多个推荐服务实例

Q3:如何提高推荐的准确性?

A

  • 特征工程:提取更丰富的用户和物品特征
  • 模型调优:使用更复杂的推荐模型
  • A/B 测试:通过 A/B 测试验证推荐效果
  • 反馈机制:收集用户反馈并优化推荐策略

最佳实践与避坑

1. 推荐系统最佳实践

  • 数据质量:确保用户行为数据的质量和完整性
  • 实时性:平衡实时性和计算复杂度
  • 多样性:保证推荐结果的多样性
  • 冷启动:为新用户提供合适的冷启动策略

2. 实时推荐优化

  • 缓存策略:使用多级缓存减少计算量
  • 增量更新:实时更新用户向量
  • 负载均衡:合理分配计算资源
  • 监控告警:实时监控推荐效果

本节小结

通过本节的学习,你已经掌握了基于 Milvus 的智能推荐系统的核心技术,包括用户行为向量化、协同过滤、内容推荐和混合推荐算法。这些技术将帮助你构建精准的个性化推荐系统。

记住,在实际应用中需要根据具体的业务场景选择合适的技术方案,并不断优化推荐质量和用户体验。建立完善的监控和反馈机制,确保系统的持续改进。

延伸阅读

  • 推荐系统算法与实现
  • 机器学习在推荐中的应用
  • 大规模推荐系统架构
  • 个性化推荐技术实践

关键词:Milvus, 推荐系统, 协同过滤, 内容推荐, 混合推荐, 个性化推荐
难度:高级
预计阅读:60 分钟


作者与出处
整理: 灏天文库整理
本站整理收录,版权归原作者/开源协议所有;欢迎通过原文链接访问源仓库。
发布者: 作者: 来自仙女座的脉冲的小龙虾 转发
评论区 (0)
U