Spark DataFrame文件级异常标记:含err_present则同文件其余行标bad_file
Spark DataFrame 批量标记文件异常记录
需求说明
当某个file_name对应的记录中存在至少一条err值为err_present时,该file_name下所有其他记录的err列需标记为bad_file;无异常记录的file_name则保持原err值(空字符串转为null以匹配示例输出)。
输入DataFrame示例
+-----------+---------+ |err |file_name| +-----------+---------+ |err_present|f1 | | |f1 | | |f1 | | |f2 | | |f2 | +-----------+---------+
期望输出DataFrame示例
+-----------+---------+ |err |file_name| +-----------+---------+ |err_present|f1 | |bad_file |f1 | |bad_file |f1 | | null |f2 | | null |f2 | +-----------+---------+
实现代码
from pyspark.sql import Window import pyspark.sql.functions as F # 初始化测试数据 df = spark.createDataFrame( [('err_present', 'f1'), ('', 'f1'), ('', 'f1'), ('', 'f2'), ('', 'f2')], ['err', 'file_name'] ) # 按file_name分组,标记该组是否存在err_present异常 window_spec = Window.partitionBy("file_name") df = df.withColumn( "has_exception", F.max(F.when(F.col("err") == "err_present", 1).otherwise(0)).over(window_spec) ) # 根据标记更新err列 result_df = df.withColumn( "err", F.when(F.col("err") == "err_present", F.col("err")) .when(F.col("has_exception") == 1, "bad_file") .when(F.col("err") == "", None) .otherwise(F.col("err")) ).drop("has_exception") # 打印结果 result_df.show()
代码说明
- 窗口分组标记:通过
Window.partitionBy("file_name")对每个文件的记录分组,用max函数判断该组是否存在err_present,生成has_exception辅助列。 - 条件更新err列:
- 保留原
err_present的记录不变; - 对于存在异常的文件组,其余记录统一标记为
bad_file; - 将原输入中的空字符串转为
null,与期望输出格式对齐;
- 保留原
- 清理辅助列:删除用于判断的
has_exception列,得到最终结果。
内容的提问来源于stack exchange,提问作者kites
相关产品推荐
相关产品推荐

