功能介绍

这是一个基于机器学习和行为分析的微信生态风控系统,具备以下高级功能:

  • ​实时交易风险检测​​:基于深度学习的交易欺诈识别
  • ​多维度用户行为分析​​:设备指纹、社交网络、交易模式分析
  • ​动态风险评估引擎​​:实时计算用户风险分数
  • ​智能规则引擎​​:支持复杂业务规则的风控策略
  • ​图神经网络检测​​:识别团伙欺诈和洗钱行为
  • ​实时预警与拦截​​:毫秒级风险响应
  • ​风控数据可视化​​:全方位风险监控仪表盘

系统架构

graph TB
    A[微信支付交易] --> B[数据采集层]
    C[用户行为数据] --> B
    D[设备指纹数据] --> B
    B --> E[实时流处理]
    E --> F[特征工程]
    F --> G[机器学习模型]
    G --> H[风险决策引擎]
    H --> I[实时预警]
    H --> J[交易拦截]
    I --> K[风控仪表盘]
    J --> L[人工审核队列]

核心代码实现

import json
import time
import hashlib
import numpy as np
import pandas as pd
from datetime import datetime, timedelta
from typing import Dict, List, Optional, Tuple
import torch
import torch.nn as nn
from torch_geometric.data import Data
from torch_geometric.nn import GCNConv
from flask import Flask, request, jsonify
from flask_socketio import SocketIO
from kafka import KafkaProducer, KafkaConsumer
from pymongo import MongoClient
from redis import Redis
from celery import Celery
import xgboost as xgb
from sklearn.ensemble import IsolationForest
from sklearn.preprocessing import StandardScaler
import plotly.express as px
import plotly.graph_objects as go
from plotly.subplots import make_subplots
from wechatpy import WeChatPay
import tensorflow as tf
from tensorflow import keras
from prometheus_client import Counter, Gauge, Histogram
import warnings
warnings.filterwarnings('ignore')

# 初始化应用
app = Flask(__name__)
socketio = SocketIO(app, cors_allowed_origins="*")

# 初始化数据库和消息队列
mongo_client = MongoClient('mongodb://localhost:27017/')
redis_client = Redis(host='localhost', port=6379, db=0)
celery = Celery('wechat_risk', broker='redis://localhost:6379/1')

# Kafka配置
kafka_producer = KafkaProducer(
    bootstrap_servers=['localhost:9092'],
    value_serializer=lambda v: json.dumps(v).encode('utf-8')
)

# 微信支付客户端
wx_pay = WeChatPay(
    mch_id='your_mch_id',
    api_key='your_api_key',
    cert_path='/path/to/cert.pem'
)

# Prometheus监控指标
RISK_SCORE = Gauge('risk_score', 'Current risk score')
FRAUD_COUNT = Counter('fraud_detected', 'Number of fraud cases detected')
PROCESSING_TIME = Histogram('processing_time_seconds', 'Request processing time')

class DeviceFingerprint:
    """设备指纹识别与验证"""
    
    def __init__(self):
        self.device_profiles = {}
        
    def generate_fingerprint(self, request_data: Dict) -> str:
        """生成设备指纹"""
        device_info = {
            'user_agent': request_data.get('User-Agent', ''),
            'screen_resolution': request_data.get('screen_resolution', ''),
            'timezone': request_data.get('timezone', ''),
            'language': request_data.get('language', ''),
            'plugins': request_data.get('plugins', []),
            'fonts': request_data.get('fonts', [])
        }
        
        fingerprint_str = json.dumps(device_info, sort_keys=True)
        return hashlib.sha256(fingerprint_str.encode()).hexdigest()
    
    def check_device_anomaly(self, fingerprint: str, user_id: str) -> Dict:
        """检查设备异常"""
        # 检查设备是否首次出现
        if fingerprint not in self.device_profiles:
            self.device_profiles[fingerprint] = {
                'first_seen': datetime.now(),
                'users': set(),
                'locations': set()
            }
            return {'risk_level': 'low', 'reason': 'new_device'}
        
        profile = self.device_profiles[fingerprint]
        
        # 检查设备是否关联过多用户
        if len(profile['users']) > 5 and user_id not in profile['users']:
            return {'risk_level': 'high', 'reason': 'device_sharing'}
        
        # 检查地理位置异常
        current_location = request_data.get('location', '')
        if current_location and len(profile['locations']) > 3:
            if current_location not in profile['locations']:
                return {'risk_level': 'medium', 'reason': 'location_hopping'}
        
        return {'risk_level': 'low', 'reason': 'normal'}

