功能介绍

这是一个面向微信生态的全链路自动化营销系统,具备以下高级功能:

  • 跨平台用户身份识别(公众号+小程序+企业微信)
  • 实时行为追踪与用户画像构建
  • 营销自动化决策树引擎
  • 多变量测试(MVT)框架
  • 智能内容生成与个性化推荐
  • 营销效果归因分析
  • 企业微信SCRM深度集成

系统架构

graph TB
    A[微信公众平台] --> B[用户行为采集]
    C[微信小程序] --> B
    D[企业微信] --> B
    B --> E[统一用户画像]
    E --> F[营销决策引擎]
    F --> G[内容生成系统]
    G --> H[多渠道分发]
    H -->|用户反馈| B
    F --> I[实验管理系统]
    I --> J[多变量测试]
    J --> F

核心代码实现

import json
from datetime import datetime, timedelta
from typing import Dict, List, Optional
import hashlib
import numpy as np
import pandas as pd
from pymongo import MongoClient
from redis import Redis
from celery import Celery
from wechatpy import WeChatClient, WeChatEnterprise
from wechatpy.oauth import WeChatOAuth
from minio import Minio
from sklearn.ensemble import RandomForestClassifier
from transformers import pipeline
import jieba.analyse
from flask import Flask, request, jsonify
import plotly.express as px
from alibabacloud_tea_openapi import models as open_api_models
from alibabacloud_dyvmsapi20211025.client import Client as DyvmsapiClient
from alibabacloud_dyvmsapi20211025 import models as dyvmsapi_models

class UnifiedUserProfile:
    """统一用户画像系统"""
    
    def __init__(self):
        self.mongo = MongoClient('mongodb://localhost:27017/')
        self.db = self.mongo['wechat_marketing']
        self.redis = Redis(host='localhost', port=6379, db=0)
        self.user_profiles = self.db.user_profiles
        self.behavior_logs = self.db.behavior_logs
        
    def identify_user(self, openid: str, unionid: str = None) -> str:
        """用户身份识别与合并"""
        if unionid:
            # 优先使用unionid作为唯一标识
            user_id = f"u_{unionid}"
            # 检查是否已有openid关联
            existing = self.user_profiles.find_one({'unionid': unionid})
            if existing and openid not in existing['openids']:
                self.user_profiles.update_one(
                    {'unionid': unionid},
                    {'$addToSet': {'openids': openid}}
                )
        else:
            # 仅使用openid
            user_id = f"o_{hashlib.md5(openid.encode()).hexdigest()}"
            
        # 初始化用户画像(如果不存在)
        if not self.user_profiles.find_one({'user_id': user_id}):
            self.user_profiles.insert_one({
                'user_id': user_id,
                'openids': [openid],
                'unionid': unionid,
                'profile': {},
                'created_at': datetime.now(),
                'updated_at': datetime.now()
            })
            
        return user_id
    
    def update_behavior(self, user_id: str, behavior: Dict):
        """记录用户行为"""
        # 记录原始行为日志
        log_entry = {
            'user_id': user_id,
            'behavior': behavior,
            'timestamp': datetime.now(),
            'processed': False
        }
        self.behavior_logs.insert_one(log_entry)
        
        # 实时更新Redis中的最新行为
        self.redis.set(f"user:recent:{user_id}", json.dumps(behavior))
        
        # 触发异步处理
        process_behavior.delay(user_id, behavior)
    
    def get_user_profile(self, user_id: str) -> Dict:
        """获取完整用户画像"""
        profile = self.user_profiles.find_one({'user_id': user_id})
        if not profile:
            return {}
            
        # 从Redis获取实时行为
        recent_behavior = self.redis.get(f"user:recent:{user_id}")
        if recent_behavior:
            profile['recent_behavior'] = json.loads(recent_behavior)
            
        return profile

