AI 驱动用户画像构建:从标签体系到实时特征工程的落地路径

cover

一、用户画像的工程困境:标签爆炸与实时性矛盾

用户画像是个性化推荐、精准营销和风险控制的基础设施。传统方案基于规则引擎,运营人员手动定义标签规则(如"近30天购买次数≥5"标记为高活跃用户),然后离线批处理生成画像数据。这种方案在标签数量较少时可行,但随着业务复杂度增长,标签数量从几十膨胀到数千,规则维护成本急剧上升,标签之间的冲突和冗余也难以发现。

更关键的是,离线批处理的画像数据存在小时级甚至天级的延迟,无法满足实时风控(如反欺诈需要秒级判断)和实时推荐(如用户刚浏览完商品立即推荐相关内容)的需求。AI 驱动的用户画像方案,通过机器学习自动挖掘标签、流式计算实时更新特征,试图解决这两个核心矛盾。

本文将系统分析 AI 用户画像的技术架构,给出从标签体系设计到实时特征工程的完整方案。

二、从规则引擎到特征平台:AI 用户画像的技术架构

AI 驱动的用户画像系统,核心变化是将"规则定义标签"转变为"特征平台统一管理"。特征平台负责特征的注册、计算、存储和服务,机器学习模型消费特征而非原始数据。

flowchart TD
    A[原始数据源] --> B[特征计算层]
    B --> C[特征存储层]
    C --> D[特征服务层]
    D --> E[模型消费]

    A --> A1[行为日志]
    A --> A2[交易数据]
    A --> A3[设备信息]

    B --> B1[离线批处理:Spark]
    B --> B2[实时流处理:Flink]

    C --> C1[Redis:实时特征]
    C --> C2[HBase:历史特征]

    D --> D1[特征点查API]
    D --> D2[特征批量导出]

    E --> E1[推荐模型]
    E --> E2[风控模型]
    E --> E3[营销模型]

    style B fill:#e1f5fe,stroke:#0288d1,stroke-width:2px
    style C fill:#fff3e0,stroke:#f57c00,stroke-width:2px

标签体系设计

标签分为三类:事实标签(直接从原始数据提取,如性别、注册时间)、统计标签(对原始数据聚合计算,如近7天购买次数)、预测标签(模型推断,如流失概率、购买意向)。AI 的核心价值在预测标签——通过模型自动发现人工规则难以捕捉的隐含模式。

实时特征计算

实时特征的计算依赖流式处理引擎。以 Flink 为例,用户行为事件进入 Kafka 后,Flink 消费事件并维护滑动窗口内的聚合状态(如最近1小时的点击次数),计算结果写入 Redis 供在线服务查询。关键挑战是窗口状态的精确维护和迟到数据的处理。

三、生产级代码实现与最佳实践

特征定义与注册

from dataclasses import dataclass
from enum import Enum
from typing import Any, Optional

class FeatureType(Enum):
    """特征类型枚举"""
    CATEGORICAL = "categorical"   # 类别型
    NUMERICAL = "numerical"       # 数值型
    BOOLEAN = "boolean"           # 布尔型

class ComputeMode(Enum):
    """计算模式"""
    BATCH = "batch"               # 离线批处理
    STREAMING = "streaming"       # 实时流处理
    HYBRID = "hybrid"             # 混合模式

@dataclass
class FeatureDefinition:
    """特征定义"""
    name: str
    description: str
    feature_type: FeatureType
    compute_mode: ComputeMode
    source_table: str
    compute_logic: str            # SQL或Python表达式
    default_value: Any = None
    ttl_seconds: Optional[int] = None  # 特征过期时间

# 示例:定义用户购买频率特征
purchase_freq_feature = FeatureDefinition(
    name="user_purchase_freq_7d",
    description="用户近7天购买次数",
    feature_type=FeatureType.NUMERICAL,
    compute_mode=ComputeMode.HYBRID,
    source_table="order_events",
    compute_logic="""
        SELECT user_id, COUNT(*) as purchase_freq_7d
        FROM order_events
        WHERE event_time >= DATE_SUB(CURRENT_TIMESTAMP, INTERVAL 7 DAY)
        GROUP BY user_id
    """,
    default_value=0,
    ttl_seconds=86400  # 24小时过期
)

实时特征计算(Flink SQL)

