Spark Scala实现读取文件更新日期列并覆写保存
解决Spark读取文本文件更新日期列并覆写的问题
问题背景
我有一个文本文件,第一列是表名,第二列是日期,列之间用空格分隔。比如employee.txt里的内容是:
organization 4-15-2018
employee 5-15-2018
需求是读取这个文件,根据业务逻辑更新日期列后,保存并覆写原文件。以下是我目前的代码片段:
object Employee { def main(args: Array[String]) { val conf = new SparkConf().setMaster("local")... } }
完整解决方案
下面是可运行的完整Scala代码,涵盖读取文件、日期处理、覆写文件的全流程,你可以根据实际业务需求调整日期处理逻辑:
import org.apache.spark.sql.SparkSession import java.time.LocalDate import java.time.format.DateTimeFormatter import java.nio.file.{Files, Paths} object Employee { def main(args: Array[String]): Unit = { // 初始化SparkSession(推荐用这个替代旧的SparkConf+SparkContext,更简洁) val spark = SparkSession.builder() .master("local[*]") // local模式下使用所有可用核心 .appName("DateUpdateJob") .getOrCreate() import spark.implicits._ // 定义目标文件路径 val filePath = "employee.txt" // 读取文本文件并拆分列:按空格分割为表名和日期两部分 val rawData = spark.read.textFile(filePath) .map(line => { // 用split(" ", 2)确保即使日期含空格(示例中没有但预留场景)也能正确拆分 val parts = line.split(" ", 2) (parts(0), parts(1)) }) .toDF("table_name", "original_date") // 定义日期格式化器:匹配输入的M-d-yyyy格式,输出格式可按需修改 val inputFormatter = DateTimeFormatter.ofPattern("M-d-yyyy") val outputFormatter = DateTimeFormatter.ofPattern("M-d-yyyy") // 核心日期处理逻辑:这里示例是把日期往后推30天,替换成你的业务逻辑即可 val updatedData = rawData.map(row => { val tableName = row.getAs[String]("table_name") val originalDateStr = row.getAs[String]("original_date") val originalDate = LocalDate.parse(originalDateStr, inputFormatter) val updatedDate = originalDate.plusDays(30) // 替换成你的业务处理逻辑 s"$tableName ${updatedDate.format(outputFormatter)}" }) // 处理本地文件覆写:Spark本地模式默认不支持直接覆写,先删原文件再处理 val targetPath = Paths.get(filePath) if (Files.exists(targetPath)) { Files.delete(targetPath) } // 先写入临时目录,避免直接覆写的问题 updatedData.coalesce(1) // 合并为一个分区,避免生成多个part文件 .write .mode("overwrite") .text(filePath + "_temp") // 把临时目录里的结果文件移到原路径,再清理临时目录 val tempDir = Paths.get(filePath + "_temp") val partFile = Files.list(tempDir).filter(p => p.getFileName.toString.startsWith("part-")).findFirst().get Files.move(partFile, targetPath) // 递归删除临时目录 Files.walk(tempDir).sorted(java.util.Comparator.reverseOrder()).forEach(Files.delete) // 停止Spark会话 spark.stop() } }
关键细节说明
- 日期处理:用Java 8的
LocalDate和DateTimeFormatter替代旧的java.util.Date,更安全且API更友好 - 文件覆写:本地文件系统中Spark无法直接覆写单个文件,所以采用临时目录中转的方式;如果是HDFS等分布式存储,直接用
.mode("overwrite")即可 - 列拆分:
split(" ", 2)保证只拆分为表名和日期两部分,避免日期字段意外拆分的问题 - 合并分区:
coalesce(1)将数据合并到一个分区,确保最终生成单个文件,方便后续替换原文件
内容的提问来源于stack exchange,提问作者shashank
相关产品推荐
相关产品推荐