class FeatureEngineer:
    """特征工程处理器"""
    
    def __init__(self):
        self.scaler = StandardScaler()
        self.feature_stats = {}
        
    def extract_transaction_features(self, transaction: Dict) -> List[float]:
        """提取交易特征"""
        features = []
        
        # 金额相关特征
        amount = transaction.get('amount', 0)
        features.extend([
            amount,
            np.log1p(amount),
            amount / (transaction.get('user_avg_amount', 1) or 1)
        ])
        
        # 时间相关特征
        transaction_time = datetime.fromtimestamp(transaction['timestamp'])
        features.extend([
            transaction_time.hour,
            transaction_time.weekday(),
            transaction_time.hour in [0, 1, 2, 3, 4, 5]  # 是否夜间
        ])
        
        # 用户行为特征
        features.extend([
            transaction.get('user_transaction_count', 0),
            transaction.get('user_success_rate', 0),
            transaction.get('user_failed_attempts', 0)
        ])
        
        # 设备特征
        features.extend([
            transaction.get('device_risk_score', 0),
            transaction.get('ip_risk_score', 0)
        ])
        
        return features
    
    def extract_graph_features(self, user_id: str, transaction: Dict) -> List[float]:
        """提取图网络特征"""
        # 从图数据库获取用户关系特征
        user_network = self.get_user_network(user_id)
        
        features = [
            len(user_network.get('direct_connections', [])),
            len(user_network.get('second_degree_connections', [])),
            user_network.get('cluster_coefficient', 0),
            user_network.get('betweenness_centrality', 0)
        ]
        
        return features
    
    def get_user_network(self, user_id: str) -> Dict:
        """获取用户社交网络信息"""
        # 这里应该是从图数据库查询,简化为示例
        return {
            'direct_connections': [],
            'second_degree_connections': [],
            'cluster_coefficient': 0.5,
            'betweenness_centrality': 0.1
        }

class DeepLearningModel(nn.Module):
    """深度学习风险预测模型"""
    
    def __init__(self, input_dim: int, hidden_dims: List[int] = [64, 32, 16]):
        super(DeepLearningModel, self).__init__()
        layers = []
        prev_dim = input_dim
        
        for hidden_dim in hidden_dims:
            layers.append(nn.Linear(prev_dim, hidden_dim))
            layers.append(nn.ReLU())
            layers.append(nn.Dropout(0.2))
            prev_dim = hidden_dim
        
        layers.append(nn.Linear(prev_dim, 1))
        layers.append(nn.Sigmoid())
        
        self.network = nn.Sequential(*layers)
    
    def forward(self, x: torch.Tensor) -> torch.Tensor:
        return self.network(x)

class GNNModel(nn.Module):
    """图神经网络欺诈检测模型"""
    
    def __init__(self, node_dim: int, hidden_dim: int = 64):
        super(GNNModel, self).__init__()
        self.conv1 = GCNConv(node_dim, hidden_dim)
        self.conv2 = GCNConv(hidden_dim, hidden_dim)
        self.classifier = nn.Linear(hidden_dim, 1)
    
    def forward(self, data: Data) -> torch.Tensor:
        x, edge_index = data.x, data.edge_index
        x = self.conv1(x, edge_index).relu()
        x = self.conv2(x, edge_index).relu()
        x = self.classifier(x)
        return torch.sigmoid(x)

