用 Spark Structured Streaming 持续统计输入文本

流式词频统计不必写成一套手工维护计数器的程序。Structured Streaming 把不断到达的数据看成持续追加行的输入表:先描述拆词与分组计数,再让 Spark 在新数据到来时增量更新结果。这个例子连接本地 TCP 文本服务,按行接收 UTF-8 文本,把累计词频输出到控制台。

本文完整翻译整理 Apache Spark 官方入门页的快速示例、编程模型、事件时间与容错说明,以 Python 路径展开;源页还提供 Scala、Java 与 R 的对应实现链接。2026 年 10 月 5 日读取的 latest 文档标题为 Spark 4.2.0,伴随脚本链接也指向 v4.2.0。本文没有安装 Spark、启动 Netcat 或运行示例,所有输出只作为文档示意。

Structured Streaming 词频模型:TCP 输入行拆分成词,增量更新累计结果表,Complete 输出全表,Update 输出改变的行。
输入表、增量状态与输出模式。未完纪据 Apache Spark 官方模型绘制,不是运行截图。

先创建 SparkSession,再描述查询

SparkSession 是 Spark 功能的入口。通过 builder 指定应用名并调用 getOrCreate(),就能取得会话。原文片段没有硬编码执行 master;实际使用本地或集群资源取决于启动配置。

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, split

spark = (
    SparkSession.builder
    .appName("StructuredNetworkWordCount")
    .getOrCreate()
)

接着从 localhost:9999 读取 socket 数据,构造流式 DataFrame。它有一列名为 value 的字符串,每收到一行文本,就相当于向输入表追加一行。

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

split 把一行按空格切成数组,explode 把数组的每个元素变成独立一行,alias("word") 给结果列命名。最后按单词分组并计数,得到代表累计词频的流式结果表。此时只是定义变换,还没有开始接收和计算数据。

编者补充:这里是最小拆词规则,没有统一大小写、去标点或处理自然语言分词。多个空格、空行等输入也应按所用 split 语义评估。不要把它描述成通用中文分词或经过清洗的日志分析方案。

启动查询,并保持进程存活

通过 writeStream 配置输出。此例选择 complete,每次结果更新时向 console sink 写出当前完整结果表:

query = (
    wordCounts.writeStream
    .outputMode("complete")
    .format("console")
    .start()
)
query.awaitTermination()

start() 才真正启动后台流式计算;返回的 query 是正在运行的查询句柄。awaitTermination() 等待查询结束,防止主进程在查询活跃时直接退出。这不等于为程序增加重启恢复能力。

完整 Python 伴随脚本

下面保留官方 v4.2.0 脚本的执行逻辑。脚本读取主机和端口两个参数,参数数量错误时打印用法并退出;端口通过 int() 转换,没有额外实现友好的范围校验。示例来自 structured_network_wordcount.py,ASF 版权及许可见文末。

import sys
from pyspark.sql import SparkSession
from pyspark.sql.functions import explode
from pyspark.sql.functions import split

if __name__ == "__main__":
    if len(sys.argv) != 3:
        print("Usage: structured_network_wordcount.py <hostname> <port>", file=sys.stderr)
        sys.exit(-1)

    host = sys.argv[1]
    port = int(sys.argv[2])

    spark = SparkSession\
        .builder\
        .appName("StructuredNetworkWordCount")\
        .getOrCreate()
    # Create DataFrame representing the stream of input lines from connection to host:port
    lines = spark\
        .readStream\
        .format('socket')\
        .option('host', host)\
        .option('port', port)\
        .load()

    # Split the lines into words
    words = lines.select(
        # explode turns each item in an array into a separate row
        explode(
            split(lines.value, ' ')
        ).alias('word')
    )
    # Generate running word count
    wordCounts = words.groupBy('word').count()

    # Start running the query that prints the running counts to the console
    query = wordCounts\
        .writeStream\
        .outputMode('complete')\
        .format('console')\
        .start()

    query.awaitTermination()

