如何用Spark-Scala逐行处理文件并按规则提取目标日期?
嘿,这个需求我之前做过类似的,用Spark-Scala可以分几步轻松实现,我给你拆解一下具体的思路和代码:
实现步骤拆解
1. 先识别所有月份列
首先我们需要从DataFrame的列中筛选出以MON开头的月份列(比如MON01、MON02),这些就是我们要处理的目标列。
import org.apache.spark.sql.functions._ // 假设你的原始DataFrame名为rawDf val monthColumns = rawDf.columns.filter(_.startsWith("MON"))
2. 把宽表转成窄表(便于统一处理)
为了避免逐个列处理的繁琐,我们可以先把原始的宽表(一行对应多个月份列)转换成窄表结构,每一行对应一个年份+月份列+月份文本值的组合:
val meltedDf = rawDf.select( col("year"), // 将所有月份列转换成struct数组,再explode拆分成行 explode(array( monthColumns.map(colName => struct( lit(colName).alias("month_col_name"), col(colName).alias("month_text") )): _* )).alias("month_info") ).select( col("year"), col("month_info.month_col_name"), col("month_info.month_text") )
3. 提取月份数字并处理日期
接下来我们从月份列名中提取两位的月份数字(比如从MON01提取01),然后拆分月份文本的每个字符,找到值为1的位置,转换成对应的日期:
val finalResultDf = meltedDf // 用正则提取月份列名末尾的两位数字 .withColumn("month", regexp_extract(col("month_col_name"), "MON(\\d{2})", 1)) // 将月份文本拆分成单个字符的数组,并用posexplode获取每个字符的位置和值 .select( col("year"), col("month"), posexplode(split(col("month_text"), "")).alias("char_index", "flag") ) // 只保留flag为"1"的记录 .filter(col("flag") === "1") // char_index从0开始,所以加1得到实际的天数 .withColumn("day", col("char_index") + 1) // 拼接成标准日期格式,并用to_date确保日期合法性(比如自动处理闰年2月的情况) .withColumn("full_date", to_date(concat_ws("-", col("year"), col("month"), col("day")))) // 可选:如果需要固定格式的字符串日期,用date_format // .withColumn("full_date_str", date_format(col("full_date"), "yyyy-MM-dd")) // 选择最终需要的列 .select("full_date")
额外说明
- 如果你的原始数据中还有其他需要保留的字段(比如用户ID、业务ID),只需要在
select的时候把这些字段带上即可,保证最终结果和原始行关联。 - 代码中用到的
posexplode会把每个字符的位置和值拆分成单独的行,正好对应我们需要的“每个1的位置对应一个日期”的需求。 to_date函数会自动处理日期的合法性,比如如果遇到2月30日这种非法日期(虽然你的需求里说文本长度等于当月天数,所以理论上不会出现),会返回null,你可以根据需要添加过滤逻辑。
如果你需要把每个原始行对应的所有日期收集成数组而不是拆分成多行,可以把posexplode换成transform+filter的组合,用collect_list收集结果,比如:
val arrayResultDf = meltedDf .withColumn("month", regexp_extract(col("month_col_name"), "MON(\\d{2})", 1)) .withColumn( "valid_days", // 遍历每个字符的位置,筛选出flag为1的位置+1作为天数 transform( filter( zip_with(split(col("month_text"), ""), sequence(lit(1), length(col("month_text"))), (flag, day) => struct(flag, day)), x => x.flag === "1" ), x => concat_ws("-", col("year"), col("month"), lpad(x.day, 2, "0")) ) ) .select("year", "month", "valid_days")
这样valid_days列就是每个月份对应的所有符合条件的日期数组。
内容的提问来源于stack exchange,提问作者Ravi Anand
相关产品推荐
相关产品推荐

