共计 1369 个字符,预计需要花费 4 分钟才能阅读完成。
业务场景与批处理局限
以电商推荐系统为例,传统 T + 1 批处理模式存在明显缺陷:

- 用户实时行为数据(如加购 / 收藏)无法及时进入特征库
- 凌晨 ETL 任务集中爆发导致资源争用
- 跨数据源 JOIN 时出现近 24 小时数据不一致
测试显示:当商品价格在日间变动时,批处理模式会导致推荐结果滞后,CTR 下降 12%。
技术方案选型
存储层对比
| 方案 | 适用场景 | AI 数据优势 |
|---|---|---|
| Delta Lake | 频繁更新的维度表 | ACID 事务保障特征一致性 |
| Iceberg | 大规模历史数据分析 | 隐藏分区提升特征遍历效率 |
| Hudi | 实时增量特征回填 | 内置索引加速点查性能 |
流批协同实现
from pyspark.sql import SparkSession
from pyspark.sql.functions import window
# 创建带 Checkpoint 的流式处理
spark = SparkSession.builder \
.config("spark.sql.shuffle.partitions", "200") \
.config("spark.sql.streaming.stateStore.providerClass", "HDFS") \
.getOrCreate()
# 状态计算示例:用户 30 天滚动特征
user_features = spark.readStream \
.format("kafka") \
.option("checkpointLocation", "/checkpoints/user_features") \
.load() \
.groupBy("user_id", window("event_time", "30 days")) \
.agg({"click": "sum", "purchase": "avg"})
数据质量保障
-- Great Expectations 示例
CREATE EXPECTATION order_data
WITH (
expectation_type = 'expect_column_values_to_not_be_null',
column = 'user_id',
action = 'quarantine'
);
-- 自动修复异常数据
INSERT INTO repaired_orders
SELECT * FROM orders WHERE user_id IS NOT NULL;
性能优化实战
分区策略影响
当特征表与用户画像 JOIN 时:
- 按
user_id%100哈希分区:减少 shuffle 数据量 45% - Z-Order 聚类
(user_id,event_date):范围查询提速 3 倍
小文件合并
# Delta Lake 自动压缩
spark.sql("""
OPTIMIZE feature_db.user_behavior
ZORDER BY (user_id)
""")
# 自定义合并策略
for partition in hive_partitions:
coalesce_files(partition, target_size="128MB")
生产环境避坑指南
- 时区陷阱
- 所有时间字段强制 UTC 存储
-
计算窗口使用
TIMESTAMP_LTZ类型 -
特征回填冲突
- 限制回填任务并发度
- 采用
COPY_ON_WRITE模式更新
开放性问题
当数据延迟达到 6 小时:
- 立即重算:消耗 300% 计算资源
- 跳过本轮:影响模型 AUC 0.02
- 折中方案:仅重算核心特征
需要根据业务场景(如金融风控 vs 内容推荐)制定分级响应策略。
正文完
