流式词频统计不必写成一套手工维护计数器的程序。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 或运行示例,所有输出只作为文档示意。

先创建 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/












暂无评论内容