当前位置:网站首页 >  教程

智能内容推送系统实战:从零搭建可扩展推荐引擎

时间:2026年06月14日 08:05:56 来源:易频IT社区

系统架构设计

我们将采用模块化架构,整个系统分为四个核心模块:数据采集层、特征处理层、推荐算法层和API服务层。每个模块独立部署,通过消息队列通信。

技术栈选型

  • 数据处理: Python 3.9 + Pandas 2.0 + Scikit-learn 1.3
  • 消息队列: RabbitMQ 3.12
  • 数据存储: PostgreSQL 15(用户画像)+ Redis 7.2(实时特征)
  • API服务: FastAPI 0.104

环境搭建与配置

开发环境准备

创建项目目录并安装依赖:

``` mkdir smart-content-push && cd smart-content-push python -m venv venv source venv/bin/activate Linux/Mac venv\Scripts\activate Windows pip install pandas==2.0.3 scikit-learn==1.3.0 fastapi==0.104.0 redis==5.0.1 ```

数据库初始化

创建PostgreSQL数据库和表结构:

``` CREATE DATABASE content_recommend; \c content_recommend; CREATE TABLE user_profiles ( user_id VARCHAR(50) PRIMARY KEY, age_group INT, interests TEXT[], last_active TIMESTAMP, click_history JSONB ); CREATE TABLE content_items ( content_id VARCHAR(50) PRIMARY KEY, title TEXT, category VARCHAR(50), tags TEXT[], publish_date DATE, feature_vector FLOAT[] ); ```

数据采集模块实现

用户行为日志收集

创建data_collector.py文件:

``` import json import time from datetime import datetime import redis class UserBehaviorCollector: def __init__(self): self.redis_client = redis.Redis(host='localhost', port=6379, db=0) def log_click(self, user_id, content_id, duration): """记录用户点击行为""" event = { 'event_type': 'click', 'user_id': user_id, 'content_id': content_id, 'duration': duration, 'timestamp': datetime.now().isoformat() } self.redis_client.lpush('user_events', json.dumps(event)) self.redis_client.expire('user_events', 86400) 保留24小时 def get_recent_events(self, limit=100): """获取最近的行为事件""" events = self.redis_client.lrange('user_events', 0, limit-1) return [json.loads(e) for e in events] ```

特征工程处理

用户特征提取

创建feature_processor.py文件:

``` import numpy as np from collections import Counter import psycopg2 class FeatureProcessor: def __init__(self, db_config): self.conn = psycopg2.connect(db_config) def extract_user_features(self, user_id): """提取用户特征向量""" cursor = self.conn.cursor() 获取用户基础信息 cursor.execute(""" SELECT age_group, interests, click_history FROM user_profiles WHERE user_id = %s """, (user_id,)) result = cursor.fetchone() if not result: return None age_group, interests, click_history = result 计算兴趣标签权重 tag_weights = {} if click_history: for item in click_history: tags = item.get('tags', []) for tag in tags: tag_weights[tag] = tag_weights.get(tag, 0) + 1 生成特征向量 features = { 'age_group': age_group or 0, 'interest_count': len(interests) if interests else 0, 'top_tags': dict(Counter(tag_weights).most_common(5)) } return features ```

推荐算法实现

协同过滤算法

创建recommender.py文件:

``` import numpy as np from sklearn.metrics.pairwise import cosine_similarity import psycopg2 class ContentRecommender: def __init__(self, db_config): self.conn = psycopg2.connect(db_config) def calculate_similarity(self, user_features, content_features): """计算用户与内容的相似度""" 将特征转换为向量 user_vector = self.features_to_vector(user_features) content_vector = np.array(content_features) 计算余弦相似度 similarity = cosine_similarity([user_vector], [content_vector])[0][0] return similarity def features_to_vector(self, features): """特征向量化""" vector = [] vector.append(features.get('age_group', 0)) vector.append(features.get('interest_count', 0)) 标签权重部分 top_tags = features.get('top_tags', {}) tag_score = sum(top_tags.values()) / 10 if top_tags else 0 vector.append(tag_score) return np.array(vector) def get_recommendations(self, user_id, limit=10): """获取推荐内容""" cursor = self.conn.cursor() 获取用户特征 cursor.execute(""" SELECT up.age_group, up.interests, up.click_history, ci.feature_vector, ci.content_id, ci.title FROM user_profiles up CROSS JOIN content_items ci WHERE up.user_id = %s """, (user_id,)) results = cursor.fetchall() recommendations = [] for row in results: age_group, interests, click_history, feature_vector, content_id, title = row user_features = { 'age_group': age_group, 'interest_count': len(interests) if interests else 0, 'top_tags': self.extract_top_tags(click_history) } similarity = self.calculate_similarity(user_features, feature_vector) recommendations.append({ 'content_id': content_id, 'title': title, 'score': similarity }) 按分数排序并返回前N个 recommendations.sort(key=lambda x: x['score'], reverse=True) return recommendations[:limit] ```

API服务部署

FastAPI应用搭建

创建main.py文件:

``` from fastapi import FastAPI, HTTPException from pydantic import BaseModel from typing import List import uvicorn from recommender import ContentRecommender from feature_processor import FeatureProcessor app = FastAPI(title="智能内容推送API") 数据库配置 DB_CONFIG = { 'host': 'localhost', 'database': 'content_recommend', 'user': 'postgres', 'password': 'your_password' } 初始化组件 recommender = ContentRecommender(DB_CONFIG) feature_processor = FeatureProcessor(DB_CONFIG) class RecommendationRequest(BaseModel): user_id: str limit: int = 10 class RecommendationResponse(BaseModel): content_id: str title: str score: float @app.post("/recommend", response_model=List[RecommendationResponse]) async def get_recommendations(request: RecommendationRequest): """获取个性化内容推荐""" try: recommendations = recommender.get_recommendations( request.user_id, request.limit ) return recommendations except Exception as e: raise HTTPException(status_code=500, detail=str(e)) @app.get("/user/{user_id}/features") async def get_user_features(user_id: str): """获取用户特征""" features = feature_processor.extract_user_features(user_id) if not features: raise HTTPException(status_code=404, detail="User not found") return features if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8000) ```

API测试

智能内容推送系统实战:从零搭建可扩展推荐引擎

启动服务后,使用curl测试:

``` 启动服务 python main.py 另一个终端执行测试 curl -X POST "http://localhost:8000/recommend" \ -H "Content-Type: application/json" \ -d '{"user_id": "user123", "limit": 5}' ```

系统部署与监控

Docker容器化部署

创建Dockerfile

``` FROM python:3.9-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt COPY . . EXPOSE 8000 CMD ["python", "main.py"] ```

创建docker-compose.yml

``` version: '3.8' services: postgres: image: postgres:15 environment: POSTGRES_PASSWORD: your_password POSTGRES_DB: content_recommend volumes: - postgres_data:/var/lib/postgresql/data ports: - "5432:5432" redis: image: redis:7.2-alpine ports: - "6379:6379" recommender: build: . ports: - "8000:8000" depends_on: - postgres - redis environment: DB_HOST: postgres REDIS_HOST: redis volumes: postgres_data: ```

启动完整系统

执行以下命令启动所有服务:

``` docker-compose up -d ```

性能优化建议

缓存策略实现

recommender.py中添加缓存:

``` import redis from functools import lru_cache class CachedRecommender(ContentRecommender): def __init__(self, db_config): super().__init__(db_config) self.redis_client = redis.Redis(host='localhost', port=6379, db=1) @lru_cache(maxsize=1000) def get_recommendations(self, user_id, limit=10): cache_key = f"recommendations:{user_id}:{limit}" 检查缓存 cached_result = self.redis_client.get(cache_key) if cached_result: return json.loads(cached_result) 计算推荐结果 result = super().get_recommendations(user_id, limit) 缓存结果(5分钟过期) self.redis_client.setex(cache_key, 300, json.dumps(result)) return result ```

批量处理优化

对于大量用户推荐,使用批量查询:

``` def batch_recommendations(self, user_ids, limit=10): """批量获取推荐""" placeholders = ','.join(['%s'] len(user_ids)) query = f""" SELECT user_id, array_agg(content_id) as recommendations FROM ( SELECT up.user_id, ci.content_id, ROW_NUMBER() OVER (PARTITION BY up.user_id ORDER BY similarity DESC) as rank FROM user_profiles up CROSS JOIN content_items ci WHERE up.user_id IN ({placeholders}) ) ranked WHERE rank <= %s GROUP BY user_id """ cursor = self.conn.cursor() cursor.execute(query, user_ids + [limit]) return dict(cursor.fetchall()) ```

故障排查指南

常见问题解决

  • 数据库连接失败:检查PostgreSQL服务状态sudo systemctl status postgresql
  • Redis连接超时:确认Redis服务运行redis-cli ping应返回PONG
  • 推荐结果为空:检查content_items表是否有数据,特征向量是否已计算
  • API响应慢:启用查询缓存,检查数据库索引是否创建

监控指标设置

main.py中添加健康检查:

``` @app.get("/health") async def health_check(): """系统健康检查""" checks = { 'database': check_database_connection(), 'redis': check_redis_connection(), 'recommender': True } status = all(checks.values()) return { 'status': 'healthy' if status else 'unhealthy', 'checks': checks } ```

至此,完整的智能内容推送系统已搭建完成。系统包含数据采集、特征处理、推荐算法和API服务全链路,可直接部署到生产环境。后续可根据业务需求调整推荐算法或增加更多特征维度。

相关推荐

最新

热门

推荐

精选

标签

易频IT社区是综合性互联网IT技术门户网站,专注分享网络技术、服务器运维、网络安全、编程开发、系统架构、云计算、大数据等行业干货,实时更新IT行业资讯、零基础教程、实战案例,为IT从业者、技术爱好者提供专业的学习交流平台。

Copyright © 2021-2026 易频IT社区. All Rights Reserved. 备案号:闽ICP备2023013482号 网站地图