class MarketingDecisionEngine:
    """营销决策引擎"""
    
    def __init__(self):
        self.client = WeChatClient('appid', 'secret')
        self.ent_client = WeChatEnterprise('corpid', 'corpsecret')
        self.nlp = pipeline("text-generation", model="uer/gpt2-chinese-cluecorpussmall")
        self.rf_model = self.load_model('rf_segment.model')
        self.minio = Minio(
            "minio.example.com",
            access_key="minioadmin",
            secret_key="minioadmin",
            secure=False
        )
        
    def load_model(self, model_name: str):
        """从MinIO加载模型"""
        try:
            self.minio.fget_object("models", model_name, f"/tmp/{model_name}")
            return joblib.load(f"/tmp/{model_name}")
        except Exception as e:
            print(f"Failed to load model {model_name}: {str(e)}")
            return None
    
    def make_decision(self, user_id: str, context: Dict = None) -> Dict:
        """制定营销决策"""
        # 获取用户画像
        profile = UnifiedUserProfile().get_user_profile(user_id)
        
        # 获取实验分组
        exp_group = self.get_experiment_group(user_id)
        
        # 决策树逻辑
        if profile.get('rfm_score', 0) > 4.5:
            # 高价值用户策略
            decision = self.handle_high_value_user(profile, exp_group)
        elif profile.get('churn_risk', 0) > 0.7:
            # 流失风险用户策略
            decision = self.handle_churn_risk_user(profile, exp_group)
        else:
            # 普通用户策略
            decision = self.handle_normal_user(profile, exp_group)
            
        # 添加实验标记
        decision['experiment_group'] = exp_group
        
        return decision
    
    def handle_high_value_user(self, profile: Dict, exp_group: str) -> Dict:
        """高价值用户处理逻辑"""
        # 个性化内容生成
        content = self.generate_personalized_content(
            profile, 
            template="high_value_template",
            exp_group=exp_group
        )
        
        return {
            'strategy': 'high_value',
            'content': content,
            'channels': ['wechat_mp', 'enterprise_wechat'],
            'timing': 'immediate',
            'incentive': {
                'type': 'coupon',
                'value': 50,
                'expiry': (datetime.now() + timedelta(days=7)).isoformat()
            }
        }
    
    def generate_personalized_content(self, profile: Dict, template: str, exp_group: str) -> str:
        """生成个性化营销内容"""
        # 从MinIO获取模板
        try:
            self.minio.fget_object("templates", f"{template}.json", "/tmp/template.json")
            with open("/tmp/template.json") as f:
                template_data = json.load(f)
        except:
            template_data = {"default": "尊敬的{name},我们为您准备了专属优惠!"}
            
        # 选择实验组特定模板
        template_text = template_data.get(exp_group, template_data['default'])
        
        # 提取用户关键词
        behaviors = [b['behavior']['type'] for b in 
                    self.db.behavior_logs.find({'user_id': profile['user_id']})]
        text = " ".join(behaviors)
        keywords = jieba.analyse.extract_tags(text, topK=3)
        
        # 生成个性化内容
        prompt = f"根据以下信息生成营销内容:\n用户标签: {profile.get('tags', [])}\n" \
                f"近期行为: {keywords}\n模板: {template_text}\n内容:"
                
        generated = self.nlp(prompt, max_length=100, num_return_sequences=1)
        return generated[0]['generated_text']

class ExperimentManager:
    """实验管理系统"""
    
    def __init__(self):
        self.experiments = {}
        self.load_experiments()
        
    def load_experiments(self):
        """加载实验配置"""
        # 这里应该是从数据库加载,简化为示例
        self.experiments = {
            'high_value_strategy': {
                'groups': {
                    'A': {'weight': 0.3, 'params': {'template': 'A'}},
                    'B': {'weight': 0.3, 'params': {'template': 'B'}},
                    'control': {'weight': 0.4, 'params': {}}
                },
                'metrics': ['conversion_rate', 'revenue_per_user']
            }
        }
    
    def assign_group(self, user_id: str, experiment_name: str) -> str:
        """分配实验分组"""
        if experiment_name not in self.experiments:
            return 'default'
            
        # 使用用户ID哈希确保一致性分配
        hash_val = int(hashlib.md5(user_id.encode()).hexdigest()[:8], 16)
        total = sum(g['weight'] for g in self.experiments[experiment_name]['groups'].values())
        point = (hash_val % 10000) / 10000 * total
        
        cumulative = 0
        for group_name, group in self.experiments[experiment_name]['groups'].items():
            cumulative += group['weight']
            if point <= cumulative:
                return group_name
                
        return 'control'

