共计 3676 个字符,预计需要花费 10 分钟才能阅读完成。
背景痛点
电商平台在大促期间面临用户行为数据爆炸式增长的挑战。这类数据具有三个显著特点:

- 高并发 :秒杀场景下 QPS 可达 10 万 +,如双 11 期间某头部电商点击事件峰值达 45 万 / 秒
- 非结构化 :包含浏览轨迹(JSON)、点击坐标(ProtoBuf)、搜索词(文本)等异构数据
- 时效敏感 :实时推荐需在 500ms 内完成从行为发生到策略更新的闭环
传统批处理方案(如 Hadoop MR)存在明显局限:
- 延迟高 :小时级的数据处理周期无法满足实时风控需求
- 资源浪费 :为应对峰值配置的集群在平时利用率不足 30%
- 状态管理弱 :难以处理用户跨会话的行为分析(如 30 分钟内加购又删除)
技术选型
消息队列对比
| 维度 | Kafka | Pulsar |
|---|---|---|
| 峰值吞吐 | 2MB/s/partition(SSD 环境) | 1.5MB/s/partition |
| 积压能力 | 依赖磁盘容量 | 支持分层存储(Tiered Storage) |
| 延迟稳定性 | 99 线在 100ms 内 | 存在 200ms 左右的毛刺 |
选择建议 :当需要处理历史数据回放时,Pulsar 的持久化订阅更优;纯实时场景选 Kafka 更经济
流处理框架
Flink 的核心优势在于:
- 状态管理 :通过 Keyed State/Operator State 实现精确一次(Exactly-Once)处理
- 事件时间 :支持 Watermark 机制处理乱序事件(相比 Spark Streaming 的微批更精准)
- 资源弹性 :TaskManager 可动态调整并行度(on YARN/K8s)
存储格式
# Parquet vs Avro 写入示例
# Parquet 适合分析场景(列裁剪 / 谓词下推)df.write \
.format("parquet") \
.option("compression", "snappy") \
.save("/path/to/user_behavior")
# Avro 适合 CDC 场景(Schema 演进友好)df.write \
.format("avro") \
.option("compression", "deflate") \
.save("/path/to/cdc_logs")
核心实现
分层存储架构
flowchart LR
A[Kafka] -->| 实时消费 | B(Flink)
B --> C{数据类型}
C -->| 热数据 | D[Redis]
C -->| 温数据 | E[Elasticsearch]
C -->| 冷数据 | F[S3+Hive]
- 热数据 :保留 7 天,存放 Redis Cluster(如用户最近浏览记录)
- 温数据 :保留 30 天,存 ES 集群(支持多维度聚合查询)
- 冷数据 :永久保存,转 Parquet 格式入数据湖
会话检测实现
public class SessionTimeoutDetector
extends KeyedProcessFunction<String, UserEvent, String> {
private ValueState<Long> lastActiveTimeState;
@Override
public void open(Configuration parameters) {lastActiveTimeState = getRuntimeContext()
.getState(new ValueStateDescriptor<>("lastActive", Long.class));
}
@Override
public void processElement(UserEvent event, Context ctx, Collector<String> out) {
// 更新状态并注册定时器
long currentTime = ctx.timestamp();
lastActiveTimeState.update(currentTime);
ctx.timerService().registerEventTimeTimer(currentTime + 30 * 60 * 1000);
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<String> out) {
// 检查是否真的超时
if (lastActiveTimeState.value() != null &&
timestamp >= lastActiveTimeState.value() + 30 * 60 * 1000) {out.collect("用户" + ctx.getCurrentKey() + "会话已超时");
}
}
}
性能优化
Watermark 调优
处理迟到数据的三种策略:
- Allowed Lateness:窗口关闭后保留一段时间
.window(TumblingEventTimeWindows.of(Time.minutes(5))) .allowedLateness(Time.minutes(1)) // 允许迟到 1 分钟 - SideOutput:将超时数据转到侧输出流
OutputTag<UserEvent> lateDataTag = new OutputTag<>("late-data"); SingleOutputStreamOperator<Result> mainStream = stream .keyBy(...) .window(...) .sideOutputLateData(lateDataTag) .process(...); - 动态 Watermark:根据历史延迟分布自适应调整
维表关联优化
# 使用 Guava Cache 防止缓存穿透
cache = CacheBuilder.newBuilder()
.maximumSize(100000)
.expireAfterWrite(5, TimeUnit.MINUTES)
.build(CacheLoader.from(lambda key: queryFromDB(key))
)
# 在 Flink 中异步访问
class AsyncRedisLookup(AsyncFunction[String, (String, String)]):
def asyncInvoke(self, key: str, result: ResultFuture[(str, str)]):
try:
# 先从本地缓存查
cached = cache.getIfPresent(key)
if cached:
result.complete([(key, cached)])
else:
# 异步查询 Redis
future = redis_client.get(key)
future.add_done_callback(lambda f: result.complete([(key, f.result())])
)
except Exception as e:
result.completeExceptionally(e)
避坑指南
Flink 反压处理
关键配置参数:
taskmanager.network.memory.fraction:网络缓冲区占比(建议 0.2)taskmanager.numberOfTaskSlots:每个 TM 的 slot 数(建议等于 CPU 核数)restart-strategy.fixed-delay.attempts:失败重试次数(生产环境建议 3)state.backend.rocksdb.memory.managed:开启 RocksDB 托管内存table.exec.state.ttl:状态存活时间(根据业务设置)
Kafka 分区设计
计算公式:
分区数 = max(
当前峰值吞吐 / 单分区处理能力,
未来 6 个月预估峰值 * 1.5 / 单分区处理能力
)
其中单分区处理能力通常为:
– 普通消费:2MB/s
– 压缩消息:4-6MB/s
数据倾斜处理
典型解决方案:
- 加盐处理 :对倾斜 Key 添加随机后缀
-- 原始 SQL SELECT user_id, COUNT(*) FROM clicks GROUP BY user_id; -- 优化后 SELECT REPLACE(user_id, '_suffix', ''), SUM(cnt) FROM ( SELECT CONCAT(user_id, '_', CAST(RAND()*10 AS INT)) AS user_id, COUNT(*) AS cnt FROM clicks GROUP BY CONCAT(user_id, '_', CAST(RAND()*10 AS INT)) ) GROUP BY REPLACE(user_id, '_suffix', ''); - 本地聚合 :在 Map 阶段先做预聚合
- 倾斜 Key 隔离 :单独处理热点用户
未来演进
如果引入 ClickHouse,可在以下方面获得提升:
- 查询性能 :单表千亿级数据秒级响应
- 存储效率 :列存压缩比可达 10:1
- 实时分析 :支持 Kafka 引擎直接消费消息
但需注意其不适合高频更新的场景,且 JOIN 操作性能较差。建议将 ClickHouse 作为实时 OLAP 层,与现有流处理架构配合使用。
结语
这套方案在某电商 618 大促期间经受住了实战检验,核心指标表现:
– 峰值处理能力:12 万事件 / 秒
– 端到端延迟:90% 请求 <300ms
– 资源成本:比原批处理方案节省 40%
读者可基于文中代码示例快速搭建原型,建议先从流量较小的非核心业务开始验证,逐步完善监控和容错机制。
正文完
发表至: 未分类
近一天内
