功能介绍

这是一个基于微信生态的智能健康管理平台,具备以下高级功能:

  • ​多模态健康数据采集​​:穿戴设备、手动录入、AI识别
  • ​智能健康风险评估​​:机器学习疾病预测模型
  • ​远程诊疗集成​​:视频问诊、电子处方、药品配送
  • ​健康知识图谱​​:症状-疾病-治疗方案关联分析
  • ​用药智能提醒​​:基于NLP的用药指导与提醒
  • ​流行病学监测​​:区域健康数据实时监控
  • ​隐私安全保护​​:医疗数据加密与联邦学习

系统架构

graph TB
    A[微信小程序] --> B[健康数据采集]
    C[穿戴设备] --> B
    D[手动录入] --> B
    B --> E[数据预处理]
    E --> F[特征工程]
    F --> G[健康风险评估]
    G --> H[远程诊疗引擎]
    H --> I[医生端微信]
    H --> J[药品配送系统]
    G --> K[健康预警]
    K --> L[用户通知]
    M[医疗知识图谱] --> G

核心代码实现

import os
import json
import time
import hashlib
import asyncio
from datetime import datetime, timedelta
from typing import Dict, List, Optional, Tuple
import numpy as np
import pandas as pd
import cv2
from PIL import Image, ImageDraw, ImageFont
import torch
import torch.nn as nn
from transformers import (
    BertTokenizer, BertModel,
    CLIPProcessor, CLIPModel,
    pipeline
)
import torchvision.transforms as transforms
from flask import Flask, request, jsonify
from flask_socketio import SocketIO
from celery import Celery
from pymongo import MongoClient
from redis import Redis
from minio import Minio
import plotly.express as px
import plotly.graph_objects as go
from plotly.subplots import make_subplots
from wechatpy import WeChatClient
import jwt
from cryptography.fernet import Fernet
import mysql.connector
from sqlalchemy import create_engine, Column, String, Float, DateTime
from sqlalchemy.ext.declarative import declarative_base
from sqlalchemy.orm import sessionmaker
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('healthcare', broker='redis://localhost:6379/1')

# 初始化微信客户端
wx_client = WeChatClient('your_appid', 'your_secret')

# 加密配置
encryption_key = Fernet.generate_key()
cipher_suite = Fernet(encryption_key)

# 医疗数据库连接
medical_engine = create_engine('mysql+mysqlconnector://user:password@localhost/medical_db')
SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=medical_engine)
Base = declarative_base()

class MedicalKnowledgeGraph:
    """医疗知识图谱管理系统"""
    
    def __init__(self):
        self.db = mongo_client['medical_kg']
        self.symptom_disease_map = {}
        self.disease_treatment_map = {}
        self.load_knowledge_graph()
    
    def load_knowledge_graph(self):
        """加载医疗知识图谱"""
        # 从数据库加载知识图谱数据
        try:
            symptoms = list(self.db.symptoms.find({}))
            diseases = list(self.db.diseases.find({}))
            treatments = list(self.db.treatments.find({}))
            
            # 构建症状-疾病映射
            for disease in diseases:
                for symptom in disease.get('symptoms', []):
                    if symptom not in self.symptom_disease_map:
                        self.symptom_disease_map[symptom] = []
                    self.symptom_disease_map[symptom].append({
                        'disease': disease['name'],
                        'probability': disease.get('probability', 0.5)
                    })
            
            # 构建疾病-治疗方案映射
            for treatment in treatments:
                for disease in treatment.get('diseases', []):
                    if disease not in self.disease_treatment_map:
                        self.disease_treatment_map[disease] = []
                    self.disease_treatment_map[disease].append(treatment)
                    
        except Exception as e:
            print(f"加载知识图谱失败: {str(e)}")
    
    def diagnose_symptoms(self, symptoms: List[str], user_data: Dict) -> List[Dict]:
        """基于症状进行初步诊断"""
        possible_diseases = []
        
        for symptom in symptoms:
            if symptom in self.symptom_disease_map:
                for disease_info in self.symptom_disease_map[symptom]:
                    # 考虑用户年龄、性别等因素调整概率
                    adjusted_prob = self.adjust_probability(
                        disease_info['probability'], user_data
                    )
                    possible_diseases.append({
                        'disease': disease_info['disease'],
                        'probability': adjusted_prob,
                        'matching_symptoms': [symptom]
                    })
        
        # 合并相同疾病
        merged_diseases = {}
        for disease in possible_diseases:
            name = disease['disease']
            if name in merged_diseases:
                merged_diseases[name]['probability'] *= 1.2  # 多个症状增加概率
                merged_diseases[name]['matching_symptoms'].extend(disease['matching_symptoms'])
            else:
                merged_diseases[name] = disease
        
        # 按概率排序
        return sorted(merged_diseases.values(), key=lambda x: x['probability'], reverse=True)[:5]
    
    def adjust_probability(self, base_prob: float, user_data: Dict) -> float:
        """根据用户数据调整疾病概率"""
        adjusted_prob = base_prob
        
        # 年龄因素
        age = user_data.get('age', 30)
        if age > 60:
            adjusted_prob *= 1.3
        elif age < 18:
            adjusted_prob *= 0.8
        
        # 性别因素
        gender = user_data.get('gender', 'unknown')
        # 这里可以根据具体疾病调整
        
        return min(adjusted_prob, 0.95)  # 最大概率限制

