跳到正文
Joeplover
学习笔记·2025-10-15·约 3 分钟阅读

Spark Streaming + Kafka + Redis 用户行为实时分析

基于 Spark Streaming 和 Kafka 的实时用户行为分析系统,消费 Kafka 消息并聚合写入 Redis。

实时用户行为

Spark Streaming + Kafka + Redis 用户行为实时分析

架构

用户行为数据 → Kafka → Spark Streaming(微批次处理)→ Redis(存储聚合结果)→ REST API(查询展示)

为什么用 Spark Streaming

这是一个接近实时的用户行为分析系统。当用户在电商平台产生浏览、点击、加购、购买等行为时,我们需要实时(秒级延迟)了解:

  • 当前热销商品排名
  • 实时 PV/UV
  • 用户行为漏斗(浏览→加购→购买转化率)

Spark Streaming 的微批次模型(Micro-Batch)在这里表现良好——几秒的延迟对数据分析来说可以接受。如果要求毫秒级延迟,则需要 Flink。

微批次原理

Spark Streaming 将实时数据流切分成固定间隔的小批次,每个批次作为一个 RDD 处理:

from pyspark.streaming import StreamingContext
from pyspark.streaming.kafka import KafkaUtils

ssc = StreamingContext(sc, batchDuration=5)  # 每 5 秒一个批次

kafkaStream = KafkaUtils.createStream(
    ssc, 
    "192.168.8.133:2181",  # ZooKeeper
    "user-behavior-group",   # 消费者组
    {"user-behavior": 3}     # 3 个分区
)

窗口聚合

# 热搜商品:统计最近 10 分钟的点击量
hotProducts = (kafkaStream
    .map(lambda msg: json.loads(msg[1]))
    .filter(lambda event: event["type"] == "click")
    .map(lambda event: (event["product_id"], 1))
    .reduceByKeyAndWindow(
        lambda a, b: a + b,       # 窗口内聚合
        lambda a, b: a - b,       # 窗口滑动时移出旧数据
        windowDuration=600,       # 窗口长度:10 分钟
        slideDuration=30          # 滑动间隔:30 秒
    ))

mapWithState:跨批次状态

# 统计各商品的历史累计点击数
def updateState(key, value, state):
    old_count = state.get() or 0
    new_count = (old_count + sum(value)) if value else old_count
    return (key, new_count)

statefulStream = kafkaStream.map(...).mapWithState(
    StateSpec.function(updateState).timeout(Durations.minutes(30))
)
# 30 分钟无更新的商品自动清除,节省内存

背压机制

# 开启背压后,Spark Streaming 自动调节消费速率
ssc = StreamingContext(sc, 5)
ssc.sparkContext.setConf("spark.streaming.backpressure.enabled", "true")

当处理速度跟不上消费速度时,背压机制会自动降低 Kafka 消费速率,防止数据积压导致 OOM。

存储到 Redis

import redis

r = redis.StrictRedis(host='192.168.8.133', port=6379, password='redis123')

def save_to_redis(rdd):
    if rdd.isEmpty(): return
    for product_id, count in rdd.collect():
        r.zadd("hot_products", {product_id: count})  # 有序集合存储排名

hotProducts.foreachRDD(save_to_redis)

使用 Redis 的 Sorted Set(有序集合)存储热搜排名,天然支持 Top-N 查询:

# 查看 Top 10 热搜商品
ZREVRANGE hot_products 0 9 WITHSCORES

端到端延迟

微批次模型的端到端延迟 = 批次间隔 + 处理时间 + 存储时间。5 秒间隔下,总延迟约 5-10 秒。

对于毫秒级实时场景,推荐 Flink(事件驱动、真正的流处理)。但对于大多数"准实时"分析需求(秒级延迟),Spark Streaming 已经足够,而且可以复用 Spark 的其余生态(MLlib、SQL)。