# Celery任务队列
celery = Celery('wechat_marketing', broker='redis://localhost:6379/1')

@celery.task
def process_behavior(user_id: str, behavior: Dict):
    """异步处理用户行为"""
    # 1. 更新用户RFM模型
    update_rfm_score(user_id, behavior)
    
    # 2. 更新流失风险预测
    update_churn_risk(user_id)
    
    # 3. 触发实时营销决策
    decision = MarketingDecisionEngine().make_decision(user_id)
    execute_decision.delay(user_id, decision)

@celery.task
def execute_decision(user_id: str, decision: Dict):
    """执行营销决策"""
    # 根据决策选择发送渠道
    if 'wechat_mp' in decision['channels']:
        send_wechat_message(user_id, decision['content'])
        
    if 'enterprise_wechat' in decision['channels']:
        send_enterprise_wechat_message(user_id, decision['content'])
        
    if 'sms' in decision['channels']:
        send_sms_message(user_id, decision['content'])

# Flask API服务
app = Flask(__name__)

@app.route('/track', methods=['POST'])
def track_behavior():
    """用户行为追踪接口"""
    data = request.json
    user_id = UnifiedUserProfile().identify_user(
        data.get('openid'),
        data.get('unionid')
    )
    UnifiedUserProfile().update_behavior(user_id, data['behavior'])
    return jsonify({'status': 'success'})

@app.route('/decision', methods=['GET'])
def get_decision():
    """获取营销决策接口"""
    user_id = request.args.get('user_id')
    decision = MarketingDecisionEngine().make_decision(user_id)
    return jsonify(decision)

if __name__ == "__main__":
    # 启动Flask应用
    app.run(host='0.0.0.0', port=5000, debug=True)

使用说明

1. 环境准备

# 安装核心依赖
pip install wechatpy pymongo redis celery minio jieba transformers scikit-learn plotly flask

# 启动Redis服务
docker run -d -p 6379:6379 redis

# 启动MongoDB服务
docker run -d -p 27017:27017 mongo

# 启动Celery worker
celery -A wechat_marketing worker --loglevel=info

2. 配置系统

创建config/marketing_config.json文件:

{
  "wechat": {
    "appid": "your_appid",
    "secret": "your_secret",
    "corpid": "your_corpid",
    "corpsecret": "your_corpsecret"
  },
  "minio": {
    "endpoint": "localhost:9000",
    "access_key": "minioadmin",
    "secret_key": "minioadmin"
  },
  "experiments": {
    "high_value_strategy": {
      "groups": {
        "A": {"weight": 0.4, "params": {}},
        "B": {"weight": 0.4, "params": {}},
        "control": {"weight": 0.2, "params": {}}
      }
    }
  }
}

3. 启动系统

python wechat_marketing_system.py

功能扩展建议

  1. ​跨渠道归因分析​​:整合微信生态外数据源
  2. ​强化学习优化​​:动态调整营销策略
  3. ​语音交互支持​​:集成微信语音消息处理
  4. ​社交裂变追踪​​:监控分享传播路径
  5. ​预测库存管理​​:结合销售预测优化库存

适用场景

✅ 电商平台精准营销
✅ 会员生命周期管理
✅ 线下零售数字化运营
✅ 金融服务客户维系
✅ 教育行业学员转化

这个系统比之前的解决方案更专注于营销自动化,整合了:

  • 跨平台用户身份识别
  • 实时行为分析与画像构建
  • 智能决策引擎
  • 多变量测试框架
  • 个性化内容生成
  • 多渠道执行能力

适合需要在微信生态内实现智能化营销的企业级应用场景。

更多推荐