Python-微信商业智能分析系统:多维度数据挖掘与预测模型
·
功能介绍
这是一个面向微信生态的商业智能分析平台,具备以下高级功能:
- 微信支付交易数据深度分析
- 用户生命周期价值(LTV)预测模型
- 商品关联规则挖掘
- 营销活动ROI计算引擎
- 实时销售仪表盘
- 客户流失预警系统
- 库存智能预测
系统架构
graph LR
A[微信支付接口] --> B[数据采集层]
C[公众号用户行为] --> B
D[小程序交互数据] --> B
B --> E[数据湖]
E --> F[批处理分析]
E --> G[流处理分析]
F --> H[预测模型]
G --> I[实时仪表盘]
H --> J[商业决策]
I --> J
核心代码实现
import pandas as pd
import numpy as np
from pyspark.sql import SparkSession
from pyspark.ml import Pipeline
from pyspark.ml.feature import VectorAssembler, StringIndexer
from pyspark.ml.regression import RandomForestRegressor
from pyspark.ml.evaluation import RegressionEvaluator
from pyspark.ml.clustering import KMeans
from pyspark.ml.fpm import FPGrowth
from flask import Flask, render_template
from flask_socketio import SocketIO
import plotly.graph_objects as go
from datetime import datetime, timedelta
import joblib
from wechatpay import WeChatPay
from minio import Minio
import json
class WeChatBusinessIntelligence:
"""微信商业智能分析核心类"""
def __init__(self):
# 初始化Spark会话
self.spark = SparkSession.builder \
.appName("WeChatBI") \
.config("spark.executor.memory", "8g") \
.config("spark.driver.memory", "4g") \
.getOrCreate()
# 初始化微信支付客户端
self.wxpay = WeChatPay(
mch_id='your_mch_id',
api_key='your_api_key',
cert_path='/path/to/cert'
)
# 初始化MinIO客户端
self.minio = Minio(
"minio.example.com",
access_key="minioadmin",
secret_key="minioadmin",
secure=False
)
# 加载预训练模型
self.load_models()
# 初始化Flask应用
self.app = Flask(__name__)
self.socketio = SocketIO(self.app)
self.setup_routes()
def load_models(self):
"""加载预训练模型"""
try:
# 从MinIO加载模型
self.minio.fget_object(
"models", "ltv_rf.model", "/tmp/ltv_rf.model"
)
self.ltv_model = joblib.load("/tmp/ltv_rf.model")
self.minio.fget_object(
"models", "churn_lr.model", "/tmp/churn_lr.model"
)
self.churn_model = joblib.load("/tmp/churn_lr.model")
except:
# 如果模型不存在,训练新模型
self.train_models()
def train_models(self):
"""训练所有预测模型"""
# 1. 准备训练数据
df = self.prepare_training_data()
# 2. 训练LTV预测模型
self.train_ltv_model(df)
# 3. 训练流失预测模型
self.train_churn_model(df)
# 4. 训练商品关联规则模型
self.train_association_rules(df)
def prepare_training_data(self):
"""准备训练数据"""
# 从微信支付获取交易数据
transactions = self.wxpay.get_transactions(
start_date=(datetime.now() - timedelta(days=365)),
end_date=datetime.now()
)
# 转换为Spark DataFrame
df = self.spark.createDataFrame(transactions)
# 数据预处理
df = self.preprocess_data(df)
return df
def preprocess_data(self, df):
"""数据预处理"""
# 处理缺失值
df = df.na.fill({
'amount': 0,
'user_age': df.agg({'user_age': 'mean'}).collect()[0][0]
})
# 特征工程
df = df.withColumn(
'purchase_freq',
df['transaction_count'] / df['days_since_first_purchase']
)
return df
def train_ltv_model(self, df):
"""训练LTV预测模型"""
# 特征选择
feature_cols = [
'user_age', 'transaction_count', 'avg_amount',
'purchase_freq', 'days_since_first_purchase'
]
# 数据转换
assembler = VectorAssembler(
inputCols=feature_cols,
outputCol="features"
)
# 随机森林回归
rf = RandomForestRegressor(
featuresCol="features",
labelCol="ltv_actual",
numTrees=100,
maxDepth=5
)
# 构建Pipeline
pipeline = Pipeline(stages=[assembler, rf])
# 训练模型
model = pipeline.fit(df)
# 保存模型
model_path = "/tmp/ltv_rf.model"
model.save(model_path)
self.minio.fput_object(
"models", "ltv_rf.model", model_path
)
self.ltv_model = model
def train_churn_model(self, df):
"""训练流失预测模型"""
# 特征选择
feature_cols = [
'inactivity_days', 'last_purchase_amount',
'purchase_freq_decline', 'service_complaints'
]
# 数据转换
assembler = VectorAssembler(
inputCols=feature_cols,
outputCol="features"
)
# 逻辑回归分类
from pyspark.ml.classification import LogisticRegression
lr = LogisticRegression(
featuresCol="features",
labelCol="churned"
)
# 构建Pipeline
pipeline = Pipeline(stages=[assembler, lr])
# 训练模型
model = pipeline.fit(df)
# 保存模型
model_path = "/tmp/churn_lr.model"
model.save(model_path)
self.minio.fput_object(
"models", "churn_lr.model", model_path
)
self.churn_model = model
def train_association_rules(self, df):
"""训练商品关联规则模型"""
# 转换数据格式
basket_df = df.groupBy("user_id") \
.agg(collect_list("product_id").alias("items"))
# FP-Growth算法
fp_growth = FPGrowth(
itemsCol="items",
minSupport=0.05,
minConfidence=0.3
)
# 训练模型
model = fp_growth.fit(basket_df)
# 保存模型
model_path = "/tmp/fpgrowth.model"
model.save(model_path)
self.minio.fput_object(
"models", "fpgrowth.model", model_path
)
self.fpgrowth_model = model
def predict_ltv(self, user_data):
"""预测用户生命周期价值"""
# 准备特征数据
features = self.prepare_features(user_data)
# 预测
prediction = self.ltv_model.transform(features)
return prediction.collect()[0]['prediction']
def predict_churn(self, user_data):
"""预测用户流失概率"""
# 准备特征数据
features = self.prepare_features(user_data)
# 预测
prediction = self.churn_model.transform(features)
return prediction.collect()[0]['probability']
def get_product_recommendations(self, product_ids):
"""获取商品关联推荐"""
# 查询关联规则
rules = self.fpgrowth_model.associationRules
# 筛选相关规则
recommendations = rules.filter(
array_contains(rules.antecedent, product_ids[0])
).orderBy("confidence", ascending=False)
return recommendations.limit(5).collect()
def setup_routes(self):
"""设置Flask路由"""
@self.app.route('/')
def dashboard():
return render_template('dashboard.html')
@self.app.route('/api/sales_trend')
def sales_trend():
# 获取销售趋势数据
data = self.get_sales_trend_data()
return json.dumps(data)
@self.socketio.on('realtime_update')
def handle_realtime_update():
# 实时推送数据更新
while True:
data = self.get_realtime_metrics()
self.socketio.emit('update', data)
time.sleep(5)
def get_sales_trend_data(self):
"""获取销售趋势数据"""
# 从微信支付API获取数据
transactions = self.wxpay.get_transactions(
start_date=datetime.now() - timedelta(days=30),
end_date=datetime.now()
)
# 按天聚合
df = pd.DataFrame(transactions)
df['date'] = pd.to_datetime(df['transaction_time']).dt.date
daily_sales = df.groupby('date')['amount'].sum()
# 转换为图表数据
return {
'dates': daily_sales.index.astype(str).tolist(),
'amounts': daily_sales.values.tolist()
}
def get_realtime_metrics(self):
"""获取实时指标"""
# 获取最近1小时的交易数据
transactions = self.wxpay.get_transactions(
start_date=datetime.now() - timedelta(hours=1),
end_date=datetime.now()
)
# 计算关键指标
df = pd.DataFrame(transactions)
metrics = {
'total_sales': df['amount'].sum(),
'transaction_count': len(df),
'avg_amount': df['amount'].mean(),
'top_products': df['product_id'].value_counts().head(3).to_dict()
}
return metrics
def run(self):
"""启动系统"""
self.socketio.run(self.app, host='0.0.0.0', port=5000)
if __name__ == "__main__":
bi_system = WeChatBusinessIntelligence()
bi_system.run()
使用说明
1. 环境准备
# 安装核心依赖
pip install pyspark flask flask-socketio plotly python-wechatpay minio joblib pandas
# 启动Spark集群 (需要预先安装Spark)
start-spark-cluster.sh
# 启动MinIO存储 (Docker方式)
docker run -p 9000:9000 -p 9001:9001 minio/minio server /data --console-address ":9001"
2. 配置文件
创建config/wechat_bi.json文件:
{
"wechat_pay": {
"mch_id": "your_merchant_id",
"api_key": "your_api_key",
"cert_path": "/path/to/cert.pem"
},
"minio": {
"endpoint": "localhost:9000",
"access_key": "minioadmin",
"secret_key": "minioadmin"
},
"spark": {
"master": "spark://localhost:7077",
"executor_memory": "8g"
}
}
3. 启动系统
python wechat_bi.py
访问 http://localhost:5000 查看商业智能仪表盘
功能扩展建议
- 供应链优化:集成库存预测与供应商管理系统
- 个性化推荐:基于用户画像的实时推荐引擎
- 欺诈检测:使用异常检测算法识别可疑交易
- 市场篮子分析:深入挖掘商品组合销售模式
- 预测性维护:针对智能硬件产品的故障预测
适用场景
✅ 微信电商平台经营分析
✅ 连锁零售业销售预测
✅ 会员制服务用户价值管理
✅ 营销活动效果实时监控
✅ 商品组合策略优化
这个系统比之前的解决方案更专注于商业智能领域,整合了:
- 微信支付数据深度分析
- Spark大数据处理
- 机器学习预测模型
- 实时数据可视化
- 商品关联规则挖掘
- 用户生命周期管理
适合需要基于微信生态进行深度商业分析的企业级应用场景。
更多推荐

所有评论(0)