该代码归 Apache Software Foundation 及贡献者,按 Apache License 2.0 提供;文末保留许可证全文与上游 NOTICE。正文代码不含执行结果。

本地演示的两个终端

原文在第一个终端使用 Netcat 作为数据服务:

nc -lk 9999

在 Spark 分发目录的第二个终端启动随附的 Python 示例:

./bin/spark-submit examples/src/main/python/sql/streaming/structured_network_wordcount.py localhost 9999

操作边界:这些是原文面向 Unix 类环境的命令。Netcat 的实现和选项在不同系统可能不同,nc -lk 9999 也可能监听多个网络接口;演示前须按本机实现将服务限制到回环接口或隔离网络,不向公网暴露这个无认证文本端口。若需强制 Spark 使用本机资源,可在确认版本后为 spark-submit 配置本地 master;这属于部署选择,不是原文已经验证的配置。Windows 用户不能假定当前系统已经提供 nc 或相同 shell 路径。

在 Netcat 终端依次输入以下两行:

apache spark
apache hadoop

按原文分成两批处理的示意,第一批累计结果是 apache=1、spark=1;第二批是 apache=2、spark=1、hadoop=1。实际批次编号、批次边界、显示顺序与时间取决于运行过程,两行也可能进入同一批。原文用“一秒”帮助说明刷新过程,但脚本没有显式设置一秒 trigger,不能据此保证严格每秒输出。

示意状态 word count
处理第一行后 apache 1
处理第一行后 spark 1
处理第二行后 apache 2
处理第二行后 spark 1
处理第二行后 hadoop 1

原文展示差异:官方页面的控制台示意表头写成 value,而本次核对的 Python 脚本明确把列命名为 word,并按它分组。本稿依照该 Python 代码使用 word;上表是逻辑演示,不是我们执行后的屏幕输出。

“无限输入表”是模型,不是把所有历史输入留在内存

Structured Streaming 的核心是把实时数据流视为不断追加的输入表。查询定义仍类似批处理中的静态表查询,Spark 则把它作为增量查询执行。每当触发一次处理,就把新行纳入计算,并更新结果表;输出到外部系统的是按输出模式选定的结果。

在词频示例中,lines 是输入表,wordCounts 是结果表。Spark 接收新行后,把新的词频增量与已有累计值合并。它不会把整个无限输入表全部物化;处理后的原始输入可以丢弃,只保留更新结果所需的中间状态,例如每个词的累计次数。

编者补充:“最少必要状态”不等于有界内存。这个无窗口累计词频查询需要保留所有已出现词的计数,词种持续增长时,状态也会增长。它适合说明模型,不能据此声称可以无限量处理任意输入而无需状态治理。

Complete、Append 与 Update 分别写出什么

输出模式 写出内容 适用边界
Complete 本次更新后的整个结果表 接收端决定如何处理整表输出;本例使用这一模式。
Append 上次触发后新追加到结果表的行 只适用于已有结果行不会再变化的查询;不能直接用于本例不断改变的累计计数。
Update 上次触发后发生改变的结果行 从 Spark 2.1.1 开始提供;无聚合时与 Append 等价。

输出模式受查询类型和 sink 支持情况共同约束。例中第二行到来时,apache 的计数改变、hadoop 是新词,而 spark 没有改变:Complete 会输出三者,Update 关注新增或更新行。这里没有把代码改成另一种模式并声称实测成功。

事件时间与迟到数据

事件时间是数据本身携带的发生时间,不是 Spark 收到它的时间。统计物联网设备每分钟产生的事件数时,通常希望按设备事件发生时间分组。在表模型中,事件时间就是一列,时间窗口聚合便是围绕这一列的特殊分组;滑动窗口中,一行可能属于多个窗口。

