共计 1935 个字符,预计需要花费 5 分钟才能阅读完成。
背景与痛点
电商用户行为数据通常具有高维度、实时性要求高等特点。以 2024 年的数据集为例,单日日志量可达 PB 级,且包含点击流、搜索词、加购行为等 20+ 维度字段。主要面临三大挑战:

- 数据倾斜 :热门商品的行为日志可能占据 80% 流量,导致常规 Hash 分片失效
- 实时性瓶颈 :从行为发生到分析结果的延迟需控制在 5 秒内(如实时推荐场景)
- 维度爆炸 :用户画像标签组合可能超过 10 万种,传统 JOIN 操作性能骤降
技术方案对比
批处理(Spark)适用场景
- 离线报表生成(如日活月活统计)
- 历史数据挖掘(用户生命周期分析)
- 特征工程(机器学习样本准备)
流处理(Flink)适用场景
- 实时风控(如刷单行为检测)
- 动态定价(竞品价格监控)
- 会话级推荐(购物车关联推荐)
选型建议 :采用 Lambda 架构,用 Spark 处理 T + 1 数据,Flink 处理实时流。关键指标(如 GMV)需双链路校验。
核心实现
数据清洗示例(Python/Pandas)
import pandas as pd
import numpy as np
def clean_behavior_data(raw_df):
# 处理缺失值:设备 ID 为空的用会话 ID 填充
df = raw_df.copy()
df['device_id'] = df['device_id'].fillna(df['session_id'])
# 过滤异常值:点击时间早于 2024 年的记录
df = df[df['click_time'] >= pd.Timestamp('2024-01-01')]
# 标准化:统一用户行为类型为小写
df['action_type'] = df['action_type'].str.lower()
return df
存储优化:Parquet 实战优势
- 压缩比 :相比 CSV 节省 60% 存储空间
- 查询性能 :列剪裁使扫描量减少 5 -10 倍
- Schema 演进 :支持向后兼容的字段增减
实时分析(Flink 窗口计算)
// 5 分钟滚动窗口统计类目点击量
DataStream<CategoryClick> clicks = env
.addSource(new KafkaSource())
.keyBy("category_id")
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new AggregateFunction<UserEvent, Long, Long>() {
@Override
public Long createAccumulator() { return 0L;}
@Override
public Long add(UserEvent event, Long acc) {return acc + 1;}
@Override
public Long getResult(Long acc) {return acc;}
@Override
public Long merge(Long acc1, Long acc2) {return acc1 + acc2;}
});
性能调优
解决 Join 倾斜(Spark 示例)
// 对大表进行随机前缀扩容
val expandedLargeTable = largeTable
.withColumn("join_key", concat(col("product_id"),
lit("_"),
(rand() * 10).cast("int")
))
// 对小表进行笛卡尔扩容
val expandedSmallTable = smallTable
.crossJoin(spark.range(0,10).toDF("prefix_id"))
.withColumn("join_key", concat(col("product_id"),
lit("_"),
col("prefix_id")
))
// 执行 Join 后聚合
val result = expandedLargeTable
.join(expandedSmallTable, "join_key")
.groupBy("product_id")
.agg(sum("click_count"))
避坑指南
- 事件时间乱序
- 对策:设置合理 Watermark(如允许 2 秒乱序)
-
代码:
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) -
状态后端配置错误
- 现象:Checkpoint 超时失败
-
解决:RocksDBStateBackend + 增量检查点
-
Flink 反压传导
- 识别:Web UI 显示高背压指标
- 优化:调整并行度或启用本地 KeyBy
开放思考
当实时处理链路出现 10 分钟级延迟时,如何在以下维度做权衡:
– 降级方案:改用近实时数据?
– 资源分配:牺牲其他业务计算资源?
– 准确性:允许部分指标暂时不准确?
期待大家在评论区分享实战经验。
正文完
发表至: 未分类
近一天内
