
数据源为前端埋点日志(JSON 格式)写入 Kafka,目标端为 HDFS 和 Redis。我们选用 Spark SS 的原因在于其声明式 API 能将流处理逻辑与批处理高度统一,且内置的 Watermark 机制和 State Store 能优雅地处理乱序数据和聚合状态。
核心开发思路:
首先定义 JSON Schema,并从 Kafka 消费数据。注意在生产环境中必须显式设置消费组和反序列化参数。
from pyspark.sql import SparkSession
from pyspark.sql.functions import *
from pyspark.sql.types import *
spark = SparkSession.builder \
.appName("RealtimeAnalytics") \
.config("spark.sql.shuffle.partitions", "400") \
.config("spark.sql.streaming.stateStore.providerClass",
"org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider") \
.getOrCreate()
# 定义埋点 Schema
schema = StructType() \
.add("uid", StringType()) \
.add("page_id", StringType()) \
.add("event_time", StringType()) \ # 原始字符串,需转换
.add("action", StringType()) \
.add("session_id", StringType())
# 读取 Kafka 流
df_raw = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "localhost:9092") \
.option("subscribe", "app_logs") \
.option("failOnDataLoss", "false") \
.load() \
.select(from_json(col("value").cast("string"), schema).alias("data")) \
.select("data.*") \
.withColumn("event_ts", to_timestamp(col("event_time"), "yyyy-MM-dd HH:mm:ss"))需求是统计 每分钟的 UV(独立访客数) 和 PV(页面浏览量),并容忍 10 分钟的延迟数据。这里的关键在于使用 withWatermark 配合 groupBy 窗口聚合。
# 设置 Watermark(允许 10 分钟乱序)
df_agg = df_raw \
.withWatermark("event_ts", "10 minutes") \
.groupBy(
window(col("event_ts"), "1 minute", "1 minute"),
col("page_id")
) \
.agg(
count("*").alias("pv"),
approx_count_distinct("uid").alias("uv") # 近似去重,高效
)对于精确去重,建议采用 countDistinct 配合 flatMapGroupsWithState,但生产环境中通常使用 HyperLogLog(如上述 approx_count_distinct)以节省内存。若业务要求极致精确,则必须开启 RocksDB 状态后端,避免 OOM。
流式任务的结果需要写入 HDFS(用于历史追溯)和 Redis(用于大屏展示)。Spark SS 原生支持 foreachBatch 实现双写逻辑。
def write_to_sinks(df_batch, epoch_id):
# 1. 写入 HDFS Parquet(动态分区)
df_batch.write \
.format("parquet") \
.mode("append") \
.partitionBy("page_id") \
.save("/data/warehouse/pv_uv/")
# 2. 写入 Redis(更新实时大屏)
rows = df_batch.collect()
import redis
r = redis.Redis(host='cache-node', port=6379, decode_responses=True)
for row in rows:
window_end = row['window'].end.strftime('%Y-%m-%d %H:%M')
key = f"dash:{row['page_id']}:{window_end}"
r.hset(key, mapping={"pv": row['pv'], "uv": row['uv']})
r.expire(key, 3600 * 24) # 保留 24 小时
# 启动流作业
query = df_agg.writeStream \
.foreachBatch(write_to_sinks) \
.outputMode("append") \
.option("checkpointLocation", "/checkpoints/pv_uv_cp") \
.trigger(processingTime="30 seconds") \
.start()checkpointLocation 至关重要,它记录了消费的 Kafka Offset 和聚合状态快照,是故障恢复的基础。
流式任务跑久了会遇到两大“杀手”:数据倾斜导致单分区处理过慢,以及 HDFS 小文件 过多压垮 NameNode。
热点 Key(如首页 page_id=home)可在聚合前加随机前缀打散(两阶段聚合),但 SS 中状态关联较复杂。更轻量的方案是优化 spark.sql.shuffle.partitions 参数,将其调大至 400~800,使 Shuffle 阶段分区更细。同时开启自适应执行(AQE):
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")写入 HDFS 时,每个微批(30 秒)都会生成少量 Parquet 文件。我们可在 foreachBatch 内执行 repartition 或 coalesce 减少输出分区数。更优方案是启用 Spark 3.x 的 动态分区插入 优化:
# 在 foreachBatch 中写入前合并
df_batch.repartition(5).write.format("parquet")... # 强制合并为 5 个文件另外,可配置 Spark 的 spark.sql.adaptive.skewJoin.enabled 解决 Join 过程中的数据倾斜。
Input Rate 和 Processing Rate,若前者持续高于后者,说明背压触发,需增加资源(Executor 数量)。checkpoint 目录,需定期清理过期 Key。通过设置 spark.sql.streaming.minBatchesToRetain 控制保留状态的最小批次,超过窗口范围的老数据状态会被自动淘汰。本文通过一个完整的 PV/UV 案例,串联了 Spark Structured Streaming 从数据接入、时间语义处理、状态存储到多端 Sink 的全流程开发。相较于传统微批处理,SS 在 Continuous 模式(实验性)下甚至能达到毫秒级延迟。但需要注意的是,流式开发的核心难点并不在 API,而在于对状态容量的预估和端到端一致性保障。建议在压测环境中模拟数据乱序和 Failover 场景,验证 Checkpoint 恢复的有效性。随着数据湖技术(如 Iceberg、Hudi)的成熟,流式写入已从 HDFS 演进到湖存储,但上述调优思想依然通用,希望能为你的大数据开发实战提供扎实的参考。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。