You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Streaming如何监控CSV文件新增行并输出至控制台?

如何让Spark Streaming监听CSV文件的内容追加(而非新增文件)

你的问题我太熟悉了——Spark Structured Streaming的文件数据源天生是监听新增文件的,不会盯着已有文件的内容变化。所以你往unit_testing.csv里追加行时,Spark根本没察觉到,自然不会输出新内容。只有当你往目录里丢一个新的CSV文件时,它才会触发读取。

下面给你两种可行的解决方案,看哪种更贴合你的需求:


方案1:用Socket流+tail命令实时转发新增行

这是最快捷的方式,不需要改动原文件的写入逻辑,只需要借助tail和netcat把新增行转发给Spark的socket流。

步骤1:修改Spark代码监听Socket

把原来的文件流改成监听本地socket端口,同时解析CSV行:

from pyspark.sql import SparkSession
from pyspark.sql.types import (
    StringType, IntegerType, StructType, StructField, BooleanType
)

spark = (
    SparkSession.builder.appName("csv-append-stream")
    .master("local[*]")
    .getOrCreate()
)

# 保持原Schema不变
schema = StructType(
    [
        StructField("Drug_Name", StringType(), True),
        StructField("Count", IntegerType(), True),
        StructField("Faulty", BooleanType(), True),
    ]
)

# 监听本地9999端口的socket流
df = (
    spark.readStream
    .format("socket")
    .option("host", "localhost")
    .option("port", 9999)
    .load()
    # 按分号分割每行字符串,映射到Schema
    .selectExpr("split(value, ';') as parts")
    .select(
        df["parts"][0].cast(StringType()).alias("Drug_Name"),
        df["parts"][1].cast(IntegerType()).alias("Count"),
        df["parts"][2].cast(BooleanType()).alias("Faulty")
    )
)

# 输出到控制台(如果要输出到CSV,看后面的补充)
query = df.writeStream
    .format("console")
    .outputMode("append")
    .queryName("csv-append-stream")
    .start()

query.awaitTermination()

步骤2:用tail+nc转发新增行

打开另一个终端,执行以下命令(确保你的系统有nc(netcat)工具,Linux/macOS默认自带,Windows可以用WSL或者类似工具):

# 只转发之后追加的行(-n 0),实时跟踪文件变化(-F),发送到本地9999端口
tail -n 0 -F csv_files/unit_testing.csv | nc -lk localhost 9999

现在,当你往unit_testing.csv里追加新行时,tail会立刻捕获到,通过socket发送给Spark,控制台就会实时输出新增的行了!


方案2:改用“小文件写入”模式

如果你能控制原CSV的写入逻辑,可以每次新增行时写入一个新的小文件(比如按时间戳命名,比如unit_testing_202405201205.csv),放在csv_files目录下。这样Spark的文件流就能自动检测到新文件,读取其中的内容。

这种方式不需要额外工具,修改写入逻辑即可,适合批量新增行的场景。比如每次5分钟把新增行写入一个新文件,Spark会自动读取并输出。


补充:输出到CSV文件

如果要把新增行输出到另一个CSV文件,只需要修改writeStream的配置:

query = df.writeStream
    .format("csv")
    .option("path", "output_csv")  # 输出目录
    .option("checkpointLocation", "checkpoint_dir")  # 必须设置检查点,用于恢复流状态
    .option("sep", ";")  # 保持和原文件一致的分隔符
    .outputMode("append")
    .start()

注意:检查点目录checkpoint_dir需要是一个空目录或者之前用过的检查点目录,Spark会用它来跟踪已经处理过的数据,避免重复输出。


内容的提问来源于stack exchange,提问作者beeteedubs

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.09 08:47:53