引言:电影市场的复杂性与预测的必要性

电影市场是一个充满变数的复杂生态系统,受到多种因素的综合影响。传统的票房预测方法往往依赖于历史数据和简单的线性模型,但这些方法在面对快速变化的市场环境时显得力不从心。随着大数据和人工智能技术的发展,我们有机会更精准地把握电影市场的脉搏和观众的真实喜好趋势。

本文将深入探讨如何利用现代数据科学方法,超越传统的猫眼票房预测模型,构建更精准的票房预测系统。我们将从数据收集、特征工程、模型构建等多个维度进行详细分析,并提供完整的代码示例,帮助读者理解并实践这些方法。

一、电影票房预测的核心挑战

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 关键成功因素

  1. 数据质量:确保数据的准确性和完整性
  2. 特征工程:深入理解业务,提取有意义的特征
  3. 模型选择:根据数据特点选择合适的模型
  4. 持续监控:建立模型监控机制,及时发现问题
  5. 业务理解:结合市场洞察,避免纯技术驱动

8.2 常见陷阱与解决方案

陷阱 描述 解决方案
过拟合 模型在训练集表现好,测试集差 使用交叉验证,正则化,早停机制
数据泄露 使用了未来信息 严格划分时间序列,避免使用滞后特征
概念漂移 市场变化导致模型失效 定期重新训练,监控模型表现
特征爆炸 特征过多导致维度灾难 特征选择,降维,业务相关性筛选

8.3 未来发展方向

  1. 多模态融合:结合文本、图像、视频分析
  2. 强化学习:动态优化营销策略
  3. 图神经网络:分析观众社交网络
  4. 因果推断:理解营销活动的真实效果
  5. 实时学习:在线学习,适应市场变化

8.4 代码仓库与工具推荐

  • 数据获取:Scrapy, BeautifulSoup, Selenium
  • 数据处理:Pandas, NumPy, Dask
  • 机器学习:Scikit-learn, XGBoost, LightGBM
  • 深度学习:PyTorch, TensorFlow
  • 超参数优化:Optuna, Hyperopt, scikit-optimize
  • 模型监控:Evidently, MLflow, Prometheus
  • 部署:FastAPI, Docker, Kubernetes

通过本文提供的完整代码示例和详细说明,读者可以构建一个超越传统猫眼票房预测的精准预测系统。关键在于理解业务本质,做好特征工程,并选择合适的模型架构。持续的监控和优化是保持模型长期有效的关键。