class HealthDataProcessor:
    """健康数据处理与分析引擎"""
    
    def __init__(self):
        self.models = self.load_ai_models()
        self.feature_scalers = {}
        self.risk_models = {}
        self.setup_risk_models()
    
    def load_ai_models(self) -> Dict:
        """加载AI模型"""
        return {
            'bert_tokenizer': BertTokenizer.from_pretrained('bert-base-chinese'),
            'bert_model': BertModel.from_pretrained('bert-base-chinese'),
            'clip_processor': CLIPProcessor.from_pretrained('openai/clip-vit-base-patch32'),
            'clip_model': CLIPModel.from_pretrained('openai/clip-vit-base-patch32'),
            'symptom_analyzer': pipeline("text-classification", 
                                       model="medical-bert-symptom")
        }
    
    def setup_risk_models(self):
        """设置风险评估模型"""
        # 这里应该是训练好的机器学习模型
        # 简化为示例
        self.risk_models = {
            'diabetes': {'threshold': 0.7, 'features': ['glucose', 'bmi', 'age']},
            'hypertension': {'threshold': 0.6, 'features': ['bp_systolic', 'bp_diastolic', 'age']},
            'heart_disease': {'threshold': 0.65, 'features': ['cholesterol', 'age', 'smoking']}
        }
    
    async def process_health_data(self, health_data: Dict) -> Dict:
        """处理健康数据"""
        processed_data = {}
        
        # 处理数值型数据
        if 'vitals' in health_data:
            processed_data['vitals'] = self.process_vitals(health_data['vitals'])
        
        # 处理症状描述
        if 'symptoms' in health_data:
            processed_data['symptoms'] = await self.analyze_symptoms(health_data['symptoms'])
        
        # 处理医疗影像
        if 'medical_images' in health_data:
            processed_data['image_analysis'] = await self.analyze_medical_images(
                health_data['medical_images']
            )
        
        # 风险评估
        processed_data['risk_assessment'] = self.assess_health_risks(health_data)
        
        return processed_data
    
    def process_vitals(self, vitals: Dict) -> Dict:
        """处理生命体征数据"""
        processed = {}
        
        # 血压分析
        if 'bp_systolic' in vitals and 'bp_diastolic' in vitals:
            systolic = vitals['bp_systolic']
            diastolic = vitals['bp_diastolic']
            processed['bp_category'] = self.categorize_blood_pressure(systolic, diastolic)
        
        # 血糖分析
        if 'glucose' in vitals:
            glucose = vitals['glucose']
            processed['glucose_category'] = self.categorize_glucose(glucose)
        
        return processed
    
    def categorize_blood_pressure(self, systolic: float, diastolic: float) -> str:
        """血压分类"""
        if systolic < 120 and diastolic < 80:
            return 'normal'
        elif systolic < 130 and diastolic < 85:
            return 'elevated'
        elif systolic < 140 or diastolic < 90:
            return 'stage1_hypertension'
        else:
            return 'stage2_hypertension'
    
    async def analyze_symptoms(self, symptoms_text: str) -> Dict:
        """分析症状描述"""
        # 使用医疗BERT模型分析症状
        analysis = self.models['symptom_analyzer'](symptoms_text[:512])
        
        return {
            'main_symptoms': [analysis[0]['label']],
            'confidence': analysis[0]['score'],
            'suggested_specialties': self.suggest_specialties(analysis[0]['label'])
        }
    
    def suggest_specialties(self, symptom: str) -> List[str]:
        """根据症状推荐科室"""
        specialty_map = {
            'headache': ['neurology', 'general'],
            'fever': ['infectious_disease', 'general'],
            'chest_pain': ['cardiology', 'emergency'],
            'abdominal_pain': ['gastroenterology', 'general']
        }
        return specialty_map.get(symptom, ['general'])
    
    def assess_health_risks(self, health_data: Dict) -> Dict:
        """评估健康风险"""
        risks = {}
        
        for disease, model_info in self.risk_models.items():
            # 提取特征
            features = []
            for feature_name in model_info['features']:
                if feature_name in health_data:
                    features.append(health_data[feature_name])
                else:
                    features.append(0)  # 缺失值处理
            
            # 简化的风险评估(实际应该使用训练好的模型)
            risk_score = sum(features) / len(features) if features else 0
            risks[disease] = {
                'score': risk_score,
                'level': 'high' if risk_score > model_info['threshold'] else 'low'
            }
        
        return risks