-- Flink SQL:实时计算用户近1小时点击次数
CREATE TABLE click_events (
    user_id BIGINT,
    item_id BIGINT,
    event_time TIMESTAMP(3),
    WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'click_events',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json'
);

-- 滑动窗口聚合:每分钟输出一次近1小时的点击计数
CREATE TABLE user_click_count_1h AS
SELECT
    user_id,
    COUNT(*) as click_count_1h,
    TUMBLE_START(event_time, INTERVAL '1' MINUTE) as window_start
FROM click_events
GROUP BY
    user_id,
    TUMBLE(event_time, INTERVAL '1' MINUTE);

特征服务 API

from fastapi import FastAPI, HTTPException
from typing import Dict, List

app = FastAPI(title="Feature Service")

class FeatureService:
    """特征服务:统一特征查询接口"""

    def __init__(self, redis_client, hbase_client):
        self.redis = redis_client
        self.hbase = hbase_client

    async def get_features(
        self, user_id: int, feature_names: List[str]
    ) -> Dict[str, Any]:
        """获取用户特征,优先查Redis,miss时查HBase"""
        result = {}

        # 先查Redis实时特征
        redis_key = f"feature:{user_id}"
        cached = self.redis.hgetall(redis_key)
        for name in feature_names:
            if name in cached:
                result[name] = cached[name]

        # 对Redis中miss的特征,查HBase历史数据
        missed = [n for n in feature_names if n not in result]
        if missed:
            hbase_data = self.hbase.get(user_id, missed)
            result.update(hbase_data)

        # 填充默认值
        for name in feature_names:
            if name not in result:
                result[name] = None

        return result

@app.get("/features/{user_id}")
async def get_user_features(user_id: int, features: str):
    """特征查询API:GET /features/123?features=purchase_freq_7d,click_count_1h"""
    feature_names = features.split(",")
    service = FeatureService(redis_client=None, hbase_client=None)
    data = await service.get_features(user_id, feature_names)
    return {"code": 0, "data": data}

四、边界分析与架构权衡

特征一致性问题

离线训练和在线推理使用不同的特征计算引擎(Spark vs Flink),可能导致特征值不一致。例如,离线用 Spark SQL 计算的"近7天购买次数"与在线 Flink 计算的结果存在微小差异(窗口对齐方式不同)。这种训练-推理偏差(Training-Serving Skew)会导致模型在线效果下降。解决方案是统一特征计算逻辑,或使用特征平台(如 Feast)保证一致性。

实时特征的存储成本

Redis 存储实时特征的成本远高于 HBase。当用户量达到千万级、特征数达到数百个时,Redis 内存需求可能超过 100GB。需要根据特征的访问频率和时效性要求,制定分级存储策略:高频实时特征存 Redis,低频历史特征存 HBase,过期特征自动清理。

隐私合规风险

用户画像涉及大量个人数据,需要遵守 GDPR、个人信息保护法等法规。特征平台必须支持数据脱敏、访问审计和用户授权管理。预测标签(如"流失概率")的输出需要可解释,避免黑箱决策引发合规风险。

维度 规则引擎 AI驱动
标签开发效率 低(手动编写规则) 高(模型自动挖掘)
标签可解释性 高(规则透明) 低(模型黑箱)
实时性 差(离线批处理) 好(流式计算)
维护成本 高(规则膨胀) 中(模型迭代)
隐私合规 易控制 需额外审计

五、总结

AI 驱动用户画像的核心价值在于自动挖掘预测标签和实时特征计算,解决了传统规则引擎的标签膨胀和延迟问题。但 AI 方案引入了训练-推理偏差、存储成本和隐私合规等新挑战。工程落地时,建议采用混合架构:事实标签和统计标签仍用规则引擎(可解释、易审计),预测标签用机器学习模型(自动挖掘、实时更新),特征平台统一管理所有特征的注册、计算和服务。理解 AI 用户画像的边界,比盲目追求全自动化更重要。

Logo

AtomGit 是由开放原子开源基金会联合 CSDN 等生态伙伴共同推出的新一代开源与人工智能协作平台。平台坚持“开放、中立、公益”的理念,把代码托管、模型共享、数据集托管、智能体开发体验和算力服务整合在一起,为开发者提供从开发、训练到部署的一站式体验。

更多推荐