class RiskEngine:
    """风险决策引擎"""
    
    def __init__(self):
        self.device_fp = DeviceFingerprint()
        self.feature_engineer = FeatureEngineer()
        self.load_models()
        self.rules_engine = self.setup_rules_engine()
        
    def load_models(self):
        """加载所有风险模型"""
        # 深度学习模型
        self.dl_model = DeepLearningModel(input_dim=50)
        try:
            self.dl_model.load_state_dict(torch.load('/models/dl_risk_model.pth'))
        except:
            print("DL model not found, using default")
        
        # XGBoost模型
        self.xgb_model = xgb.Booster()
        try:
            self.xgb_model.load_model('/models/xgb_risk_model.json')
        except:
            print("XGBoost model not found")
        
        # 孤立森林异常检测
        self.isolation_forest = IsolationForest(contamination=0.1)
        
        # 图神经网络
        self.gnn_model = GNNModel(node_dim=10)
        
    def setup_rules_engine(self) -> Dict:
        """设置规则引擎"""
        return {
            'high_risk_rules': [
                {'condition': 'amount > 10000 and hour in [0,6]', 'score': 0.8},
                {'condition': 'new_device and amount > 5000', 'score': 0.7},
                {'condition': 'ip_country != user_country', 'score': 0.6}
            ],
            'medium_risk_rules': [
                {'condition': 'amount > user_avg_amount * 3', 'score': 0.5},
                {'condition': 'transaction_count > 10 in 1h', 'score': 0.4}
            ]
        }
    
    @PROCESSING_TIME.time()
    def assess_risk(self, transaction: Dict) -> Dict:
        """综合风险评估"""
        start_time = time.time()
        
        # 1. 设备指纹验证
        device_risk = self.device_fp.check_device_anomaly(
            transaction.get('device_fingerprint', ''),
            transaction.get('user_id', '')
        )
        
        # 2. 特征工程
        features = self.feature_engineer.extract_transaction_features(transaction)
        graph_features = self.feature_engineer.extract_graph_features(
            transaction.get('user_id', ''), transaction
        )
        all_features = features + graph_features
        
        # 3. 模型预测
        dl_score = self.predict_dl_risk(all_features)
        xgb_score = self.predict_xgb_risk(all_features)
        anomaly_score = self.detect_anomaly(all_features)
        
        # 4. 规则引擎
        rule_score = self.apply_rules(transaction)
        
        # 5. 综合评分
        final_score = self.compute_final_score({
            'dl': dl_score,
            'xgb': xgb_score,
            'anomaly': anomaly_score,
            'rule': rule_score,
            'device': self.convert_risk_level(device_risk['risk_level'])
        })
        
        # 6. 决策
        decision = self.make_decision(final_score, transaction)
        
        processing_time = time.time() - start_time
        RISK_SCORE.set(final_score)
        
        if decision['action'] == 'block':
            FRAUD_COUNT.inc()
        
        return {
            'risk_score': final_score,
            'decision': decision,
            'processing_time': processing_time,
            'model_scores': {
                'dl': float(dl_score),
                'xgb': float(xgb_score),
                'anomaly': float(anomaly_score),
                'rule': float(rule_score)
            },
            'device_risk': device_risk
        }
    
    def predict_dl_risk(self, features: List[float]) -> float:
        """深度学习模型预测"""
        features_tensor = torch.FloatTensor(features).unsqueeze(0)
        with torch.no_grad():
            score = self.dl_model(features_tensor).item()
        return score
    
    def predict_xgb_risk(self, features: List[float]) -> float:
        """XGBoost模型预测"""
        dmatrix = xgb.DMatrix([features])
        score = self.xgb_model.predict(dmatrix)[0]
        return score
    
    def detect_anomaly(self, features: List[float]) -> float:
        """异常检测"""
        # 这里应该是使用训练好的孤立森林模型
        # 简化为示例
        return 0.0
    
    def apply_rules(self, transaction: Dict) -> float:
        """应用业务规则"""
        total_score = 0.0
        rule_count = 0
        
        for rule in self.rules_engine['high_risk_rules'] + self.rules_engine['medium_risk_rules']:
            if self.evaluate_condition(rule['condition'], transaction):
                total_score += rule['score']
                rule_count += 1
        
        return total_score / max(rule_count, 1)
    
    def evaluate_condition(self, condition: str, data: Dict) -> bool:
        """评估规则条件"""
        # 这里应该是完整的条件解析器,简化为示例
        try:
            return eval(condition, {}, data)
        except:
            return False
    
    def compute_final_score(self, scores: Dict) -> float:
        """计算最终风险分数"""
        weights = {'dl': 0.3, 'xgb': 0.3, 'anomaly': 0.2, 'rule': 0.1, 'device': 0.1}
        final_score = 0.0
        
        for model, score in scores.items():
            final_score += score * weights.get(model, 0)
        
        return min(max(final_score, 0), 1)
    
    def convert_risk_level(self, level: str) -> float:
        """转换风险等级到分数"""
        risk_map = {'high': 0.8, 'medium': 0.5, 'low': 0.1}
        return risk_map.get(level, 0.1)
    
    def make_decision(self, risk_score: float, transaction: Dict) -> Dict:
        """根据风险分数做出决策"""
        amount = transaction.get('amount', 0)
        
        if risk_score > 0.8:
            return {'action': 'block', 'reason': 'high_risk'}
        elif risk_score > 0.6:
            return {'action': 'review', 'reason': 'suspicious'}
        elif risk_score > 0.4 and amount > 5000:
            return {'action': 'verify', 'reason': 'large_amount'}
        else:
            return {'action': 'approve', 'reason': 'low_risk'}

