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

