Apache Spark Structured Streaming
输入数据流 --> 无界表 --> 查询逻辑 --> 结果表 --> 输出 (每批追加) (声明式转换) (持续更新) (Sink)
- 每触发一次(Trigger),Spark 读取新到达的数据,将其视为追加行。
- 在结果表中执行用户定义的增量计算(如聚合、join 等)。
- 将更新后的结果按指定的输出模式写入外部系统。
## 快速开始:从 Socket 读取单词计数
以下用经典的实时单词计数示例演示基本流程。使用 Spark 2.4+ 版本,语言选择 Scala 或 Python(本教程使用 Python 示例,Scala 逻辑相同)。
### 创建 SparkSession
```python
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("StructuredNetworkWordCount") \
.getOrCreate()
spark.sparkContext.setLogLevel("WARN")
读取流数据源
# 从 socket 读取文本流,每行一个数据
lines = spark.readStream \
.format("socket") \
.option("host", "localhost") \
.option("port", 9999) \
.load()
此时 lines 是一个流式 DataFrame,包含一列 value(字符串)。
定义转换逻辑
使用与批处理相同的 DataFrame 操作:
from pyspark.sql.functions import explode, split
words = lines.select(
explode(split(lines.value, " ")).alias("word")
)
wordCounts = words.groupBy("word").count()
启动流查询并输出到控制台
query = wordCounts.writeStream \
.outputMode("complete") \
.format("console") \
.option("truncate", "false") \
.start()
query.awaitTermination()
outputMode("complete"):每次触发器执行后,将完整的结果表全部输出(因为使用了聚合,且没有窗口)。format("console"):将结果打印到控制台,适合调试。
运行前需先在终端用 nc -lk 9999 启动一个 socket 服务器,然后输入单词行,即可看到实时词频统计。
核心概念详解
输入源(Source)
Structured Streaming 内置支持多种可靠的数据源:
| 源 | 是否容错 | 说明 |
|---|---|---|
| File Source | 是 | 支持 JSON, CSV, Parquet, ORC, text 等。监控目录,处理新增文件。 |
| Kafka Source | 是 | 消费 Kafka 主题,支持 offset 管理和批次读取。 |
| Socket Source | 否 | 仅测试用,无容错保证。 |
| Rate Source | 否 | 自动生成数据,用于测试和基准测试。 |
示例:读取 Kafka 数据
df = spark.readStream \
.format("kafka") \
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") \
.option("subscribe", "topic1") \
.load()
df.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)")
输出模式(Output Mode)
定义结果表如何写入外部 Sink。必须在启动流查询时指定,并与查询类型适配:
- Append 模式:只输出自上次触发以来新增的行到结果表。适用于仅做选择、过滤、映射等不改变已有行的操作。聚合查询不支持,除非指定了 watermark 并使用窗口聚合。
- Complete 模式:每次触发都将整个结果表完全输出。仅适用于包含聚合(且无状态算子需要旧行)的查询。
- Update 模式:只输出结果表中有更新的行(自上次触发以来发生变化)。适用于有聚合且没有 watermarked 窗口聚合的场景。从 Spark 2.1.1 起可用。
典型对应关系:
| 查询类型 | 支持的模式 |
|---|---|
| 无聚合(select, where, map 等) | Append, Update |
| 聚合(groupBy) | Complete, Update |
| 带有 watermark 的窗口聚合 + groupBy | Append, Update |
触发器(Trigger)
控制流查询的执行频率:
- 默认(微批次):上一批次处理完毕后立即启动下一批次。
- 固定间隔微批次:
trigger(processingTime="5 seconds") - 一次性微批次:
trigger(once=True)处理所有可用数据后停止。 - 持续处理(实验性):
trigger(continuous="1 second")实现毫秒级延迟。
# 每10秒触发一次
query = wordCounts.writeStream \
.trigger(processingTime="10 seconds") \
... \
.start()
输出 Sink
常用 Sink 类型:
| Sink | 说明 |
|---|---|
console |
打印到控制台,用于调试。 |
kafka |
将结果写入 Kafka 主题。 |
foreach / foreachBatch |
自定义写入逻辑,可复用外部连接池。 |
file |
写入文件(支持 JSON, Parquet 等),滚动输出。 |
memory |
将结果保存到内存表,供交互式查询。 |
示例:写入 Parquet 文件,带分区列:
query = df.writeStream \
.outputMode("append") \
.format("parquet") \
.option("path", "/output/path/") \
.option("checkpointLocation", "/checkpoint/path/") \
.start()
窗口操作与事件时间
事件时间与 Watermark
Structured Streaming 支持基于事件时间(数据自身携带的时间戳)的窗口聚合,并能处理延迟数据。通过 withWatermark() 指定允许的最大延迟阈值,引擎会维护状态并自动丢弃过期的聚合状态,避免状态无限增长。
from pyspark.sql.functions import window
events = spark.readStream \
.format("...") \
.load()
# 假设数据包含 timestamp 列 'event_time'
# 添加 watermark 并定义 10 分钟滚动窗口,滑动步长 5 分钟
windowedCounts = events \
.withWatermark("event_time", "10 minutes") \
.groupBy(
window(events.event_time, "10 minutes", "5 minutes"),
events.word
) \
.count()
此时可以安全地使用 Append 或 Update 输出模式,因为 watermark 保证了超时的聚合窗口不会被再次更新,引擎可以安全地将“已完结”的窗口结果输出。
Watermark 的本质:定义一个时间阈值 T,引擎认为任何时间戳 max(event_time) - T 之前的数据不会再出现,因此可以丢弃维护的中间状态并输出最终结果。选择 T 需平衡延迟容忍和状态大小。
窗口类型
- 滚动窗口:
window(column, "10 minutes")窗口长度 = 滑动步长。 - 滑动窗口:
window(column, "10 minutes", "5 minutes")窗口长度 > 滑动步长,有重叠。
窗口例: {start: 00:00, end: 00:10}, {start: 00:05, end: 00:15} 等。
流式 Dataset/DataFrame 操作
几乎所有 Spark SQL 和 DataFrame 的操作都可以应用于流式 DataFrames,但有一些限制:
- 支持:投影、过滤、选择、join(有限制)、聚合(带或不带窗口)、UDF(标记为 SQL 函数)等。
- 不支持:排序(除非在聚合后输出完整模式下)、limit、take、distinct 在无状态流上(或需要相应模式)。
- 流-流 Join 必须使用 watermark 和时间区间条件来管理状态。
- 流-静态 Join 可以无限制地对流数据和静态 DataFrame 做 join。
流-流 Join
仅当满足如下条件时才允许:
- 两流都定义了 watermark。
- Join 条件中包含时间范围约束(例如事件时间间隔 < 某个值)。
from pyspark.sql.functions import expr
impressionsWithWatermark = impressions \
.withWatermark("impression_time", "2 hours")
clicksWithWatermark = clicks \
.withWatermark("click_time", "3 hours")
# 内连接,要求点击时间在曝光时间之后且不超过1小时
joined = impressionsWithWatermark.join(
clicksWithWatermark,
expr("""
ad_id = click_ad_id AND
click_time >= impression_time AND
click_time <= impression_time + interval 1 hour
"""),
"inner"
)
管理流查询
启动 start() 返回一个 StreamingQuery 对象,用于监控和管理:
query.id、query.name、query.status获取执行状态。query.awaitTermination()阻塞直到查询终止。query.stop()主动停止。- 使用
spark.streams.active查看所有活跃查询,spark.streams.get(id)获取特定查询。
query = df.writeStream.format("console").start()
# ... 其他代码
query.stop()
容错与 Exactly-Once 语义
Structured Streaming 通过 预写日志 (Checkpoint) 和 幂等输出 Sink 提供端到端的 exactly-once 保证。
Checkpoint 位置
每个流查询必须指定一个 checkpoint 目录,用于存储:
- 已处理数据的偏移量(例如 Kafka offset)。
- 中间聚合状态。
- 查询元数据。
因此即使失败重启,也能从断点正确恢复,不丢不重。
query = df.writeStream \
.option("checkpointLocation", "/checkpointdir") \
.start()
支持的 Sink 的保证级别
| Sink | 语义 | 说明 |
|---|---|---|
| File | exactly-once | 使用幂等文件写入和事务提交。 |
| Kafka | at-least-once 或 exactly-once | 若 Kafka 版本和配置允许幂等写入,可实现 exactly-once。 |
| foreach | at-least-once | 自定义逻辑需自行处理幂等。 |
高级特性:foreach 与 foreachBatch
foreachBatch
允许以微批次内的 DataFrame 为单位,复用批处理作业的写入逻辑(如写入多个表、自定义更新等)。每个微批次都会调用用户定义的函数,传入 batchDF 和 batchId。
def write_batch(df, epoch_id):
# 将这一小批数据写入现有 Delta 表,或调用其他批处理 API
df.write.format("parquet").mode("append").save("/data/batch")
pass
df.writeStream.foreachBatch(write_batch).start()
foreach
对结果表的每一行执行用户自定义操作,适合写入外部存储。需实现 open, process, close 方法。
class MySQLWriter:
def open(self, partition_id, epoch_id):
# 打开连接
self.connection = ...
return True
def process(self, row):
# 写入一行
self.connection.execute(...)
def close(self, error):
self.connection.close()
query = df.writeStream.foreach(MySQLWriter()).start()