class RiskDashboard:
    """风险监控仪表盘"""
    
    def __init__(self):
        self.db = mongo_client['risk_monitoring']
        
    def get_realtime_metrics(self) -> Dict:
        """获取实时监控指标"""
        # 最近1小时数据
        one_hour_ago = datetime.now() - timedelta(hours=1)
        
        metrics = {
            'total_transactions': self.db.transactions.count_documents({}),
            'blocked_transactions': self.db.transactions.count_documents({'decision.action': 'block'}),
            'avg_risk_score': self.db.transactions.aggregate([
                {'$match': {'timestamp': {'$gte': one_hour_ago}}},
                {'$group': {'_id': None, 'avg_score': {'$avg': '$risk_score'}}}
            ]).next().get('avg_score', 0),
            'top_risky_users': list(self.db.transactions.aggregate([
                {'$match': {'risk_score': {'$gt': 0.7}}},
                {'$group': {'_id': '$user_id', 'count': {'$sum': 1}}},
                {'$sort': {'count': -1}},
                {'$limit': 5}
            ]))
        }
        
        return metrics
    
    def create_dashboard(self) -> Dict:
        """创建监控仪表盘"""
        metrics = self.get_realtime_metrics()
        
        # 风险分数分布图
        risk_scores = list(self.db.transactions.find(
            {}, {'risk_score': 1, '_id': 0}
        ).limit(1000))
        df = pd.DataFrame(risk_scores)
        
        fig1 = px.histogram(df, x='risk_score', title='风险分数分布')
        
        # 实时风险趋势
        hourly_data = list(self.db.transactions.aggregate([
            {'$group': {
                '_id': {'$hour': '$timestamp'},
                'avg_risk': {'$avg': '$risk_score'},
                'count': {'$sum': 1}
            }},
            {'$sort': {'_id': 1}}
        ]))
        
        fig2 = make_subplots(specs=[[{"secondary_y": True}]])
        fig2.add_trace(
            go.Scatter(x=[d['_id'] for d in hourly_data], 
                      y=[d['avg_risk'] for d in hourly_data], 
                      name="平均风险分数"),
            secondary_y=False
        )
        fig2.add_trace(
            go.Bar(x=[d['_id'] for d in hourly_data], 
                  y=[d['count'] for d in hourly_data], 
                  name="交易数量", opacity=0.3),
            secondary_y=True
        )
        fig2.update_layout(title='每小时风险趋势')
        
        return {
            'metrics': metrics,
            'charts': {
                'risk_distribution': fig1.to_html(),
                'hourly_trend': fig2.to_html()
            }
        }

