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()
关键说明
- 动态路径:用
java.time.LocalDate.now().toString自动生成当日日期字符串,确保输入路径始终指向当日目录 - latest启动位置:
startingPosition="latest"会让流任务启动时忽略目录中已存在的文件,只处理后续新增的文件 - 按日checkpoint:将checkpoint目录按日期拆分,避免每日任务之间的状态干扰,也方便后续清理历史checkpoint
- 分区字段处理:如果数据中的
day字段是从日志本身提取的(比如日志生成日期),就用原字段;如果没有则用current_date()生成当日日期,保证分区目录是当日的
内容的提问来源于stack exchange,提问作者Etisha
相关产品推荐
相关产品推荐