同一类查询可以作用于已经收集好的静态日志,也可以作用于实时流。数据即使晚于预期到达,Spark 仍可以更新旧聚合。为了控制中间状态,需要明确什么时候可以清理旧窗口。原文指出 Spark 从 2.1 开始支持 watermark,用于声明迟到阈值并帮助清理状态;具体窗口与水位线规则应继续阅读官方窗口操作章节。本例没有事件时间字段、窗口或 watermark,不能把这些能力当作已经启用。

容错语义取决于源、引擎与接收端共同配合

Structured Streaming 的设计目标之一是端到端 exactly-once 语义。原文描述的前提是:源能够用 offset 等位置标记跟踪读取进度,并支持重放;引擎通过 checkpoint 和预写日志记录每次触发所处理的范围;sink 具备幂等处理重放的能力。这些条件配合,才能在重启或重处理之后维持目标语义。

本例不具备这些生产前提。socket 源不提供可重放的持久输入,console sink 用于观察,不是持久业务输出,脚本也没有显式设置持久 checkpoint 位置。不能因为看到累计结果,就宣称已验证端到端 exactly-once 恢复。若转向生产,需要选择匹配的源与 sink、可靠检查点存储和恢复策略,并进行失败恢复测试。

静态审核结论与来源

本次只逐行检查代码与命令:没有发现硬编码秘密,也没有发现把网络文本拼成 shell 命令或动态执行代码的路径;并不等于不存在其他漏洞。现实风险包括无认证的演示端口、任意指定连接目标、恶意超大输入或高基数词导致资源压力,以及控制台可能暴露输入内容。运行应限于隔离演示数据,不能接入生产日志或敏感文本。

原文:Structured Streaming Programming Guide — Getting Started。伴随脚本:v4.2.0 Python 示例。作者归属:Apache Software Foundation 及 Spark 贡献者;未完纪中文翻译整理。正文调整了结构、统一 Python 列名,并增加已标识的版本与静态审核说明。

代码许可:Apache License 2.0;上游 NOTICE 声明 Apache Spark,Copyright 2014 and onwards The Apache Software Foundation。相关许可与 NOTICE 见文末;原始代码头见链接源码。原创图注明其独立绘制性质。

Apache Spark 示例代码许可及 NOTICE

