引言:电影市场的复杂性与预测的必要性
电影市场是一个充满变数的复杂生态系统,受到多种因素的综合影响。传统的票房预测方法往往依赖于历史数据和简单的线性模型,但这些方法在面对快速变化的市场环境时显得力不从心。随着大数据和人工智能技术的发展,我们有机会更精准地把握电影市场的脉搏和观众的真实喜好趋势。
本文将深入探讨如何利用现代数据科学方法,超越传统的猫眼票房预测模型,构建更精准的票房预测系统。我们将从数据收集、特征工程、模型构建等多个维度进行详细分析,并提供完整的代码示例,帮助读者理解并实践这些方法。
一、电影票房预测的核心挑战
1.1 多维度影响因素
电影票房受到众多因素的影响,包括但不限于:
- 影片本身因素:导演、演员阵容、制作成本、类型、IP影响力
- 市场环境因素:档期选择、竞争对手、市场饱和度
- 营销推广因素:预告片传播量、社交媒体热度、媒体报道
- 观众反馈因素:点映口碑、评分、评论情感分析
- 外部环境因素:经济环境、政策法规、突发事件
1.2 数据的时效性与稀疏性
电影市场数据具有很强的时效性,需要实时更新。同时,某些类型电影的数据可能相对稀疏,给模型训练带来挑战。
1.3 非线性关系
票房与各影响因素之间往往存在复杂的非线性关系,简单的线性模型难以捕捉这些关系。
二、数据收集与预处理
2.1 数据来源
构建精准的票房预测模型需要多源数据:
import pandas as pd
import numpy as np
import requests
from bs4 import BeautifulSoup
import json
from datetime import datetime, timedelta
class MovieDataCollector:
def __init__(self):
self.sources = {
'box_office': '票房数据',
'movie_info': '影片信息',
'social_media': '社交媒体数据',
'reviews': '评论数据'
}
def collect_box_office_data(self, start_date, end_date):
"""
收集每日票房数据
"""
# 模拟API调用(实际使用时需要真实的API接口)
dates = pd.date_range(start=start_date, end=end_date)
data = []
for date in dates:
# 这里应该是真实的API调用
# response = requests.get(f'https://api.boxoffice.com/daily?date={date}')
# 实际数据结构会更复杂
mock_data = {
'date': date,
'total_box_office': np.random.randint(10000000, 50000000),
'movies': [
{'title': f'Movie_{i}', 'box_office': np.random.randint(1000000, 10000000)}
for i in range(5)
]
}
data.append(mock_data)
return pd.DataFrame(data)
def collect_movie_info(self, movie_id):
"""
收集影片详细信息
"""
# 模拟从数据库或API获取影片信息
movie_info = {
'movie_id': movie_id,
'title': f'电影_{movie_id}',
'director': '知名导演',
'cast': ['演员A', '演员B', '演员C'],
'genre': ['动作', '科幻'],
'release_date': '2024-01-01',
'budget': 50000000,
'production_company': 'XX影业',
'ip_power': 8.5 # IP影响力评分
}
return movie_info
# 使用示例
collector = MovieDataCollector()
box_office_data = collector.collect_box_office_data('2024-01-01', '2024-01-31')
print("收集到的票房数据:")
print(box_office_data.head())
2.2 数据清洗与标准化
class DataPreprocessor:
def __init__(self):
self.scalers = {}
def clean_box_office_data(self, df):
"""
清洗票房数据,处理异常值和缺失值
"""
# 处理缺失值
df = df.dropna(subset=['total_box_office'])
# 处理异常值(使用IQR方法)
Q1 = df['total_box_office'].quantile(0.25)
Q3 = df['total_box_office'].quantile(0.75)
IQR = Q3 - Q1
lower_bound = Q1 - 1.5 * IQR
upper_bound = Q3 + 1.5 * IQR
df = df[(df['total_box_office'] >= lower_bound) &
(df['total_box_office'] <= upper_bound)]
return df
def normalize_features(self, df, feature_columns):
"""
标准化特征数据
"""
from sklearn.preprocessing import StandardScaler
scaler = StandardScaler()
df[feature_columns] = scaler.fit_transform(df[feature_columns])
self.scalers['features'] = scaler
return df
def encode_categorical_features(self, df, categorical_columns):
"""
编码分类特征
"""
df_encoded = pd.get_dummies(df, columns=categorical_columns, drop_first=True)
return df_encoded
# 使用示例
preprocessor = DataPreprocessor()
cleaned_data = preprocessor.clean_box_office_data(box_office_data)
print("清洗后的数据:")
print(cleaned_data.head())
2.3 特征工程:构建预测模型的基础
特征工程是票房预测的关键步骤。我们需要从原始数据中提取有意义的特征:
class FeatureEngineer:
def __init__(self):
self.feature_list = []
def create_time_features(self, df, date_column='date'):
"""
创建时间相关特征
"""
df[date_column] = pd.to_datetime(df[date_column])
# 提取时间特征
df['year'] = df[date_column].dt.year
df['month'] = df[date_column].dt.month
df['day'] = df[date_column].dt.day
df['day_of_week'] = df[date_column].dt.dayofweek
df['is_weekend'] = df['day_of_week'].isin([5, 6]).astype(int)
# 节假日特征(简化版)
df['is_holiday'] = df['day'].isin([1, 2, 3, 4, 5, 6, 7, 15, 16, 17, 18, 19, 20, 21]).astype(int)
self.feature_list.extend(['year', 'month', 'day', 'day_of_week', 'is_weekend', 'is_holiday'])
return df
def create_movie_features(self, movie_info):
"""
创建影片相关特征
"""
features = {}
# 演员影响力(简化计算)
cast_power = len(movie_info['cast']) * 0.3 # 演员数量
# 导演影响力(简化计算)
director_power = 1.0 # 可以根据历史数据调整
# 类型特征
genre_features = {f'genre_{g}': 1 for g in movie_info['genre']}
# IP影响力
ip_power = movie_info.get('ip_power', 0)
# 预算特征(标准化)
budget = movie_info['budget']
budget_normalized = np.log1p(budget) # 对数变换
features.update({
'cast_power': cast_power,
'director_power': director_power,
'ip_power': ip_power,
'budget_log': budget_normalized,
'production_company_encoded': hash(movie_info['production_company']) % 1000
})
features.update(genre_features)
return features
def create_social_media_features(self, social_data):
"""
创建社交媒体特征
"""
features = {}
# 提取关键指标
features['weibo_mentions'] = social_data.get('weibo_mentions', 0)
features['douyin_views'] = social_data.get('douyin_views', 0)
features['xiaohongshu_notes'] = social_data.get('xiaohongshu_notes', 0)
features['total_engagement'] = (
features['weibo_mentions'] * 0.3 +
features['douyin_views'] * 0.0001 +
features['xiaohongshu_notes'] * 0.5
)
# 增长率特征
features['mention_growth_rate'] = social_data.get('mention_growth_rate', 0)
self.feature_list.extend(['weibo_mentions', 'douyin_views', 'xiaohongshu_notes',
'total_engagement', 'mention_growth_rate'])
return features
def create_review_features(self, review_data):
"""
创建评论相关特征
"""
features = {}
# 评分特征
features['avg_rating'] = review_data.get('avg_rating', 0)
features['rating_count'] = review_data.get('rating_count', 0)
# 情感分析(简化版)
positive_ratio = review_data.get('positive_ratio', 0.5)
features['sentiment_score'] = positive_ratio * 2 - 1 # 转换到[-1, 1]区间
# 评论关键词提取(简化)
keywords = review_data.get('keywords', [])
features['keyword_diversity'] = len(set(keywords))
self.feature_list.extend(['avg_rating', 'rating_count', 'sentiment_score', 'keyword_diversity'])
return features
# 使用示例
feature_engineer = FeatureEngineer()
# 创建时间特征
df_with_time = feature_engineer.create_time_features(cleaned_data)
# 模拟影片特征
movie_info = {
'cast': ['演员A', '演员B'],
'genre': ['动作', '科幻'],
'budget': 50000000,
'ip_power': 8.5,
'production_company': 'XX影业'
}
movie_features = feature_engineer.create_movie_features(movie_info)
# 模拟社交媒体数据
social_data = {
'weibo_mentions': 15000,
'douyin_views': 5000000,
'xiaohongshu_notes': 2000,
'mention_growth_rate': 0.3
}
social_features = feature_engineer.create_social_media_features(social_data)
# 模拟评论数据
review_data = {
'avg_rating': 8.2,
'rating_count': 50000,
'positive_ratio': 0.75,
'keywords': ['特效', '剧情', '演员', '节奏']
}
review_features = feature_engineer.create_review_features(review_data)
print("提取的特征列表:", feature_engineer.feature_list)
print("影片特征:", movie_features)
print("社交媒体特征:", social_features)
print("评论特征:", review_features)
三、构建预测模型
3.1 基础模型选择
对于票房预测,我们通常需要尝试多种模型:
from sklearn.model_selection import train_test_split, cross_val_score
from sklearn.linear_model import LinearRegression, Ridge, Lasso
from sklearn.ensemble import RandomForestRegressor, GradientBoostingRegressor
from sklearn.svm import SVR
from sklearn.metrics import mean_absolute_error, mean_squared_error, r2_score
import xgboost as xgb
import lightgbm as lgb
class BoxOfficePredictor:
def __init__(self):
self.models = {
'linear': LinearRegression(),
'ridge': Ridge(alpha=1.0),
'lasso': Lasso(alpha=0.1),
'random_forest': RandomForestRegressor(n_estimators=100, random_state=42),
'gradient_boosting': GradientBoostingRegressor(n_estimators=100, random_state=42),
'xgboost': xgb.XGBRegressor(n_estimators=100, random_state=42),
'lightgbm': lgb.LGBMRegressor(n_estimators=100, random_state=42)
}
self.best_model = None
self.best_score = float('-inf')
def prepare_training_data(self, df, feature_columns, target_column):
"""
准备训练数据
"""
X = df[feature_columns]
y = df[target_column]
# 划分训练集和测试集
X_train, X_test, y_train, y_test = train_test_split(
X, y, test_size=0.2, random_state=42
)
return X_train, X_test, y_train, y_test
def train_and_evaluate_models(self, X_train, X_test, y_train, y_test):
"""
训练并评估多个模型
"""
results = {}
for name, model in self.models.items():
print(f"训练模型: {name}")
# 训练模型
model.fit(X_train, y_train)
# 预测
y_pred = model.predict(X_test)
# 计算指标
mae = mean_absolute_error(y_test, y_pred)
mse = mean_squared_error(y_test, y_pred)
r2 = r2_score(y_test, y_pred)
# 交叉验证
cv_scores = cross_val_score(model, X_train, y_train, cv=5, scoring='r2')
results[name] = {
'model': model,
'mae': mae,
'mse': mse,
'r2': r2,
'cv_mean': cv_scores.mean(),
'cv_std': cv_scores.std()
}
print(f" MAE: {mae:.2f}, MSE: {mse:.2f}, R2: {r2:.4f}")
print(f" CV Score: {cv_scores.mean():.4f} (+/- {cv_scores.std() * 2:.4f})")
# 更新最佳模型
if cv_scores.mean() > self.best_score:
self.best_score = cv_scores.mean()
self.best_model = model
return results
def predict_with_confidence(self, model, X, confidence=0.95):
"""
带置信区间的预测
"""
from scipy import stats
# 获取预测值
predictions = model.predict(X)
# 计算残差标准差(简化)
# 实际应用中应该使用验证集残差
residual_std = np.std(predictions) * 0.1 # 简化假设
# 计算置信区间
z_score = stats.norm.ppf((1 + confidence) / 2)
margin_of_error = z_score * residual_std
lower_bound = predictions - margin_of_error
upper_bound = predictions + margin_of_error
return predictions, lower_bound, upper_bound
# 使用示例
predictor = BoxOfficePredictor()
# 准备训练数据(模拟)
# 这里需要实际的数据集,我们创建一个模拟数据集
np.random.seed(42)
n_samples = 1000
# 创建模拟特征数据
simulated_data = pd.DataFrame({
'total_box_office': np.random.lognormal(17, 0.5, n_samples), # 票房
'cast_power': np.random.uniform(0, 10, n_samples),
'director_power': np.random.uniform(0, 10, n_samples),
'ip_power': np.random.uniform(0, 10, n_samples),
'budget_log': np.random.uniform(15, 20, n_samples),
'total_engagement': np.random.uniform(0, 100000, n_samples),
'avg_rating': np.random.uniform(6, 9.5, n_samples),
'sentiment_score': np.random.uniform(-1, 1, n_samples),
'is_weekend': np.random.randint(0, 2, n_samples),
'is_holiday': np.random.randint(0, 2, n_samples)
})
# 添加一些相关性
simulated_data['total_box_office'] += (
simulated_data['cast_power'] * 500000 +
simulated_data['ip_power'] * 800000 +
simulated_data['total_engagement'] * 10 +
simulated_data['avg_rating'] * 200000
)
feature_columns = ['cast_power', 'director_power', 'ip_power', 'budget_log',
'total_engagement', 'avg_rating', 'sentiment_score',
'is_weekend', 'is_holiday']
X_train, X_test, y_train, y_test = predictor.prepare_training_data(
simulated_data, feature_columns, 'total_box_office'
)
# 训练和评估
results = predictor.train_and_evaluate_models(X_train, X_test, y_train, y_test)
print(f"\n最佳模型: {type(predictor.best_model).__name__}")
print(f"最佳CV分数: {predictor.best_score:.4f}")
3.2 高级模型:集成学习与时间序列分析
对于票房预测,时间序列特征非常重要。我们可以结合集成学习和时间序列分析:
class AdvancedBoxOfficePredictor:
def __init__(self):
self.models = {}
self.feature_importance = {}
def create_time_series_features(self, df, date_col='date', target_col='total_box_office'):
"""
创建时间序列特征
"""
df = df.copy()
df[date_col] = pd.to_datetime(df[date_col])
df = df.sort_values(date_col)
# 滞后特征
for lag in [1, 7, 14, 30]:
df[f'{target_col}_lag_{lag}'] = df[target_col].shift(lag)
# 滚动统计特征
for window in [7, 14, 30]:
df[f'{target_col}_rolling_mean_{window}'] = df[target_col].rolling(window=window).mean()
df[f'{target_col}_rolling_std_{window}'] = df[target_col].rolling(window=window).std()
# 扩展窗口特征
df[f'{target_col}_expanding_mean'] = df[target_col].expanding().mean()
# 增长率特征
df[f'{target_col}_growth_rate'] = df[target_col].pct_change()
# 填充NaN值
df = df.fillna(0)
return df
def build_ensemble_model(self, X_train, y_train):
"""
构建集成模型
"""
# XGBoost
xgb_model = xgb.XGBRegressor(
n_estimators=200,
max_depth=6,
learning_rate=0.1,
subsample=0.8,
colsample_bytree=0.8,
random_state=42
)
# LightGBM
lgb_model = lgb.LGBMRegressor(
n_estimators=200,
max_depth=6,
learning_rate=0.1,
subsample=0.8,
colsample_bytree=0.8,
random_state=42
)
# 随机森林
rf_model = RandomForestRegressor(
n_estimators=200,
max_depth=10,
random_state=42
)
# 训练基础模型
print("训练XGBoost...")
xgb_model.fit(X_train, y_train)
print("训练LightGBM...")
lgb_model.fit(X_train, y_train)
print("训练随机森林...")
rf_model.fit(X_train, y_train)
self.models = {
'xgboost': xgb_model,
'lightgbm': lgb_model,
'random_forest': rf_model
}
# 计算特征重要性
self._calculate_feature_importance(X_train.columns)
return self.models
def _calculate_feature_importance(self, feature_names):
"""
计算特征重要性
"""
for name, model in self.models.items():
if hasattr(model, 'feature_importances_'):
importance = model.feature_importances_
self.feature_importance[name] = dict(zip(feature_names, importance))
def ensemble_predict(self, X, method='weighted'):
"""
集成预测
"""
predictions = {}
for name, model in self.models.items():
predictions[name] = model.predict(X)
if method == 'weighted':
# 根据模型性能加权
weights = {'xgboost': 0.4, 'lightgbm': 0.35, 'random_forest': 0.25}
final_pred = sum(predictions[name] * weights[name] for name in predictions)
elif method == 'average':
# 简单平均
final_pred = np.mean(list(predictions.values()), axis=0)
elif method == 'median':
# 中位数
final_pred = np.median(np.column_stack(list(predictions.values())), axis=1)
else:
raise ValueError("Method must be 'weighted', 'average', or 'median'")
return final_pred, predictions
def get_feature_importance_df(self, top_n=10):
"""
获取特征重要性DataFrame
"""
importance_dfs = []
for model_name, importance_dict in self.feature_importance.items():
df = pd.DataFrame(list(importance_dict.items()), columns=['feature', 'importance'])
df['model'] = model_name
importance_dfs.append(df)
combined_importance = pd.concat(importance_dfs)
combined_importance = combined_importance.groupby('feature')['importance'].mean().reset_index()
combined_importance = combined_importance.sort_values('importance', ascending=False).head(top_n)
return combined_importance
# 使用示例
advanced_predictor = AdvancedBoxOfficePredictor()
# 添加时间序列特征
df_with_ts = advanced_predictor.create_time_series_features(simulated_data)
# 准备数据
feature_columns_ts = [col for col in df_with_ts.columns if col not in ['total_box_office', 'date']]
X_train_ts, X_test_ts, y_train_ts, y_test_ts = predictor.prepare_training_data(
df_with_ts, feature_columns_ts, 'total_box_office'
)
# 训练集成模型
models = advanced_predictor.build_ensemble_model(X_train_ts, y_train_ts)
# 进行预测
ensemble_pred, individual_preds = advanced_predictor.ensemble_predict(X_test_ts, method='weighted')
# 评估
print("\n集成模型评估:")
print(f"MAE: {mean_absolute_error(y_test_ts, ensemble_pred):.2f}")
print(f"R2: {r2_score(y_test_ts, ensemble_pred):.4f}")
# 显示特征重要性
importance_df = advanced_predictor.get_feature_importance_df()
print("\n特征重要性(前10名):")
print(importance_df)
3.3 深度学习模型:LSTM与注意力机制
对于时间序列票房预测,深度学习模型能够捕捉更复杂的模式:
import torch
import torch.nn as nn
import torch.optim as optim
from torch.utils.data import Dataset, DataLoader
class MovieDataset(Dataset):
def __init__(self, X, y, sequence_length=30):
self.X = torch.FloatTensor(X.values) if isinstance(X, pd.DataFrame) else torch.FloatTensor(X)
self.y = torch.FloatTensor(y.values) if isinstance(y, pd.Series) else torch.FloatTensor(y)
self.seq_len = sequence_length
def __len__(self):
return len(self.X) - self.seq_len + 1
def __getitem__(self, idx):
return self.X[idx:idx+self.seq_len], self.y[idx+self.seq_len-1]
class LSTMAttentionModel(nn.Module):
def __init__(self, input_size, hidden_size=128, num_layers=2, dropout=0.2):
super(LSTMAttentionModel, self).__init__()
self.hidden_size = hidden_size
self.num_layers = num_layers
# LSTM层
self.lstm = nn.LSTM(
input_size=input_size,
hidden_size=hidden_size,
num_layers=num_layers,
batch_first=True,
dropout=dropout if num_layers > 1 else 0
)
# 注意力机制
self.attention = nn.Sequential(
nn.Linear(hidden_size, hidden_size),
nn.Tanh(),
nn.Linear(hidden_size, 1)
)
# 输出层
self.fc = nn.Sequential(
nn.Linear(hidden_size, 64),
nn.ReLU(),
nn.Dropout(dropout),
nn.Linear(64, 1)
)
def forward(self, x):
# LSTM前向传播
lstm_out, (hidden, cell) = self.lstm(x)
# 注意力权重
attention_weights = self.attention(lstm_out)
attention_weights = torch.softmax(attention_weights, dim=1)
# 加权求和
context_vector = torch.sum(attention_weights * lstm_out, dim=1)
# 全连接层
output = self.fc(context_vector)
return output.squeeze(-1)
class DeepLearningPredictor:
def __init__(self, input_size, hidden_size=128, num_layers=2, lr=0.001):
self.model = LSTMAttentionModel(input_size, hidden_size, num_layers)
self.device = torch.device('cuda' if torch.cuda.is_available() else 'cpu')
self.model.to(self.device)
self.optimizer = optim.Adam(self.model.parameters(), lr=lr)
self.criterion = nn.MSELoss()
self.history = {'train_loss': [], 'val_loss': []}
def train(self, train_loader, val_loader, epochs=50, patience=10):
"""
训练模型
"""
best_val_loss = float('inf')
patience_counter = 0
for epoch in range(epochs):
# 训练阶段
self.model.train()
train_loss = 0
for batch_X, batch_y in train_loader:
batch_X, batch_y = batch_X.to(self.device), batch_y.to(self.device)
self.optimizer.zero_grad()
outputs = self.model(batch_X)
loss = self.criterion(outputs, batch_y)
loss.backward()
self.optimizer.step()
train_loss += loss.item()
# 验证阶段
self.model.eval()
val_loss = 0
with torch.no_grad():
for batch_X, batch_y in val_loader:
batch_X, batch_y = batch_X.to(self.device), batch_y.to(self.device)
outputs = self.model(batch_X)
loss = self.criterion(outputs, batch_y)
val_loss += loss.item()
train_loss /= len(train_loader)
val_loss /= len(val_loader)
self.history['train_loss'].append(train_loss)
self.history['val_loss'].append(val_loss)
if (epoch + 1) % 10 == 0:
print(f'Epoch [{epoch+1}/{epochs}], Train Loss: {train_loss:.4f}, Val Loss: {val_loss:.4f}')
# 早停机制
if val_loss < best_val_loss:
best_val_loss = val_loss
patience_counter = 0
# 保存最佳模型
torch.save(self.model.state_dict(), 'best_model.pth')
else:
patience_counter += 1
if patience_counter >= patience:
print(f"Early stopping at epoch {epoch+1}")
break
# 加载最佳模型
self.model.load_state_dict(torch.load('best_model.pth'))
return self.history
def predict(self, X, sequence_length=30):
"""
预测
"""
self.model.eval()
# 创建数据集
dataset = MovieDataset(X, np.zeros(len(X)), sequence_length)
loader = DataLoader(dataset, batch_size=32, shuffle=False)
predictions = []
with torch.no_grad():
for batch_X, _ in loader:
batch_X = batch_X.to(self.device)
outputs = self.model(batch_X)
predictions.extend(outputs.cpu().numpy())
return np.array(predictions)
# 使用示例
# 准备时间序列数据
def prepare_lstm_data(df, feature_cols, target_col, seq_length=30):
"""
准备LSTM训练数据
"""
from sklearn.preprocessing import StandardScaler
# 标准化
scaler = StandardScaler()
scaled_features = scaler.fit_transform(df[feature_cols])
scaled_target = scaler.fit_transform(df[[target_col]])
# 创建序列
X, y = [], []
for i in range(len(df) - seq_length):
X.append(scaled_features[i:i+seq_length])
y.append(scaled_target[i+seq_length])
return np.array(X), np.array(y).flatten(), scaler
# 假设我们有足够的数据
if len(simulated_data) > 50:
X_lstm, y_lstm, target_scaler = prepare_lstm_data(
simulated_data, feature_columns_ts, 'total_box_office', seq_length=30
)
# 划分数据集
train_size = int(0.8 * len(X_lstm))
X_train_lstm = X_lstm[:train_size]
y_train_lstm = y_lstm[:train_size]
X_val_lstm = X_lstm[train_size:]
y_val_lstm = y_lstm[train_size:]
# 创建数据加载器
train_dataset = MovieDataset(pd.DataFrame(X_train_lstm.reshape(X_train_lstm.shape[0], -1)), y_train_lstm, sequence_length=1)
val_dataset = MovieDataset(pd.DataFrame(X_val_lstm.reshape(X_val_lstm.shape[0], -1)), y_val_lstm, sequence_length=1)
train_loader = DataLoader(train_dataset, batch_size=32, shuffle=True)
val_loader = DataLoader(val_dataset, batch_size=32, shuffle=False)
# 训练模型
dl_predictor = DeepLearningPredictor(input_size=X_train_lstm.shape[2], hidden_size=64, num_layers=2)
history = dl_predictor.train(train_loader, val_loader, epochs=30)
# 预测
dl_predictions = dl_predictor.predict(pd.DataFrame(X_val_lstm.reshape(X_val_lstm.shape[0], -1)))
# 反标准化
dl_predictions = target_scaler.inverse_transform(dl_predictions.reshape(-1, 1)).flatten()
y_val_original = target_scaler.inverse_transform(y_val_lstm.reshape(-1, 1)).flatten()
print("\n深度学习模型评估:")
print(f"MAE: {mean_absolute_error(y_val_original, dl_predictions):.2f}")
print(f"R2: {r2_score(y_val_original, dl_predictions):.4f}")
else:
print("数据量不足,跳过LSTM模型训练")
四、模型优化与超参数调优
4.1 贝叶斯优化
使用贝叶斯优化进行超参数调优,比网格搜索更高效:
from skopt import BayesSearchCV
from skopt.space import Real, Integer, Categorical
class HyperparameterOptimizer:
def __init__(self):
self.best_params = None
self.best_score = None
def optimize_xgboost(self, X, y, n_iter=50):
"""
优化XGBoost超参数
"""
# 定义搜索空间
search_space = {
'n_estimators': Integer(100, 500),
'max_depth': Integer(3, 10),
'learning_rate': Real(0.01, 0.3, prior='log-uniform'),
'subsample': Real(0.6, 1.0),
'colsample_bytree': Real(0.6, 1.0),
'gamma': Real(0, 5),
'reg_alpha': Real(0, 10),
'reg_lambda': Real(0, 10)
}
# 创建模型
model = xgb.XGBRegressor(random_state=42, objective='reg:squarederror')
# 贝叶斯搜索
opt = BayesSearchCV(
model,
search_space,
n_iter=n_iter,
cv=3,
n_jobs=-1,
random_state=42,
scoring='r2'
)
print("开始贝叶斯优化...")
opt.fit(X, y)
self.best_params = opt.best_params_
self.best_score = opt.best_score_
print(f"最佳参数: {opt.best_params_}")
print(f"最佳分数: {opt.best_score_:.4f}")
return opt.best_estimator_
def optimize_lightgbm(self, X, y, n_iter=50):
"""
优化LightGBM超参数
"""
search_space = {
'n_estimators': Integer(100, 500),
'max_depth': Integer(3, 10),
'learning_rate': Real(0.01, 0.3, prior='log-uniform'),
'num_leaves': Integer(20, 100),
'subsample': Real(0.6, 1.0),
'colsample_bytree': Real(0.6, 1.0),
'reg_alpha': Real(0, 10),
'reg_lambda': Real(0, 10)
}
model = lgb.LGBMRegressor(random_state=42)
opt = BayesSearchCV(
model,
search_space,
n_iter=n_iter,
cv=3,
n_jobs=-1,
random_state=42,
scoring='r2'
)
print("开始LightGBM贝叶斯优化...")
opt.fit(X, y)
print(f"最佳参数: {opt.best_params_}")
print(f"最佳分数: {opt.best_score_:.4f}")
return opt.best_estimator_
# 使用示例
optimizer = HyperparameterOptimizer()
# 优化XGBoost
best_xgb = optimizer.optimize_xgboost(X_train, y_train, n_iter=20)
# 优化LightGBM
best_lgb = optimizer.optimize_lightgbm(X_train, y_train, n_iter=20)
# 评估优化后的模型
for name, model in [('XGBoost', best_xgb), ('LightGBM', best_lgb)]:
y_pred = model.predict(X_test)
print(f"\n{name}优化后评估:")
print(f"MAE: {mean_absolute_error(y_test, y_pred):.2f}")
print(f"R2: {r2_score(y_test, y_pred):.4f}")
4.2 模型集成与Stacking
from sklearn.base import BaseEstimator, RegressorMixin
from sklearn.model_selection import KFold
class StackingEnsemble(BaseEstimator, RegressorMixin):
def __init__(self, base_models, meta_model, n_folds=5):
self.base_models = base_models
self.meta_model = meta_model
self.n_folds = n_folds
self.models = []
self.meta_features = None
def fit(self, X, y):
"""
训练Stacking集成模型
"""
kf = KFold(n_splits=self.n_folds, shuffle=True, random_state=42)
# 存储每个基础模型在每个fold的预测
meta_features = np.zeros((X.shape[0], len(self.base_models)))
for i, model in enumerate(self.base_models):
print(f"训练基础模型 {i+1}/{len(self.base_models)}")
for train_idx, val_idx in kf.split(X):
X_train_fold, X_val_fold = X.iloc[train_idx], X.iloc[val_idx]
y_train_fold = y.iloc[train_idx]
# 训练模型
model_clone = clone(model)
model_clone.fit(X_train_fold, y_train_fold)
# 预测验证集
pred = model_clone.predict(X_val_fold)
meta_features[val_idx, i] = pred
# 存储模型用于后续预测
if len(self.models) <= i:
self.models.append([])
self.models[i].append(model_clone)
# 训练元模型
self.meta_model.fit(meta_features, y)
self.meta_features = meta_features
return self
def predict(self, X):
"""
预测
"""
# 获取基础模型预测
meta_test = np.zeros((X.shape[0], len(self.base_models)))
for i, model_list in enumerate(self.models):
# 对每个fold的模型取平均
fold_preds = []
for model in model_list:
fold_preds.append(model.predict(X))
meta_test[:, i] = np.mean(fold_preds, axis=0)
# 元模型预测
return self.meta_model.predict(meta_test)
# 使用示例
# 创建基础模型
base_models = [
xgb.XGBRegressor(n_estimators=100, random_state=42),
lgb.LGBMRegressor(n_estimators=100, random_state=42),
RandomForestRegressor(n_estimators=100, random_state=42),
GradientBoostingRegressor(n_estimators=100, random_state=42)
]
# 创建元模型
meta_model = Ridge(alpha=1.0)
# 创建Stacking集成
stacking_model = StackingEnsemble(base_models, meta_model, n_folds=5)
# 训练
stacking_model.fit(X_train, y_train)
# 预测
stacking_pred = stacking_model.predict(X_test)
# 评估
print("\nStacking集成模型评估:")
print(f"MAE: {mean_absolute_error(y_test, stacking_pred):.2f}")
print(f"R2: {r2_score(y_test, stacking_pred):.4f}")
五、实时预测与监控系统
5.1 实时数据流处理
import asyncio
import aiohttp
from datetime import datetime, timedelta
import redis
class RealTimePredictor:
def __init__(self, model, redis_client=None):
self.model = model
self.redis_client = redis_client or redis.Redis(host='localhost', port=6379, db=0)
self.feature_cache = {}
async def fetch_real_time_data(self, movie_id):
"""
异步获取实时数据
"""
async with aiohttp.ClientSession() as session:
# 模拟API调用
# 实际应用中需要真实的API接口
tasks = [
self._fetch_social_media(session, movie_id),
self._fetch_box_office(session, movie_id),
self._fetch_reviews(session, movie_id)
]
results = await asyncio.gather(*tasks)
# 合并数据
combined_data = {}
for result in results:
combined_data.update(result)
return combined_data
async def _fetch_social_media(self, session, movie_id):
"""
获取社交媒体数据
"""
# 模拟API调用
await asyncio.sleep(0.1) # 模拟网络延迟
return {
'weibo_mentions': np.random.randint(1000, 5000),
'douyin_views': np.random.randint(100000, 500000),
'mention_growth_rate': np.random.uniform(0.1, 0.5)
}
async def _fetch_box_office(self, session, movie_id):
"""
获取票房数据
"""
await asyncio.sleep(0.1)
return {
'daily_box_office': np.random.randint(1000000, 5000000),
'cumulative_box_office': np.random.randint(10000000, 50000000)
}
async def _fetch_reviews(self, session, movie_id):
"""
获取评论数据
"""
await asyncio.sleep(0.1)
return {
'avg_rating': np.random.uniform(7.0, 9.0),
'rating_count': np.random.randint(1000, 10000),
'positive_ratio': np.random.uniform(0.6, 0.9)
}
def update_features(self, movie_id, real_time_data):
"""
更新特征
"""
# 从缓存获取历史数据
cache_key = f"movie:{movie_id}:features"
cached_data = self.redis_client.get(cache_key)
if cached_data:
features = json.loads(cached_data)
else:
features = {}
# 更新实时特征
features.update(real_time_data)
features['last_update'] = datetime.now().isoformat()
# 存入缓存
self.redis_client.setex(cache_key, 3600, json.dumps(features))
return features
async def predict_next_day(self, movie_id):
"""
预测下一天票房
"""
# 获取实时数据
real_time_data = await self.fetch_real_time_data(movie_id)
# 更新特征
features = self.update_features(movie_id, real_time_data)
# 转换为模型输入格式
feature_vector = self._vectorize_features(features)
# 预测
prediction = self.model.predict(feature_vector.reshape(1, -1))[0]
# 计算置信区间
confidence = self._calculate_confidence(features)
return {
'movie_id': movie_id,
'predicted_box_office': float(prediction),
'confidence_interval': confidence,
'prediction_time': datetime.now().isoformat(),
'features': features
}
def _vectorize_features(self, features):
"""
将特征字典转换为特征向量
"""
# 这里需要根据模型期望的特征顺序
expected_features = ['cast_power', 'director_power', 'ip_power', 'budget_log',
'total_engagement', 'avg_rating', 'sentiment_score',
'is_weekend', 'is_holiday']
vector = []
for feature in expected_features:
value = features.get(feature, 0)
vector.append(value)
return np.array(vector)
def _calculate_confidence(self, features):
"""
计算置信区间(简化版)
"""
# 基于特征质量的置信度
quality_score = 0
if features.get('rating_count', 0) > 1000:
quality_score += 0.3
if features.get('weibo_mentions', 0) > 1000:
quality_score += 0.3
if features.get('avg_rating', 0) > 7:
quality_score += 0.2
# 基础误差范围
base_error = 0.15 # 15%
# 根据质量调整
error_range = base_error * (1 - quality_score * 0.5)
return error_range
# 使用示例
async def main():
# 假设我们有一个训练好的模型
# model = best_xgb # 使用之前训练的模型
# 创建实时预测器
# predictor = RealTimePredictor(model)
# 预测单部电影
# result = await predictor.predict_next_day('movie_123')
# print(json.dumps(result, indent=2, ensure_ascii=False))
# 模拟运行
print("实时预测系统示例:")
print("系统已启动,等待实时数据流...")
# 模拟预测
mock_result = {
'movie_id': 'movie_123',
'predicted_box_office': 25000000,
'confidence_interval': 0.12,
'prediction_time': datetime.now().isoformat(),
'features': {
'cast_power': 7.5,
'director_power': 8.0,
'ip_power': 9.0,
'total_engagement': 45000,
'avg_rating': 8.2
}
}
print(json.dumps(mock_result, indent=2, ensure_ascii=False))
# 运行示例
# asyncio.run(main())
print("提示:实际运行需要配置Redis和真实API接口")
5.2 模型监控与漂移检测
from evidently import ColumnMapping
from evidently.report import Report
from evidently.metric_preset import DataDriftPreset, TargetDriftPreset
import warnings
warnings.filterwarnings('ignore')
class ModelMonitor:
def __init__(self, reference_data, model_name):
self.reference_data = reference_data
self.model_name = model_name
self.drift_history = []
def detect_data_drift(self, current_data, threshold=0.15):
"""
检测数据漂移
"""
from evidently.metrics import DataDriftTable
# 创建漂移报告
report = Report(metrics=[DataDriftTable()])
# 计算漂移
report.run(
reference_data=self.reference_data,
current_data=current_data
)
# 获取漂移分数
drift_score = report.as_dict()['metrics'][0]['result']['drift_share']
is_drift = drift_score > threshold
# 记录历史
self.drift_history.append({
'timestamp': datetime.now(),
'drift_score': drift_score,
'is_drift': is_drift
})
return {
'drift_detected': is_drift,
'drift_score': drift_score,
'threshold': threshold,
'features_drifted': report.as_dict()['metrics'][0]['result']['drifted_features']
}
def detect_target_drift(self, current_data):
"""
检测目标变量漂移
"""
from evidently.metrics import TargetDriftTable
report = Report(metrics=[TargetDriftTable()])
report.run(
reference_data=self.reference_data,
current_data=current_data
)
# 获取目标分布变化
target_drift = report.as_dict()['metrics'][0]['result']['drift_detected']
return {
'target_drift_detected': target_drift,
'details': report.as_dict()['metrics'][0]['result']
}
def generate_monitoring_report(self, current_data):
"""
生成完整监控报告
"""
data_drift = self.detect_data_drift(current_data)
target_drift = self.detect_target_drift(current_data)
report = {
'model_name': self.model_name,
'timestamp': datetime.now().isoformat(),
'data_drift': data_drift,
'target_drift': target_drift,
'recommendation': self._generate_recommendation(data_drift, target_drift)
}
return report
def _generate_recommendation(self, data_drift, target_drift):
"""
根据漂移检测结果生成建议
"""
recommendations = []
if data_drift['drift_detected']:
recommendations.append("数据分布发生显著变化,建议重新训练模型")
recommendations.append("检查特征工程逻辑是否需要调整")
if target_drift['target_drift_detected']:
recommendations.append("目标变量分布发生变化,可能需要调整模型")
recommendations.append("分析市场环境是否发生重大变化")
if not recommendations:
recommendations.append("模型表现稳定,继续监控")
return recommendations
# 使用示例
# 创建参考数据(训练集)
reference_data = pd.DataFrame({
'cast_power': np.random.uniform(0, 10, 100),
'director_power': np.random.uniform(0, 10, 100),
'ip_power': np.random.uniform(0, 10, 100),
'total_box_office': np.random.lognormal(17, 0.5, 100)
})
# 创建当前数据(新数据)
current_data = pd.DataFrame({
'cast_power': np.random.uniform(0, 10, 50),
'director_power': np.random.uniform(0, 10, 50),
'ip_power': np.random.uniform(0, 10, 50),
'total_box_office': np.random.lognormal(17, 0.5, 50)
})
# 创建监控器
monitor = ModelMonitor(reference_data, "BoxOfficePredictor_v1")
# 检测漂移
report = monitor.generate_monitoring_report(current_data)
print("模型监控报告:")
print(json.dumps(report, indent=2, ensure_ascii=False))
六、观众真实喜好趋势分析
6.1 基于评论的情感分析
import jieba
from collections import Counter
from sklearn.feature_extraction.text import TfidfVectorizer
from textblob import TextBlob
import re
class AudiencePreferenceAnalyzer:
def __init__(self):
self.tfidf_vectorizer = TfidfVectorizer(max_features=1000, stop_words=None)
self.sentiment_lexicon = {
'positive': ['好', '棒', '精彩', '喜欢', '推荐', '感动', '震撼', '优秀', '完美', '值得'],
'negative': ['差', '烂', '失望', '无聊', '尴尬', '垃圾', '难看', '后悔', '浪费', '糟糕']
}
def preprocess_text(self, text):
"""
文本预处理
"""
# 去除特殊字符
text = re.sub(r'[^\w\s]', '', text)
# 分词
words = jieba.lcut(text)
# 去除停用词
stopwords = ['的', '了', '是', '在', '我', '就', '都', '和', '这', '那', '也', '很', '看', '觉得']
words = [w for w in words if w not in stopwords and len(w) > 1]
return words
def analyze_sentiment(self, text):
"""
分析文本情感
"""
words = self.preprocess_text(text)
# 基于词典的情感分析
positive_count = sum(1 for w in words if w in self.sentiment_lexicon['positive'])
negative_count = sum(1 for w in words if w in self.sentiment_lexicon['negative'])
# 使用TextBlob进行辅助分析
blob = TextBlob(text)
textblob_score = blob.sentiment.polarity
# 综合得分
sentiment_score = (positive_count - negative_count) * 0.1 + textblob_score
return {
'score': sentiment_score,
'positive_count': positive_count,
'negative_count': negative_count,
'label': 'positive' if sentiment_score > 0.1 else 'negative' if sentiment_score < -0.1 else 'neutral'
}
def extract_key_topics(self, reviews, top_n=10):
"""
提取关键话题
"""
# TF-IDF分析
reviews_processed = [' '.join(self.preprocess_text(r)) for r in reviews]
tfidf_matrix = self.tfidf_vectorizer.fit_transform(reviews_processed)
feature_names = self.tfidf_vectorizer.get_feature_names_out()
# 计算平均TF-IDF分数
mean_tfidf = np.array(tfidf_matrix.mean(axis=0)).flatten()
top_indices = mean_tfidf.argsort()[-top_n:][::-1]
topics = [(feature_names[i], mean_tfidf[i]) for i in top_indices]
return topics
def analyze_audience_preference(self, reviews, ratings):
"""
综合分析观众喜好
"""
# 情感分析
sentiments = [self.analyze_sentiment(review) for review in reviews]
avg_sentiment = np.mean([s['score'] for s in sentiments])
# 话题提取
topics = self.extract_key_topics(reviews)
# 评分分布
rating_dist = Counter(ratings)
# 构建偏好画像
preference_profile = {
'overall_sentiment': avg_sentiment,
'sentiment_distribution': {
'positive': sum(1 for s in sentiments if s['label'] == 'positive'),
'negative': sum(1 for s in sentiments if s['label'] == 'negative'),
'neutral': sum(1 for s in sentiments if s['label'] == 'neutral')
},
'top_topics': topics,
'rating_distribution': dict(rating_dist),
'key_phrases': self._extract_key_phrases(reviews)
}
return preference_profile
def _extract_key_phrases(self, reviews, top_n=15):
"""
提取关键短语
"""
all_phrases = []
for review in reviews:
words = self.preprocess_text(review)
# 提取2-3词短语
for i in range(len(words)-1):
all_phrases.append(words[i] + words[i+1])
for i in range(len(words)-2):
all_phrases.append(words[i] + words[i+1] + words[i+2])
# 统计频率
phrase_counts = Counter(all_phrases)
return phrase_counts.most_common(top_n)
# 使用示例
analyzer = AudiencePreferenceAnalyzer()
# 模拟评论数据
sample_reviews = [
"这部电影特效很棒,剧情也很精彩,强烈推荐!",
"演员演技在线,但剧情有点拖沓,整体还行",
"非常失望,浪费时间和金钱,不建议观看",
"视觉效果震撼,故事感人,值得二刷",
"中规中矩,没有特别出彩的地方",
"导演功力深厚,每个镜头都很美",
"节奏太慢了,看到后面有点困",
"演员阵容强大,表演出色,但剧本一般",
"超出预期,比预告片好看多了",
"适合全家观看,孩子很喜欢"
]
sample_ratings = [9, 7, 3, 10, 6, 8, 5, 7, 9, 8]
# 分析
preference = analyzer.analyze_audience_preference(sample_reviews, sample_ratings)
print("观众喜好分析报告:")
print(json.dumps(preference, indent=2, ensure_ascii=False))
6.2 基于聚类的观众细分
from sklearn.cluster import KMeans
from sklearn.decomposition import PCA
import matplotlib.pyplot as plt
import seaborn as sns
class AudienceSegmentation:
def __init__(self, n_clusters=4):
self.n_clusters = n_clusters
self.kmeans = KMeans(n_clusters=n_clusters, random_state=42)
self.pca = PCA(n_components=2)
self.segment_labels = None
def create_audience_features(self, user_data):
"""
创建观众特征
"""
features = []
for user in user_data:
feature_vector = [
user.get('avg_rating', 0),
user.get('review_count', 0),
user.get('genre_preference_score', 0),
user.get('frequency_score', 0),
user.get('social_influence', 0)
]
features.append(feature_vector)
return np.array(features)
def fit_segmentation(self, audience_features):
"""
训练聚类模型
"""
# 标准化
from sklearn.preprocessing import StandardScaler
scaler = StandardScaler()
scaled_features = scaler.fit_transform(audience_features)
# 聚类
clusters = self.kmeans.fit_predict(scaled_features)
# 降维用于可视化
self.pca_features = self.pca.fit_transform(scaled_features)
# 为每个聚类分配有意义的标签
self.segment_labels = self._assign_segment_labels(audience_features, clusters)
return clusters
def _assign_segment_labels(self, features, clusters):
"""
为聚类分配标签
"""
labels = {}
for cluster_id in range(self.n_clusters):
cluster_mask = clusters == cluster_id
cluster_features = features[cluster_mask]
# 计算聚类中心特征
center = cluster_features.mean(axis=0)
# 根据特征分配标签
if center[0] > 7.5 and center[2] > 0.7: # 高评分,类型偏好强
label = "核心影迷"
elif center[0] > 6.5 and center[3] > 0.6: # 中等评分,频率高
label = "常客观众"
elif center[0] < 6.0 and center[4] > 0.5: # 低评分,社交影响高
label = "意见领袖"
else:
label = "普通观众"
labels[cluster_id] = label
return labels
def visualize_segments(self, clusters):
"""
可视化观众细分
"""
plt.figure(figsize=(12, 8))
# 创建颜色映射
colors = ['#FF6B6B', '#4ECDC4', '#45B7D1', '#96CEB4', '#FFEAA7']
for cluster_id in range(self.n_clusters):
mask = clusters == cluster_id
plt.scatter(
self.pca_features[mask, 0],
self.pca_features[mask, 1],
c=colors[cluster_id % len(colors)],
label=self.segment_labels[cluster_id],
alpha=0.7,
s=50
)
plt.title('观众细分聚类可视化', fontsize=16)
plt.xlabel(f'PC1 ({self.pca.explained_variance_ratio_[0]:.2%} variance)')
plt.ylabel(f'PC2 ({self.pca.explained_variance_ratio_[1]:.2%} variance)')
plt.legend()
plt.grid(True, alpha=0.3)
plt.tight_layout()
plt.show()
def get_segment_insights(self, clusters, user_data):
"""
获取细分洞察
"""
insights = {}
for cluster_id in range(self.n_clusters):
mask = clusters == cluster_id
segment_users = [user_data[i] for i in np.where(mask)[0]]
# 计算统计信息
avg_rating = np.mean([u.get('avg_rating', 0) for u in segment_users])
avg_review_count = np.mean([u.get('review_count', 0) for u in segment_users])
size = len(segment_users)
insights[self.segment_labels[cluster_id]] = {
'size': size,
'percentage': size / len(clusters) * 100,
'avg_rating': avg_rating,
'avg_review_count': avg_review_count,
'characteristics': self._describe_segment(segment_users)
}
return insights
def _describe_segment(self, users):
"""
描述细分特征
"""
# 简化的特征描述
if len(users) == 0:
return "暂无特征描述"
# 计算主要特征
avg_rating = np.mean([u.get('avg_rating', 0) for u in users])
avg_reviews = np.mean([u.get('review_count', 0) for u in users])
if avg_rating > 7.5:
return "对电影质量要求高,喜欢深度评论"
elif avg_rating > 6.5:
return "观影频率高,对类型片有偏好"
elif avg_rating < 6.0:
return "喜欢发表观点,社交活跃"
else:
return "随性观影,受口碑影响"
# 使用示例
# 模拟观众数据
np.random.seed(42)
n_users = 200
audience_data = []
for i in range(n_users):
user = {
'user_id': i,
'avg_rating': np.random.normal(7.5, 1.5),
'review_count': np.random.poisson(5),
'genre_preference_score': np.random.beta(2, 2),
'frequency_score': np.random.beta(1.5, 2),
'social_influence': np.random.beta(1, 3)
}
audience_data.append(user)
# 创建特征矩阵
segmenter = AudienceSegmentation(n_clusters=4)
features = segmenter.create_audience_features(audience_data)
# 训练聚类
clusters = segmenter.fit_segmentation(features)
# 可视化
# segmenter.visualize_segments(clusters) # 需要matplotlib环境
# 获取洞察
insights = segmenter.get_segment_insights(clusters, audience_data)
print("观众细分洞察:")
for segment, info in insights.items():
print(f"\n{segment} (占比 {info['percentage']:.1f}%):")
print(f" 平均评分: {info['avg_rating']:.2f}")
print(f" 平均评论数: {info['avg_review_count']:.1f}")
print(f" 特征: {info['characteristics']}")
6.3 基于协同过滤的推荐系统
from scipy.sparse import csr_matrix
from sklearn.neighbors import NearestNeighbors
class MovieRecommender:
def __init__(self):
self.user_movie_matrix = None
self.model = None
self.movie_index_map = {}
self.user_index_map = {}
def build_user_movie_matrix(self, ratings_data):
"""
构建用户-电影评分矩阵
"""
# 创建索引映射
unique_users = ratings_data['user_id'].unique()
unique_movies = ratings_data['movie_id'].unique()
self.user_index_map = {user: idx for idx, user in enumerate(unique_users)}
self.movie_index_map = {movie: idx for idx, movie in enumerate(unique_movies)}
# 构建稀疏矩阵
rows = [self.user_index_map[user] for user in ratings_data['user_id']]
cols = [self.movie_index_map[movie] for movie in ratings_data['movie_id']]
values = ratings_data['rating'].values
self.user_movie_matrix = csr_matrix((values, (rows, cols)),
shape=(len(unique_users), len(unique_movies)))
return self.user_movie_matrix
def train_knn_model(self, n_neighbors=20):
"""
训练KNN模型
"""
# 使用余弦相似度
self.model = NearestNeighbors(metric='cosine', algorithm='brute', n_neighbors=n_neighbors)
self.model.fit(self.user_movie_matrix)
return self.model
def recommend_for_user(self, user_id, top_n=10):
"""
为用户推荐电影
"""
if user_id not in self.user_index_map:
return []
user_idx = self.user_index_map[user_id]
user_ratings = self.user_movie_matrix[user_idx]
# 找到相似用户
distances, indices = self.model.kneighbors(user_ratings, n_neighbors=top_n+1)
# 获取相似用户的评分
similar_users_ratings = self.user_movie_matrix[indices[0][1:]] # 排除自己
# 计算推荐分数(加权平均)
sim_scores = 1 - distances[0][1:] # 转换为相似度
weighted_ratings = similar_users_ratings.multiply(sim_scores[:, np.newaxis])
recommended_scores = np.array(weighted_ratings.sum(axis=0)).flatten()
# 获取top N推荐
top_indices = recommended_scores.argsort()[-top_n:][::-1]
# 转换回电影ID
idx_to_movie = {v: k for k, v in self.movie_index_map.items()}
recommendations = [(idx_to_movie[idx], recommended_scores[idx]) for idx in top_indices]
return recommendations
def recommend_similar_movies(self, movie_id, top_n=10):
"""
推荐相似电影
"""
if movie_id not in self.movie_index_map:
return []
movie_idx = self.movie_index_map[movie_id]
# 计算电影相似度(基于用户评分模式)
movie_ratings = self.user_movie_matrix[:, movie_idx]
# 转换为稠密矩阵进行计算
if movie_ratings.nnz == 0:
return []
# 计算所有电影与目标电影的相似度
similarities = []
for other_idx in range(self.user_movie_matrix.shape[1]):
if other_idx == movie_idx:
continue
other_ratings = self.user_movie_matrix[:, other_idx]
# 计算共同评分用户的相似度
common_users = movie_ratings.multiply(other_ratings)
if common_users.nnz > 2: # 至少有2个共同评分用户
# 计算余弦相似度
dot_product = common_users.power(2).sum()
norm_movie = movie_ratings.power(2).sum()
norm_other = other_ratings.power(2).sum()
if norm_movie > 0 and norm_other > 0:
similarity = dot_product / (np.sqrt(norm_movie) * np.sqrt(norm_other))
similarities.append((other_idx, similarity))
# 排序
similarities.sort(key=lambda x: x[1], reverse=True)
# 转换回电影ID
idx_to_movie = {v: k for k, v in self.movie_index_map.items()}
recommendations = [(idx_to_movie[idx], sim) for idx, sim in similarities[:top_n]]
return recommendations
# 使用示例
# 模拟评分数据
np.random.seed(42)
n_users = 100
n_movies = 50
ratings_data = pd.DataFrame({
'user_id': np.random.randint(0, n_users, 500),
'movie_id': np.random.randint(0, n_movies, 500),
'rating': np.random.randint(1, 11, 500) # 1-10分
})
# 创建推荐器
recommender = MovieRecommender()
# 构建矩阵
matrix = recommender.build_user_movie_matrix(ratings_data)
# 训练模型
recommender.train_knn_model(n_neighbors=10)
# 为用户推荐
user_recommendations = recommender.recommend_for_user(user_id=5, top_n=5)
print("用户5的电影推荐:")
for movie_id, score in user_recommendations:
print(f" 电影{movie_id}: 相似度 {score:.3f}")
# 推荐相似电影
similar_movies = recommender.recommend_similar_movies(movie_id=10, top_n=5)
print("\n电影10的相似电影:")
for movie_id, sim in similar_movies:
print(f" 电影{movie_id}: 相似度 {sim:.3f}")
七、完整案例:从数据到预测
7.1 端到端的票房预测流程
class EndToEndBoxOfficePredictor:
def __init__(self):
self.data_collector = MovieDataCollector()
self.preprocessor = DataPreprocessor()
self.feature_engineer = FeatureEngineer()
self.predictor = BoxOfficePredictor()
self.advanced_predictor = AdvancedBoxOfficePredictor()
self.optimizer = HyperparameterOptimizer()
self.monitor = None
def run_full_pipeline(self, movie_info, social_data, review_data, historical_data=None):
"""
运行完整的预测流程
"""
print("=" * 60)
print("开始端到端票房预测流程")
print("=" * 60)
# 1. 特征工程
print("\n1. 特征工程...")
features = {}
features.update(self.feature_engineer.create_movie_features(movie_info))
features.update(self.feature_engineer.create_social_media_features(social_data))
features.update(self.feature_engineer.create_review_features(review_data))
# 添加时间特征(假设今天是上映前7天)
from datetime import datetime
today = datetime.now()
features['days_before_release'] = 7
features['is_weekend'] = 1 if today.weekday() in [5, 6] else 0
features['is_holiday'] = 0 # 根据实际情况调整
print("提取的特征:")
for k, v in list(features.items())[:10]:
print(f" {k}: {v}")
# 2. 转换为模型输入
feature_vector = pd.DataFrame([features])
# 3. 模型预测(使用集成模型)
print("\n2. 模型预测...")
# 模拟训练好的模型(实际使用时需要训练)
# 这里使用简单的线性模型作为示例
from sklearn.linear_model import LinearRegression
# 模拟训练数据
if historical_data is None:
# 创建模拟历史数据
np.random.seed(42)
n_samples = 100
historical_features = pd.DataFrame({
'cast_power': np.random.uniform(0, 10, n_samples),
'director_power': np.random.uniform(0, 10, n_samples),
'ip_power': np.random.uniform(0, 10, n_samples),
'budget_log': np.random.uniform(15, 20, n_samples),
'total_engagement': np.random.uniform(0, 100000, n_samples),
'avg_rating': np.random.uniform(6, 9.5, n_samples),
'sentiment_score': np.random.uniform(-1, 1, n_samples),
'is_weekend': np.random.randint(0, 2, n_samples),
'is_holiday': np.random.randint(0, 2, n_samples),
'total_box_office': np.random.lognormal(17, 0.5, n_samples) +
np.random.uniform(0, 5000000, n_samples)
})
# 训练模拟模型
model = LinearRegression()
model.fit(historical_features.drop('total_box_office', axis=1),
historical_features['total_box_office'])
else:
# 使用历史数据训练
model = LinearRegression()
model.fit(historical_data.drop('total_box_office', axis=1),
historical_data['total_box_office'])
# 预测
prediction = model.predict(feature_vector)[0]
# 计算置信区间(简化)
confidence_range = prediction * 0.15 # 15%误差范围
print(f"预测票房: {prediction:,.0f} 元")
print(f"置信区间: {prediction - confidence_range:,.0f} - {prediction + confidence_range:,.0f} 元")
# 3. 风险评估
print("\n3. 风险评估...")
risk_factors = self._assess_risk_factors(features)
for factor, level in risk_factors.items():
print(f" {factor}: {level}")
# 4. 生成建议
print("\n4. 生成建议...")
suggestions = self._generate_suggestions(features, prediction)
for suggestion in suggestions:
print(f" - {suggestion}")
# 5. 构建完整报告
report = {
'prediction': {
'point_estimate': float(prediction),
'confidence_interval': [float(prediction - confidence_range),
float(prediction + confidence_range)],
'confidence_level': 0.85
},
'features': features,
'risk_assessment': risk_factors,
'suggestions': suggestions,
'timestamp': datetime.now().isoformat()
}
print("\n" + "=" * 60)
print("预测完成!")
print("=" * 60)
return report
def _assess_risk_factors(self, features):
"""
评估风险因素
"""
risks = {}
# 社交媒体热度不足
if features.get('total_engagement', 0) < 10000:
risks['社交媒体热度'] = "低"
# 评分偏低
if features.get('avg_rating', 0) < 7.0:
risks['观众口碑'] = "中等"
# IP影响力弱
if features.get('ip_power', 0) < 5.0:
risks['IP影响力'] = "低"
# 演员阵容一般
if features.get('cast_power', 0) < 5.0:
risks['演员阵容'] = "中等"
if not risks:
risks['整体风险'] = "低"
return risks
def _generate_suggestions(self, features, prediction):
"""
生成改进建议
"""
suggestions = []
if features.get('total_engagement', 0) < 10000:
suggestions.append("加强社交媒体营销,增加话题讨论度")
if features.get('avg_rating', 0) < 7.0:
suggestions.append("考虑增加点映场次,收集早期口碑")
if features.get('sentiment_score', 0) < 0:
suggestions.append("关注负面评论,及时调整宣传策略")
if features.get('is_weekend', 0) == 0:
suggestions.append("工作日上映,关注长尾效应")
if not suggestions:
suggestions.append("当前策略较为完善,保持现有宣传节奏")
return suggestions
# 使用示例
pipeline = EndToEndBoxOfficePredictor()
# 准备输入数据
movie_info = {
'cast': ['演员A', '演员B', '演员C'],
'genre': ['动作', '科幻'],
'budget': 80000000,
'ip_power': 7.5,
'production_company': 'XX影业'
}
social_data = {
'weibo_mentions': 25000,
'douyin_views': 8000000,
'xiaohongshu_notes': 3500,
'mention_growth_rate': 0.4
}
review_data = {
'avg_rating': 8.0,
'rating_count': 80000,
'positive_ratio': 0.78,
'keywords': ['特效', '剧情', '演员', '节奏', '视觉']
}
# 运行完整流程
report = pipeline.run_full_pipeline(movie_info, social_data, review_data)
# 保存报告
import json
with open('box_office_prediction_report.json', 'w', encoding='utf-8') as f:
json.dump(report, f, ensure_ascii=False, indent=2)
print("\n完整报告已保存到 box_office_prediction_report.json")
八、总结与最佳实践
8.1 关键成功因素
- 数据质量:确保数据的准确性和完整性
- 特征工程:深入理解业务,提取有意义的特征
- 模型选择:根据数据特点选择合适的模型
- 持续监控:建立模型监控机制,及时发现问题
- 业务理解:结合市场洞察,避免纯技术驱动
8.2 常见陷阱与解决方案
| 陷阱 | 描述 | 解决方案 |
|---|---|---|
| 过拟合 | 模型在训练集表现好,测试集差 | 使用交叉验证,正则化,早停机制 |
| 数据泄露 | 使用了未来信息 | 严格划分时间序列,避免使用滞后特征 |
| 概念漂移 | 市场变化导致模型失效 | 定期重新训练,监控模型表现 |
| 特征爆炸 | 特征过多导致维度灾难 | 特征选择,降维,业务相关性筛选 |
8.3 未来发展方向
- 多模态融合:结合文本、图像、视频分析
- 强化学习:动态优化营销策略
- 图神经网络:分析观众社交网络
- 因果推断:理解营销活动的真实效果
- 实时学习:在线学习,适应市场变化
8.4 代码仓库与工具推荐
- 数据获取:Scrapy, BeautifulSoup, Selenium
- 数据处理:Pandas, NumPy, Dask
- 机器学习:Scikit-learn, XGBoost, LightGBM
- 深度学习:PyTorch, TensorFlow
- 超参数优化:Optuna, Hyperopt, scikit-optimize
- 模型监控:Evidently, MLflow, Prometheus
- 部署:FastAPI, Docker, Kubernetes
通过本文提供的完整代码示例和详细说明,读者可以构建一个超越传统猫眼票房预测的精准预测系统。关键在于理解业务本质,做好特征工程,并选择合适的模型架构。持续的监控和优化是保持模型长期有效的关键。
