ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Spark Structured Streaming 编程指南:从编程模型到端到端 Exactly-Once 流处理

Apache Spark Structured Streaming 编程指南:从编程模型到端到端 Exactly-Once 流处理 Apache Spark Structured Streaming 编程指南从编程模型到端到端 Exactly-Once 流处理【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/sparkStructured Streaming 是 Apache Spark 中构建在 Spark SQL 引擎之上的可扩展、容错的流处理引擎它允许开发者用与批处理完全相同的 Dataset/DataFrame API 来表达流式计算。本指南围绕仓库内 docs/streaming/index.md 及 docs/streaming/getting-started.md 等章节系统讲解其核心编程模型、快速示例、事件时间与水印、容错语义、输出模式与触发器等关键主题并结合仓库源码与示例给出可运行的实战细节帮助读者快速掌握这套无需关心流式细节的流处理方案。Structured Streaming 是什么Structured Streaming 是一个构建在 Spark SQL 引擎之上的可扩展、容错的流处理引擎。它的核心设计理念是你可以像对静态数据做批处理一样来表达流式计算——使用 Scala、Java、Python 或 R 中的 Dataset/DataFrame API 来表达流式聚合aggregation、事件时间窗口event-time windows、流批连接stream-to-batch joins等操作然后由 Spark SQL 引擎负责增量地、持续地执行这些计算并在数据不断到达时更新最终结果。其最关键的承诺是系统通过checkpointing检查点与Write-Ahead Logs预写日志WAL确保端到端 Exactly-Once的容错保证。用文档中的原话概括Structured Streaming 提供了快速、可扩展、容错、端到端 Exactly-Once 的流处理能力而用户无需自己操心流式的细节。两种三种执行模式内部实现上Structured Streaming 查询默认由微批处理micro-batch processing引擎执行它把数据流拆成一系列小批量的批作业batch job来运行从而获得低至100 毫秒的端到端延迟以及 Exactly-Once 容错保证。自 Spark 2.3 起引入了新的低延迟处理模式Continuous Processing连续处理可将端到端延迟降低到低至1 毫秒但只提供 At-Least-Once 保证自 Spark 4.1.0 起仓库中新增了Real-time Mode实时模式见 docs/streaming/real-time-mode.md面向毫秒级端到端延迟的运行型工作负载如欺诈检测、实时告警、实时个性化在保持 Exactly-Once 处理语义的同时支持无状态查询与部分有状态查询。无需修改查询中的 Dataset/DataFrame 操作你只需根据应用需求选择合适的执行模式。下文先以默认的微批处理模型为主线讲解编程模型与 API再讨论 Continuous Processing 与 Real-time Mode。快速上手流式词频统计让我们从最经典的示例——流式词频统计streaming word count开始。完整的可运行代码位于仓库的 examples 目录中Pythonexamples/src/main/python/sql/streaming/structured_network_wordcount.pyScalaexamples/src/main/scala/org/apache/spark/examples/sql/streaming/StructuredNetworkWordCount.scalaJavaexamples/src/main/java/org/apache/spark/examples/sql/streaming/JavaStructuredNetworkWordCount.javaRexamples/src/main/r/streaming/structured_network_wordcount.R第 1 步创建 SparkSession所有 Spark 功能的起点都是一个本地的SparkSessionfrom pyspark.sql import SparkSession from pyspark.sql.functions import explode from pyspark.sql.functions import split spark SparkSession \ .builder \ .appName(StructuredNetworkWordCount) \ .getOrCreate()Scala 版本与之等价注意引入import spark.implicits._Java 版本通过SparkSession.builder().appName(...).getOrCreate()创建R 版本则用sparkR.session(appName StructuredNetworkWordCount)。第 2 步创建流式 DataFrame 并计算词频从监听localhost:9999的服务器接收文本行并将其转换为流式 DataFrame再切词、分组计数# 表示来自 localhost:9999 连接的输入行流的 DataFrame lines spark \ .readStream \ .format(socket) \ .option(host, localhost) \ .option(port, 9999) \ .load() # 将行切分为单词 words lines.select( explode( split(lines.value, ) ).alias(word) ) # 生成运行中的词频统计 wordCounts words.groupBy(word).count()这段代码有几个关键点需要理解linesDataFrame 代表一张无界表unbounded table它包含一个名为value的字符串列流式文本数据中的每一行都成为这张表的一行。此时尚未真正接收任何数据——我们只是在搭建转换逻辑。代码使用两个内置 SQL 函数split把每行按空格切分为单词数组explode把数组展开为多行每行一个单词alias将新列命名为word。wordCounts通过对单词列分组计数得到它本身也是一个流式 DataFrame表示流上持续更新的词频结果。Scala 版本使用类型化的方式表达同样的逻辑val lines spark.readStream .format(socket) .option(host, localhost) .option(port, 9999) .load() // 转换为 Dataset[String] 后应用 flatMap 切词 val words lines.as[String].flatMap(_.split( )) val wordCounts words.groupBy(value).count()第 3 步启动查询设置输出模式为complete每次都把完整的最新结果表输出到控制台然后调用start()启动query wordCounts \ .writeStream \ .outputMode(complete) \ .format(console) \ .start() query.awaitTermination()query对象是活动流式查询的句柄awaitTermination()会阻塞主进程防止查询还在运行时进程就退出。Scala/Java 的写法一致wordCounts.writeStream.outputMode(complete).format(console).start()R 中对应write.stream(wordCounts, console, outputMode complete)与awaitTermination(query)。第 4 步运行示例先在一个终端中用 Netcat 启动数据服务器大多数类 Unix 系统都自带该工具$ nc -lk 9999然后在另一个终端运行示例# Python $ ./bin/spark-submit examples/src/main/python/sql/streaming/structured_network_wordcount.py localhost 9999 # Scala $ ./bin/run-example org.apache.spark.examples.sql.streaming.StructuredNetworkWordCount localhost 9999 # Java $ ./bin/run-example org.apache.spark.examples.sql.streaming.JavaStructuredNetworkWordCount localhost 9999 # R $ ./bin/spark-submit examples/src/main/r/streaming/structured_network_wordcount.R localhost 9999在 Netcat 终端输入apache spark和apache hadoop另一个终端会每秒打印更新后的计数------------------------------------------- Batch: 0 ------------------------------------------- ----------- | value|count| ----------- |apache| 1| | spark| 1| ----------- ------------------------------------------- Batch: 1 ------------------------------------------- ----------- | value|count| ----------- |apache| 2| | spark| 1| |hadoop| 1| -----------可以看到apache的计数从 1 累加到 2——这正是增量查询的直观体现Spark 持续检查 socket 连接上的新数据一旦有新数据就运行一次增量查询把之前的运行计数与新数据合并计算出更新后的计数。上述 Scala 完整实现的源码在 StructuredNetworkWordCount.scala 中其主流程与文档中的分段代码一一对应。编程模型把流当作持续追加的表Structured Streaming 的关键思想是把实时数据流看作一张被持续追加continuously appended的表。这带来了一个与批处理模型非常相似的流处理模型输入数据流 →Input Table输入表流上到达的每个数据项都相当于向输入表追加一行对输入执行的查询 →Result Table结果表每个触发间隔trigger interval比如每 1 秒新行被追加到输入表进而更新结果表每当结果表被更新就把发生变化的行写入外部存储sink——这就是Output输出。注意Structured Streaming 并不会物化整张表。它从流式数据源读取最新可用的数据增量地处理以更新结果然后丢弃源数据它只保留更新结果所必需的最小中间state状态数据例如上面例子中的中间计数。这个模型与许多其他流处理引擎有本质区别很多流系统要求用户自己维护运行中的聚合从而不得不操心容错和数据一致性at-least-once / at-most-once / exactly-once而在本模型中Spark 负责在有新数据时更新结果表用户无需再操心这些问题。输出模式Output Modes输出定义为写入外部存储的内容可以在以下三种模式中定义Complete Mode完整模式——每次把整个更新后的结果表写入外部存储。如何处理整张表的写入由存储连接器决定。Append Mode追加模式——只把自上次触发以来结果表中新增的行写入外部存储。仅适用于结果表中已有行不会发生变化的查询。Update Mode更新模式——只把自上次触发以来结果表中被更新的行写入外部存储自 Spark 2.1.1 起可用。它与 Complete Mode 的区别在于只输出发生变化的行如果查询不含聚合它等价于 Append 模式。每种模式只适用于特定类型的查询详细的兼容矩阵见 docs/streaming/apis-on-dataframes-and-datasets.md 的 Output Modes 一节例如含窗口聚合的查询支持 Append/Update/Complete而纯流连接查询仅支持 Append。处理事件时间与延迟数据事件时间event-time是数据本身携带的时间。对很多应用而言你可能希望基于事件时间操作比如统计 IoT 设备每分钟产生的事件数应当使用数据产生的时间即数据中的 event-time而不是 Spark 收到数据的时间。事件时间在该模型中表达得极为自然——每个设备事件是表中的一行事件时间就是行中的一个列值。这使得基于窗口的聚合例如每分钟的事件数只是对事件时间列的一种特殊分组聚合每个时间窗口是一个分组每行可以属于多个窗口。因此这类查询既可以在静态数据集例如从设备事件日志中收集的数据上定义也可以在数据流上定义且定义方式完全一致。更关键的是该模型天然处理晚到数据。由于 Spark 负责更新结果表它有完全的控制权既能用晚到的数据更新旧的聚合结果也能清理旧聚合以限制中间状态的大小。自 Spark 2.1 起支持watermarking水印允许用户指定晚到数据的阈值并让引擎据此清理旧状态。关于窗口操作的完整讲解见 docs/streaming/apis-on-dataframes-and-datasets.md 的 Window Operations on Event Time 一节。水印的基本用法在聚合之前调用withWatermark声明事件时间列与允许的延迟阈值words ... # schema 为 { timestamp: Timestamp, word: String } 的流式 DataFrame windowedCounts words \ .withWatermark(timestamp, 10 minutes) \ .groupBy( window(words.timestamp, 10 minutes, 5 minutes), words.word) \ .count()在这个例子中水印定义在timestamp列上允许数据延迟 10 分钟。引擎会为窗口T结束时刻维护状态并允许晚到数据更新该状态直到(引擎看到的最大事件时间 - 延迟阈值 T)。换言之阈值内的晚到数据会被聚合超过阈值的数据开始被丢弃。水印的语义保证docs/streaming/apis-on-dataframes-and-datasets.md是单向严格的设withWatermark的延迟为 2 hours则引擎保证绝不丢弃延迟小于 2 小时的数据——任何比已处理最新数据落后不到 2 小时按事件时间计的数据都保证被聚合但延迟超过 2 小时的数据并不保证被丢弃——它可能被聚合也可能不被聚合数据越延迟被处理的可能性越低。时间窗口的类型Spark 支持三种时间窗口详见 docs/streaming/apis-on-dataframes-and-datasets.mdTumbling滚动窗口一系列固定大小、不重叠、连续的时间区间一条输入只能属于一个窗口Sliding滑动窗口与滚动窗口同为固定大小但当滑动步长小于窗口时长时窗口会重叠此时一条输入可属于多个窗口滚动与滑动窗口都使用window函数Session会话窗口窗口长度随输入动态变化——会话窗口从一个输入开始若在 gap 时长内收到后续输入则自我扩展静态 gap 时长下超过 gap 时长没有新输入后窗口关闭。会话窗口使用session_window函数且支持基于输入行的动态 gap 时长表达式负值或零值 gap 的行会被过滤掉。容错语义端到端 Exactly-Once提供端到端 Exactly-Once 语义是 Structured Streaming 设计的关键目标之一。为实现这一点Structured Streaming 的源source、汇sink和执行引擎被设计为能可靠地跟踪处理的确切进度从而可以通过重启和/或重处理来应对任何类型的故障每个流式源都被假定有offset偏移量类似 Kafka offsets 或 Kinesis sequence numbers来跟踪流中的读取位置引擎使用checkpointing检查点和write-ahead logs预写日志记录每个触发周期内被处理数据的 offset 范围流式 sink 被设计为幂等的以处理重处理场景。将可重放的源与幂等的 sink结合Structured Streaming 可以在任何故障下确保端到端 Exactly-Once 语义。在 docs/streaming/apis-on-dataframes-and-datasets.md 中还有更具体的说明内置的File sink提供 Exactly-Once 保证Kafka sink与Foreach sink提供 At-Least-Once 保证详见 docs/streaming/structured-streaming-kafka-integration.mdConsole sink 与 Memory sink 不具备容错能力仅用于调试。通过 Checkpointing 从故障中恢复在故障或主动停机的情况下你可以通过checkpointing 和 write-ahead logs恢复之前查询的进度与状态并从中断处继续。配置方式是在DataStreamWriter中设置checkpointLocation选项指向HDFS 兼容文件系统中的一个目录aggDF .writeStream .outputMode(complete) .option(checkpointLocation, path/to/HDFS/dir) .format(memory) .start()查询会把所有进度信息每个触发周期处理的 offset 范围以及运行中的聚合结果例如快速示例中的词频保存到该检查点位置。重启后允许/不允许的查询变更从同一检查点位置重启时部分查询变更是被禁止的或语义未定义的docs/streaming/apis-on-dataframes-and-datasets.md输入源的数量或类型变更不允许有状态操作的变更聚合、去重、流流连接、mapGroupsWithState/flatMapGroupsWithState等有状态操作的 schema 变更分组键、聚合、去重列、连接列、状态类型等在重启之间不允许若确实需要支持状态 schema 演进可将复杂状态编码为字节如 Avro 编码自行实现迁移部分 sink 参数变更如文件 sink 输出目录不允许而 Kafka 输出 topic 变更、ForeachWriter代码变更等是允许的。不可修改的配置项docs/streaming/additional-information.md 明确提醒有若干配置在查询运行后不可修改如需更改必须丢弃检查点并启动新查询spark.sql.shuffle.partitions——因为状态是按 key 哈希分区的状态的分区数必须保持不变若想对有状态操作减少任务数可使用coalesce避免不必要的重新分区spark.sql.streaming.stateStore.providerClass——要正确读取查询先前的状态状态存储提供方类必须不变spark.sql.streaming.multipleWatermarkPolicy——修改会导致查询含多个水印时水印值不一致。输入源、输出汇与流式 API 速览内置输入源通过SparkSession.readStream()R 中为read.stream()返回的DataStreamReader创建流式 DataFrame。内置输入源完整选项表见 docs/streaming/apis-on-dataframes-and-datasets.md源关键选项容错说明File sourcepath、maxFilesPerTrigger、maxBytesPerTrigger、latestFirst、fileNameOnly、maxFileAge、cleanSourcearchive/delete/off等是按文件修改时间顺序读取目录中的新文件支持 text/CSV/JSON/ORC/Parquet文件必须原子地放入目录Kafka sourcekafka.bootstrap.servers、subscribe/subscribePattern/assign是兼容 Kafka broker 0.10.0 及以上版本见 docs/streaming/structured-streaming-kafka-integration.mdSocket source仅测试host、port否从 socket 连接读取 UTF8 文本监听 socket 在 driver 端不提供端到端容错Rate source仅测试rowsPerSecond、rampUpTime、numPartitions是按指定速率生成数据每行包含timestamp与从 0 开始的valueRate Per Micro-Batch source仅测试rowsPerBatch、numPartitions、startTimestamp、advanceMillisPerBatch是每个微批生成固定行数如批 0 产生 0~999、批 1 产生 1000~1999结果与查询执行情况无关流式 DataFrame 的常用操作流式 DataFrame/Dataset 支持大多数常规操作select、where、groupBy、map、filter、flatMap等也可以注册为临时视图后用 SQL 查询还能通过df.isStreaming判断是否流式数据。但部分操作不支持会抛出AnalysisException: ... is not supported with streaming DataFrames/Datasetslimit/取前 N 行、distinct不支持排序仅在聚合之后且 Complete 模式下支持部分外连接类型不支持见 docs/streaming/apis-on-dataframes-and-datasets.md 的 join 支持矩阵Update/Complete 模式下不允许链式多个有状态操作count()用ds.groupBy().count()代替、foreach()用writeStream.foreach(...)代替、show()用 console sink 代替等会立即执行的动作不适用。此外基于文件的源默认要求用户显式指定 schema以保证查询在故障时使用一致的 schema如需临时启用 schema 推断可设置spark.sql.streaming.schemaInferencetrue。分区发现/keyvalue/子目录同样受支持。内置输出汇汇支持输出模式容错说明File sinkAppend是Exactly-Once写目录path必填可设retentionTTLKafka sinkAppend/Update/Complete是At-Least-Once写一个或多个 topicForeach sinkAppend/Update/Complete是At-Least-Once通过open/process/close自定义逐行写入逻辑ForeachBatch sinkAppend/Update/Complete取决于实现对每个微批执行任意操作可用于复用批数据写入器、写多个位置Console sinkAppend/Update/Complete否每触发一次打印到控制台numRows默认 20、truncate默认 true仅限低数据量调试Memory sinkAppend/Complete否输出存为内存表表名即查询名需要调用start()才会真正启动查询执行其返回的StreamingQuery对象是持续运行执行的句柄可用于stop()、awaitTermination()、id()/runId()/name()、explain()、exception()、lastProgress/recentProgress/status等管理与监控操作。多个查询可在同一个 SparkSession 中并发运行可通过spark.streamsStreamingQueryManager统一管理也可以注册StreamingQueryListener异步接收查询启动、进度、终止的回调或通过设置spark.sql.streaming.metricsEnabledtrue将流式指标通过 Dropwizard 上报到 Ganglia/Graphite/JMX 等。触发器Trigger一览触发设置定义了流式数据处理的时机docs/streaming/apis-on-dataframes-and-datasets.md触发器说明未指定默认微批模式上一微批完成后立即生成下一微批固定间隔微批trigger(processingTime2 seconds)ScalaTrigger.ProcessingTime(2 seconds)若上一微批在间隔内完成则等待间隔结束若超时则上一微批一完成立即开始下一微批无新数据则不触发One-time 微批已弃用只执行一个微批处理所有可用数据后自行停止建议迁移到 Available-nowAvailable-now 微批处理所有可用数据后停止但会按源选项如maxFilesPerTrigger拆成多个微批扩展性更好、保证更强且会执行 no-data batch 推进水印Continuous实验性低延迟连续处理模式trigger(continuous1 second)Continuous Processing毫秒级低延迟模式Continuous Processing是 Spark 2.3 引入的实验性流执行模式详见 docs/streaming/performance-tips.md可实现约1 毫秒端到端延迟的 At-Least-Once 容错与默认微批引擎的 Exactly-Once 保证约 100ms 延迟上限形成对比。对某些查询类型无需修改 DataFrame/Dataset 操作即可切换执行模式——只需在写入端指定 continuous triggerimport org.apache.spark.sql.streaming.Trigger spark .readStream .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(subscribe, topic1) .load() .selectExpr(CAST(key AS STRING), CAST(value AS STRING)) .writeStream .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(topic, topic1) .trigger(Trigger.Continuous(1 second)) // 查询中唯一的改动 .start()Continuous 模式下 checkpoint interval此处 1 秒表示引擎每秒记录一次查询进度生成的检查点与微批引擎格式兼容因此查询可以在两种模式间相互切换重启注意切到 Continuous 即获得 At-Least-Once 保证。截至目前该模式的支持范围有限操作仅支持 map 类操作select、map、flatMap、mapPartitions等投影与where/filter选择全部 SQL 函数除聚合函数、current_timestamp()、current_date()外均支持源Kafka source全部选项、Rate source仅numPartitions与rowsPerSecond汇Kafka sink、Memory sink、Console sink。使用注意事项Continuous 引擎会为每个输入分区启动常驻任务因此启动前必须确保集群有足够核心让所有任务并行例如读取 10 个分区的 Kafka topic 至少需要 10 个核停止连续流可能产生无谓的任务终止警告可安全忽略目前没有失败任务自动重试机制。异步进度跟踪Async Progress Tracking同样是 docs/streaming/performance-tips.md 中的特性允许流式查询异步、并行地checkpoint 进度从而降低维护 offset log 与 commit log 带来的延迟。在微批模式下offset 管理直接落在处理关键路径上异步进度跟踪让查询不再被这些操作阻塞val query stream.writeStream .format(kafka) .option(topic, out) .option(checkpointLocation, /tmp/checkpoint) .option(asyncProgressTrackingEnabled, true) .start()选项取值默认说明asyncProgressTrackingEnabledtrue/falsefalse启用或禁用异步进度跟踪asyncProgressTrackingCheckpointIntervalMs毫秒1000提交 offset 与完成提交的间隔其限制包括初始版本仅支持使用 Kafka sink 的无状态查询且由于失败时批次 offset 范围可能变化异步进度跟踪不提供端到端 Exactly-Once。若在启用后想关闭它可能抛出java.lang.IllegalStateException: batch x doesnt exist解决办法是重新启用并将asyncProgressTrackingCheckpointIntervalMs设为 0运行至少两个微批后再安全关闭。Real-time Mode面向低延迟运行负载的实时模式自 Spark 4.1.0 起仓库引入了Real-time Modedocs/streaming/real-time-mode.md面向必须数据到达即响应的运行型工作负载如欺诈检测、实时告警、实时个性化。它同样通过writeStream.trigger()上的 Real-time trigger 启用查询其余部分不变import org.apache.spark.sql.streaming.Trigger spark .readStream .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(subscribe, input-topic) .load() .selectExpr(CAST(key AS STRING), CAST(value AS STRING)) .writeStream .format(kafka) .option(kafka.bootstrap.servers, host1:port1,host2:port2) .option(topic, output-topic) .option(checkpointLocation, /path/to/checkpoint) .outputMode(update) .trigger(Trigger.RealTime(5 minutes)) // 启用 Real-time Mode .start()最重要的一点传给 trigger 的时长默认 5 分钟主要是 checkpoint 间隔而不是输入驱动输出的延迟目标。记录可以持续处理并发出而不必等待批边界。与 Continuous Processing仍为实验性、仅 map 类操作、At-Least-Once相比Real-time Mode 复用 Spark 成熟的状态管理、Catalyst 优化器与 SQL 算子提供Exactly-Once 处理语义并支持无状态查询投影、过滤、union、流静态连接及自 Spark 4.3.0 起dropDuplicates、流式聚合与 JVMtransformWithState等有状态操作。启动前提包括输出模式必须为update、必须设置checkpointLocation、批时长至少为spark.sql.streaming.realTimeMode.minBatchDuration默认 5000ms等。仓库中的更多资源完整 API 指南输入源、窗口与水印、join 支持矩阵、去重、任意有状态操作、输出模式/汇/触发器、查询管理与监控、checkpoint 恢复语义docs/streaming/apis-on-dataframes-and-datasets.mdKafka 集成指南Kafka 源/汇的全部选项、部署与安全docs/streaming/structured-streaming-kafka-integration.mdtransformWithState高级有状态算子指南docs/streaming/structured-streaming-transform-with-state.md状态数据源实验性读取检查点中的状态存储docs/streaming/structured-streaming-state-data-source.md迁移指南从旧版 Structured Streaming / DStreams 迁移docs/streaming/ss-migration-guide.md更多可运行示例窗口聚合、会话窗口、Kafka 词频等examples/src/main/scala/org/apache/spark/examples/sql/streaming、examples/src/main/python/sql/streaming、examples/src/main/java/org/apache/spark/examples/sql/streaming、examples/src/main/r/streaming通用 Dataset/DataFrame 编程指南docs/sql-programming-guide.md流式查询指标监控配置docs/monitoring.md围绕 docs/streaming/index.md 展开的这套体系覆盖了从流即无界表的编程模型、快速上手的词频统计、事件时间与水印到 Exactly-Once 容错、四种触发器和三种执行模式的完整知识链路。无论你是要搭建实时告警、事件时间窗口聚合还是流批一体的数据处理管道Structured Streaming 都提供了像写批处理一样写流处理的路径而无需亲自处理底层的流式细节。【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址: https://gitcode.com/gh_mirrors/sp/spark创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表