Apache License
                           Version 2.0, January 2004
                        http://www.apache.org/licenses/

   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION

   1. Definitions.

      "License" shall mean the terms and conditions for use, reproduction,
      and distribution as defined by Sections 1 through 9 of this document.

      "Licensor" shall mean the copyright owner or entity authorized by
      the copyright owner that is granting the License.
      "Legal Entity" shall mean the union of the acting entity and all
      other entities that control, are controlled by, or are under common
      control with that entity. For the purposes of this definition,
      "control" means (i) the power, direct or indirect, to cause the
      direction or management of such entity, whether by contract or
      otherwise, or (ii) ownership of fifty percent (50%) or more of the
      outstanding shares, or (iii) beneficial ownership of such entity.
      "You" (or "Your") shall mean an individual or Legal Entity
      exercising permissions granted by this License.

      "Source" form shall mean the preferred form for making modifications,
      including but not limited to software source code, documentation
      source, and configuration files.
      "Object" form shall mean any form resulting from mechanical
      transformation or translation of a Source form, including but
      not limited to compiled object code, generated documentation,
      and conversions to other media types.

      "Work" shall mean the work of authorship, whether in Source or
      Object form, made available under the License, as indicated by a
      copyright notice that is included in or attached to the work
      (an example is provided in the Appendix below).
      "Derivative Works" shall mean any work, whether in Source or Object
      form, that is based on (or derived from) the Work and for which the
      editorial revisions, annotations, elaborations, or other modifications
      represent, as a whole, an original work of authorship. For the purposes
      of this License, Derivative Works shall not include works that remain
      separable from, or merely link (or bind by name) to the interfaces of,
      the Work and Derivative Works thereof.
      "Contribution" shall mean any work of authorship, including
      the original version of the Work and any modifications or additions
      to that Work or Derivative Works thereof, that is intentionally
      submitted to Licensor for inclusion in the Work by the copyright owner
      or by an individual or Legal Entity authorized to submit on behalf of
      the copyright owner. For the purposes of this definition, "submitted"
      means any form of electronic, verbal, or written communication sent
      to the Licensor or its representatives, including but not limited to
      communication on electronic mailing lists, source code control systems,
      and issue tracking systems that are managed by, or on behalf of, the
      Licensor for the purpose of discussing and improving the Work, but
      excluding communication that is conspicuously marked or otherwise
      designated in writing by the copyright owner as "Not a Contribution."
      "Contributor" shall mean Licensor and any individual or Legal Entity
      on behalf of whom a Contribution has been received by Licensor and
      subsequently incorporated within the Work.
   2. Grant of Copyright License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      copyright license to reproduce, prepare Derivative Works of,
      publicly display, publicly perform, sublicense, and distribute the
      Work and such Derivative Works in Source or Object form.
   3. Grant of Patent License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      (except as stated in this section) patent license to make, have made,
      use, offer to sell, sell, import, and otherwise transfer the Work,
      where such license applies only to those patent claims licensable
      by such Contributor that are necessarily infringed by their
      Contribution(s) alone or by combination of their Contribution(s)
      with the Work to which such Contribution(s) was submitted. If You
      institute patent litigation against any entity (including a
      cross-claim or counterclaim in a lawsuit) alleging that the Work
      or a Contribution incorporated within the Work constitutes direct
      or contributory patent infringement, then any patent licenses
      granted to You under this License for that Work shall terminate
      as of the date such litigation is filed.
   4. Redistribution. You may reproduce and distribute copies of the
      Work or Derivative Works thereof in any medium, with or without
      modifications, and in Source or Object form, provided that You
      meet the following conditions:

      (a) You must give any other recipients of the Work or
          Derivative Works a copy of this License; and

      (b) You must cause any modified files to carry prominent notices
          stating that You changed the files; and
      (c) You must retain, in the Source form of any Derivative Works
          that You distribute, all copyright, patent, trademark, and
          attribution notices from the Source form of the Work,
          excluding those notices that do not pertain to any part of
          the Derivative Works; and
      (d) If the Work includes a "NOTICE" text file as part of its
          distribution, then any Derivative Works that You distribute must
          include a readable copy of the attribution notices contained
          within such NOTICE file, excluding those notices that do not
          pertain to any part of the Derivative Works, in at least one
          of the following places: within a NOTICE text file distributed
          as part of the Derivative Works; within the Source form or
          documentation, if provided along with the Derivative Works; or,
          within a display generated by the Derivative Works, if and
          wherever such third-party notices normally appear. The contents
          of the NOTICE file are for informational purposes only and
          do not modify the License. You may add Your own attribution
          notices within Derivative Works that You distribute, alongside
          or as an addendum to the NOTICE text from the Work, provided
          that such additional attribution notices cannot be construed
          as modifying the License.
      You may add Your own copyright statement to Your modifications and
      may provide additional or different license terms and conditions
      for use, reproduction, or distribution of Your modifications, or
      for any such Derivative Works as a whole, provided Your use,
      reproduction, and distribution of the Work otherwise complies with
      the conditions stated in this License.
   5. Submission of Contributions. Unless You explicitly state otherwise,
      any Contribution intentionally submitted for inclusion in the Work
      by You to the Licensor shall be under the terms and conditions of
      this License, without any additional terms or conditions.
      Notwithstanding the above, nothing herein shall supersede or modify
      the terms of any separate license agreement you may have executed
      with Licensor regarding such Contributions.
   6. Trademarks. This License does not grant permission to use the trade
      names, trademarks, service marks, or product names of the Licensor,
      except as required for reasonable and customary use in describing the
      origin of the Work and reproducing the content of the NOTICE file.
   7. Disclaimer of Warranty. Unless required by applicable law or
      agreed to in writing, Licensor provides the Work (and each
      Contributor provides its Contributions) on an "AS IS" BASIS,
      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
      implied, including, without limitation, any warranties or conditions
      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
      PARTICULAR PURPOSE. You are solely responsible for determining the
      appropriateness of using or redistributing the Work and assume any
      risks associated with Your exercise of permissions under this License.
   8. Limitation of Liability. In no event and under no legal theory,
      whether in tort (including negligence), contract, or otherwise,
      unless required by applicable law (such as deliberate and grossly
      negligent acts) or agreed to in writing, shall any Contributor be
      liable to You for damages, including any direct, indirect, special,
      incidental, or consequential damages of any character arising as a
      result of this License or out of the use or inability to use the
      Work (including but not limited to damages for loss of goodwill,
      work stoppage, computer failure or malfunction, or any and all
      other commercial damages or losses), even if such Contributor
      has been advised of the possibility of such damages.
   9. Accepting Warranty or Additional Liability. While redistributing
      the Work or Derivative Works thereof, You may choose to offer,
      and charge a fee for, acceptance of support, warranty, indemnity,
      or other liability obligations and/or rights consistent with this
      License. However, in accepting such obligations, You may act only
      on Your own behalf and on Your sole responsibility, not on behalf
      of any other Contributor, and only if You agree to indemnify,
      defend, and hold each Contributor harmless for any liability
      incurred by, or claims asserted against, such Contributor by reason
      of your accepting any such warranty or additional liability.
   END OF TERMS AND CONDITIONS

   APPENDIX: How to apply the Apache License to your work.
      To apply the Apache License to your work, attach the following
      boilerplate notice, with the fields enclosed by brackets "[]"
      replaced with your own identifying information. (Don't include
      the brackets!)  The text should be enclosed in the appropriate
      comment syntax for the file format. We also recommend that a
      file or class name and description of purpose be included on the
      same "printed page" as the copyright notice for easier
      identification within third-party archives.
   Copyright [yyyy] [name of copyright owner]

   Licensed under the Apache License, Version 2.0 (the "License");
   you may not use this file except in compliance with the License.
   You may obtain a copy of the License at

       http://www.apache.org/licenses/LICENSE-2.0
   Unless required by applicable law or agreed to in writing, software
   distributed under the License is distributed on an "AS IS" BASIS,
   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
   See the License for the specific language governing permissions and
   limitations under the License.
