Apache Spark Structured Streaming

FreeGuideOnline 最新 2026-07-12

输入数据流 --> 无界表 --> 查询逻辑 --> 结果表 --> 输出 (每批追加) (声明式转换) (持续更新) (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

仅当满足如下条件时才允许:

  1. 两流都定义了 watermark。
  2. 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.idquery.namequery.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 为单位,复用批处理作业的写入逻辑(如写入多个表、自定义更新等)。每个微批次都会调用用户定义的函数,传入 batchDFbatchId

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