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

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()

代码说明

  1. 窗口分组标记:通过Window.partitionBy("file_name")对每个文件的记录分组,用max函数判断该组是否存在err_present,生成has_exception辅助列。
  2. 条件更新err列:
    • 保留原err_present的记录不变;
    • 对于存在异常的文件组,其余记录统一标记为bad_file;
    • 将原输入中的空字符串转为null,与期望输出格式对齐;
  3. 清理辅助列:删除用于判断的has_exception列,得到最终结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 00:16:07