class RemoteConsultationSystem:
    """远程诊疗系统"""
    
    def __init__(self):
        self.db = mongo_client['telemedicine']
        self.doctor_pool = self.load_doctor_pool()
        self.consultation_sessions = {}
    
    def load_doctor_pool(self) -> Dict:
        """加载医生资源池"""
        doctors = list(self.db.doctors.find({'status': 'available'}))
        return {doc['_id']: doc for doc in doctors}
    
    async def initiate_consultation(self, user_id: str, symptoms: List[str], 
                                  urgency: str = 'normal') -> Dict:
        """发起远程咨询"""
        # 匹配适合的医生
        suitable_doctors = self.match_doctors(symptoms, urgency)
        
        if not suitable_doctors:
            return {'status': 'error', 'message': 'No available doctors'}
        
        # 选择医生
        selected_doctor = suitable_doctors[0]
        
        # 创建咨询会话
        session_id = self.create_consultation_session(user_id, selected_doctor['_id'])
        
        # 通知医生
        await self.notify_doctor(selected_doctor['_id'], session_id, symptoms)
        
        return {
            'status': 'success',
            'session_id': session_id,
            'doctor': selected_doctor,
            'wait_time': self.estimate_wait_time(urgency)
        }
    
    def match_doctors(self, symptoms: List[str], urgency: str) -> List[Dict]:
        """匹配医生"""
        matched_doctors = []
        
        for doctor_id, doctor in self.doctor_pool.items():
            # 检查专业匹配
            if self.check_specialty_match(doctor['specialties'], symptoms):
                # 检查可用性
                if self.check_availability(doctor_id, urgency):
                    matched_doctors.append(doctor)
        
        # 按评分排序
        return sorted(matched_doctors, key=lambda x: x.get('rating', 0), reverse=True)
    
    def check_specialty_match(self, doctor_specialties: List[str], symptoms: List[str]) -> bool:
        """检查专业匹配"""
        # 简化的匹配逻辑
        symptom_specialties = set()
        for symptom in symptoms:
            if symptom in ['headache', 'dizziness']:
                symptom_specialties.add('neurology')
            elif symptom in ['fever', 'cough']:
                symptom_specialties.add('internal_medicine')
        
        return bool(symptom_specialties.intersection(doctor_specialties))
    
    def create_consultation_session(self, user_id: str, doctor_id: str) -> str:
        """创建咨询会话"""
        session_id = hashlib.md5(f"{user_id}_{doctor_id}_{time.time()}".encode()).hexdigest()
        
        session_data = {
            'session_id': session_id,
            'user_id': user_id,
            'doctor_id': doctor_id,
            'status': 'waiting',
            'created_at': datetime.now(),
            'symptoms': [],
            'prescriptions': [],
            'messages': []
        }
        
        self.db.consultations.insert_one(session_data)
        self.consultation_sessions[session_id] = session_data
        
        return session_id
    
    async def notify_doctor(self, doctor_id: str, session_id: str, symptoms: List[str]):
        """通知医生"""
        # 通过微信通知医生
        doctor = self.doctor_pool[doctor_id]
        message = f"新的远程咨询请求:症状{', '.join(symptoms)},会话ID:{session_id}"
        
        # 这里应该是调用微信API发送消息
        print(f"通知医生 {doctor['name']}: {message}")