Apache Spark
Copyright 2014 and onwards The Apache Software Foundation.

This product includes software developed at
The Apache Software Foundation (http://www.apache.org/).


Export Control Notice
---------------------
This distribution includes cryptographic software. The country in which you currently reside may have
restrictions on the import, possession, use, and/or re-export to another country, of encryption software.
BEFORE using any encryption software, please check your country's laws, regulations and policies concerning
the import, possession, or use, and re-export of encryption software, to see if this is permitted. See
<http://www.wassenaar.org/> for more information.
The U.S. Government Department of Commerce, Bureau of Industry and Security (BIS), has classified this
software as Export Commodity Control Number (ECCN) 5D002.C.1, which includes information security software
using or performing cryptographic functions with asymmetric algorithms. The form and manner of this Apache
Software Foundation distribution makes it eligible for export under the License Exception ENC Technology
Software Unrestricted (TSU) exception (see the BIS Export Administration Regulations, Section 740.13) for
both object code and source code.
The following provides more details on the included cryptographic software:

This software uses Apache Commons Crypto (https://commons.apache.org/proper/commons-crypto/) to
support authentication, and encryption and decryption of data sent across the network between
services.


Metrics
Copyright 2010-2013 Coda Hale and Yammer, Inc.

This product includes software developed by Coda Hale and Yammer, Inc.
This product includes code derived from the JSR-166 project (ThreadLocalRandom, Striped64,
LongAdder), which was released with the following comments:

    Written by Doug Lea with assistance from members of JCP JSR-166
    Expert Group and released to the public domain, as explained at
    http://creativecommons.org/publicdomain/zero/1.0/
© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容