# 初始化风险引擎
risk_engine = RiskEngine()
dashboard = RiskDashboard()

# API路由
@app.route('/risk/assess', methods=['POST'])
def assess_risk():
    """风险评估接口"""
    data = request.json
    result = risk_engine.assess_risk(data)
    
    # 保存评估结果
    mongo_client.risk_monitoring.transactions.insert_one({
        **data,
        **result,
        'timestamp': datetime.now()
    })
    
    # 发送到Kafka
    kafka_producer.send('risk-assessments', result)
    
    return jsonify(result)

@app.route('/dashboard')
def risk_dashboard():
    """风险监控仪表盘"""
    dashboard_data = dashboard.create_dashboard()
    return jsonify(dashboard_data)

@socketio.on('connect')
def handle_connect():
    """实时数据推送"""
    while True:
        dashboard_data = dashboard.create_dashboard()
        socketio.emit('dashboard_update', dashboard_data)
        time.sleep(5)

# Celery任务
@celery.task
def process_batch_risk_assessment(batch_data: List[Dict]):
    """批量风险评估"""
    results = []
    for transaction in batch_data:
        result = risk_engine.assess_risk(transaction)
        results.append(result)
    
    # 保存到数据库
    mongo_client.risk_monitoring.batch_assessments.insert_many(results)
    
    return results

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

使用说明

1. 环境准备

# 安装核心依赖
pip install torch torch-geometric xgboost scikit-learn flask flask-socketio kafka-python pymongo redis celery plotly prometheus-client wechatpy tensorflow

# 启动依赖服务
docker run -d -p 27017:27017 mongo
docker run -d -p 6379:6379 redis
docker run -d -p 9092:9092 apache/kafka:latest
docker run -d -p 9090:9090 prom/prometheus

2. 模型训练与部署

  1. ​准备训练数据​​:收集历史交易数据作为训练集
  2. ​训练深度学习模型​​:
python train_dl_model.py --data-path /data/transactions.csv
  1. ​训练XGBoost模型​​:
python train_xgb_model.py --data-path /data/transactions.csv
  1. ​部署模型​​:将训练好的模型文件放到/models/目录

3. 配置系统

创建config/risk_config.json文件:

{
  "kafka": {
    "bootstrap_servers": ["localhost:9092"],
    "topic": "risk-assessments"
  },
  "wechat_pay": {
    "mch_id": "your_merchant_id",
    "api_key": "your_api_key"
  },
  "risk_thresholds": {
    "block": 0.8,
    "review": 0.6,
    "verify": 0.4
  }
}

4. 启动系统

# 启动风险引擎
python wechat_risk_engine.py

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

# 访问监控仪表盘
# http://localhost:5000/dashboard

功能扩展建议

  1. ​联邦学习​​:在保护隐私的前提下联合多家机构训练模型
  2. ​强化学习​​:动态调整风险阈值和策略
  3. ​区块链存证​​:将高风险交易上链存证
  4. ​多方安全计算​​:保护数据隐私的风险评估
  5. ​边缘计算​​:在用户设备上进行初步风险检测

适用场景

✅ 微信支付交易风控
✅ 互联网金融反欺诈
✅ 电商平台交易安全
✅ 数字货币交易监控
✅ 跨境支付风险控制

这个系统整合了:

  • 多模态机器学习模型
  • 实时流处理
  • 图神经网络
  • 复杂规则引擎
  • 实时监控告警
  • 可视化分析

适合需要构建企业级微信支付风控系统的场景。

更多推荐