本节导读:掌握基于 Milvus 的智能推荐系统设计,实现协同过滤、内容推荐和混合推荐算法,打造精准的个性化推荐体验。
推荐系统是 Milvus 的核心应用场景之一:
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
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)
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
A:
A:
A:
通过本节的学习,你已经掌握了基于 Milvus 的智能推荐系统的核心技术,包括用户行为向量化、协同过滤、内容推荐和混合推荐算法。这些技术将帮助你构建精准的个性化推荐系统。
记住,在实际应用中需要根据具体的业务场景选择合适的技术方案,并不断优化推荐质量和用户体验。建立完善的监控和反馈机制,确保系统的持续改进。
关键词:Milvus, 推荐系统, 协同过滤, 内容推荐, 混合推荐, 个性化推荐
难度:高级
预计阅读:60 分钟