class MedicationManager:
    """用药管理系统"""
    
    def __init__(self):
        self.db = mongo_client['medication']
        self.drug_knowledge_base = self.load_drug_knowledge()
    
    def load_drug_knowledge(self) -> Dict:
        """加载药品知识库"""
        drugs = list(self.db.drugs.find({}))
        return {drug['name']: drug for drug in drugs}
    
    async def create_medication_plan(self, prescription: Dict, user_data: Dict) -> Dict:
        """创建用药计划"""
        medication_plan = {
            'user_id': user_data['user_id'],
            'prescription': prescription,
            'schedule': [],
            'reminders': [],
            'created_at': datetime.now(),
            'status': 'active'
        }
        
        # 生成用药时间表
        medication_plan['schedule'] = self.generate_schedule(
            prescription['medications'], 
            prescription['duration']
        )
        
        # 生成提醒
        medication_plan['reminders'] = self.generate_reminders(
            medication_plan['schedule'],
            user_data['preferences']
        )
        
        # 保存到数据库
        self.db.medication_plans.insert_one(medication_plan)
        
        return medication_plan
    
    def generate_schedule(self, medications: List[Dict], duration_days: int) -> List[Dict]:
        """生成用药时间表"""
        schedule = []
        current_time = datetime.now()
        
        for day in range(duration_days):
            for medication in medications:
                for time_slot in medication['times_per_day']:
                    schedule.append({
                        'medication': medication['name'],
                        'dosage': medication['dosage'],
                        'scheduled_time': current_time.replace(
                            hour=time_slot['hour'],
                            minute=time_slot['minute']
                        ),
                        'status': 'pending'
                    })
            current_time += timedelta(days=1)
        
        return schedule
    
    def generate_reminders(self, schedule: List[Dict], user_preferences: Dict) -> List[Dict]:
        """生成用药提醒"""
        reminders = []
        
        for schedule_item in schedule:
            reminder_time = schedule_item['scheduled_time'] - timedelta(minutes=15)
            
            reminders.append({
                'medication': schedule_item['medication'],
                'dosage': schedule_item['dosage'],
                'reminder_time': reminder_time,
                'message': self.generate_reminder_message(schedule_item, user_preferences),
                'status': 'pending'
            })
        
        return reminders
    
    def generate_reminder_message(self, schedule_item: Dict, preferences: Dict) -> str:
        """生成提醒消息"""
        return f"用药提醒:请服用{schedule_item['medication']} {schedule_item['dosage']}"

# 初始化核心组件
medical_kg = MedicalKnowledgeGraph()
health_processor = HealthDataProcessor()
consultation_system = RemoteConsultationSystem()
medication_manager = MedicationManager()

# API路由
@app.route('/health/assessment', methods=['POST'])
async def health_assessment():
    """健康评估接口"""
    data = request.json
    user_data = data.get('user_data', {})
    health_data = data.get('health_data', {})
    
    # 处理健康数据
    processed_data = await health_processor.process_health_data(health_data)
    
    # 症状诊断
    if 'symptoms' in processed_data:
        symptoms = processed_data['symptoms']['main_symptoms']
        diagnosis = medical_kg.diagnose_symptoms(symptoms, user_data)
        processed_data['diagnosis'] = diagnosis
    
    # 加密存储
    encrypted_data = cipher_suite.encrypt(json.dumps(processed_data).encode())
    
    # 保存到数据库
    mongo_client.health_assessments.insert_one({
        'user_id': user_data.get('user_id'),
        'encrypted_data': encrypted_data,
        'created_at': datetime.now(),
        'risk_level': processed_data.get('risk_assessment', {}).get('overall', 'low')
    })
    
    return jsonify(processed_data)

@app.route('/telemedicine/consult', methods=['POST'])
async def telemedicine_consult():
    """远程诊疗接口"""
    data = request.json
    user_id = data.get('user_id')
    symptoms = data.get('symptoms', [])
    urgency = data.get('urgency', 'normal')
    
    result = await consultation_system.initiate_consultation(user_id, symptoms, urgency)
    return jsonify(result)

