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)。