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

Spark Structured Streaming读取S3当日最新文件并按日分区写入求助

解决Spark Structured Streaming读取S3当日文件并按日分区写入的问题

针对问题1:仅读取当日新增文件

要避免读取S3上的历史文件,需从路径配置和流读取规则两方面调整:

  • 动态生成当日日期格式的输入路径,确保只监听当日的文件目录
  • 将startingPosition从earliest改为latest,让流任务启动时跳过已存在的历史文件,只处理后续新增的文件
  • 可选配置maxFilesPerTrigger,限制每次触发读取的文件数量,避免单次处理压力过大

针对问题2:按当日分区写入S3

如果数据本身带有day字段(格式为yyyy-MM-dd),直接用partitionBy("day")即可自动生成day=2022-08-26这类分区目录;如果数据没有该字段,需要先添加当日日期字段再分区:

修改后的完整代码

import org.apache.spark.sql.functions.current_date
import org.apache.spark.sql.streaming.Trigger

// 初始化SparkSession
val spark = SparkSession.builder().appName("raw_data").enableHiveSupport().getOrCreate()

// 动态生成当日输入路径(格式:s3://<path>/2022-08-26)
val today = java.time.LocalDate.now().toString
val inputPath = s"s3://<path>/$today"

// 读取流数据:仅监听当日路径,只处理新增文件
val df = spark.readStream
  .option("startingPosition", "latest") // 跳过历史文件,只读新增
  .option("maxFilesPerTrigger", 50) // 可选:每次触发最多读50个文件,按需调整
  .schema(LogSchema)
  .json(inputPath)

// 若数据无day字段,添加当日日期字段;已有则可跳过此步
val dfWithDay = df.withColumn("day", current_date().cast("string"))

// 按日分区写入S3
val query = dfWithDay
  .writeStream
  .outputMode("append")
  .partitionBy("day") // 按day字段分区,自动生成day=yyyy-MM-dd目录
  .format("parquet")
  .option("path", s"s3://<path>/raw_data/data/") // 根路径,分区目录会自动生成在下方
  .option("checkpointLocation", s"s3://<path>/raw_data/checkpoint/$today") // 按日单独设置checkpoint,避免跨日干扰
  .trigger(Trigger.ProcessingTime("300 seconds"))
  .start()

query.awaitTermination()

关键说明

  1. 动态路径:用java.time.LocalDate.now().toString自动生成当日日期字符串,确保输入路径始终指向当日目录
  2. latest启动位置:startingPosition="latest"会让流任务启动时忽略目录中已存在的文件,只处理后续新增的文件
  3. 按日checkpoint:将checkpoint目录按日期拆分,避免每日任务之间的状态干扰,也方便后续清理历史checkpoint
  4. 分区字段处理:如果数据中的day字段是从日志本身提取的(比如日志生成日期),就用原字段;如果没有则用current_date()生成当日日期,保证分区目录是当日的

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 02:54:37