@app.route('/medication/plan', methods=['POST'])
async def create_medication_plan():
    """创建用药计划接口"""
    data = request.json
    prescription = data.get('prescription')
    user_data = data.get('user_data')
    
    plan = await medication_manager.create_medication_plan(prescription, user_data)
    return jsonify(plan)

@app.route('/emergency/alert', methods=['POST'])
async def emergency_alert():
    """紧急情况预警接口"""
    data = request.json
    user_id = data.get('user_id')
    health_data = data.get('health_data')
    
    # 分析紧急情况
    processed_data = await health_processor.process_health_data(health_data)
    is_emergency = self.detect_emergency(processed_data)
    
    if is_emergency:
        # 通知紧急联系人
        await self.notify_emergency_contacts(user_id, processed_data)
        # 自动呼叫急救
        await self.call_emergency_service(user_id, processed_data)
    
    return jsonify({
        'is_emergency': is_emergency,
        'actions_taken': ['notify_contacts', 'call_emergency'] if is_emergency else []
    })

# Celery任务
@celery.task
async def process_batch_health_data(batch_data: List[Dict]):
    """批量处理健康数据"""
    results = []
    for health_data in batch_data:
        try:
            processed = await health_processor.process_health_data(health_data)
            results.append(processed)
        except Exception as e:
            results.append({'error': str(e)})
    
    return results

@celery.task
async def send_medication_reminders():
    """发送用药提醒任务"""
    current_time = datetime.now()
    upcoming_reminders = list(mongo_client.medication_plans.find({
        'reminders.reminder_time': {'$lte': current_time + timedelta(minutes=5)},
        'reminders.status': 'pending'
    }))
    
    for plan in upcoming_reminders:
        for reminder in plan['reminders']:
            if reminder['status'] == 'pending' and reminder['reminder_time'] <= current_time:
                # 发送微信提醒
                await send_wechat_reminder(plan['user_id'], reminder['message'])
                # 更新状态
                mongo_client.medication_plans.update_one(
                    {'_id': plan['_id'], 'reminders.message': reminder['message']},
                    {'$set': {'reminders.$.status': 'sent'}}
                )

async def send_wechat_reminder(user_id: str, message: str):
    """发送微信提醒"""
    # 这里应该是调用微信API发送模板消息
    print(f"发送微信提醒给用户 {user_id}: {message}")

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

使用说明

1. 环境准备

# 安装核心依赖
pip install torch transformers flask flask-socketio celery pymongo redis minio plotly opencv-python pillow cryptography mysql-connector-python sqlalchemy

# 安装医疗相关模型
pip install medbert transformers[medical]

# 启动依赖服务
docker run -d -p 27017:27017 mongo
docker run -d -p 6379:6379 redis
docker run -d -p 3306:3306 mysql

2. 数据库配置

  1. 创建医疗数据库:
CREATE DATABASE medical_db;
CREATE TABLE symptoms (id INT, name VARCHAR(255), category VARCHAR(100));
CREATE TABLE diseases (id INT, name VARCHAR(255), symptoms JSON);
CREATE TABLE treatments (id INT, disease_id INT, medication JSON, procedure TEXT);
  1. 导入医疗知识图谱数据

3. 微信配置

  1. 在微信公众平台配置服务器地址
  2. 获取API密钥和访问令牌
  3. 配置模板消息和客服接口

4. 启动系统

# 启动健康管理系统
python health_management_system.py

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

# 启动定时任务
celery -A healthcare beat --loglevel=info

功能扩展建议

  1. ​区块链医疗记录​​:使用区块链技术确保医疗数据不可篡改
  2. ​AI辅助诊断​​:集成更先进的医疗AI诊断模型
  3. ​物联网设备集成​​:连接更多智能医疗设备
  4. ​多语言支持​​:扩展支持多语言医疗咨询
  5. ​心理健康模块​​:增加心理健康评估和咨询功能

适用场景

✅ 个人健康管理与监测
✅ 慢性病患者远程监护
✅ 基层医疗机构诊疗支持
✅ 突发疾病紧急响应
✅ 用药依从性管理

这个系统整合了:

  • 医疗知识图谱与智能诊断
  • 多模态健康数据分析
  • 远程诊疗工作流
  • 用药智能管理
  • 紧急情况响应
  • 隐私安全保护

适合医疗机构、健康科技公司和个人用户使用。

更多推荐