基于其他列值数量更新DataFrame的ErrorCode列
问题背景
我的DataFrame会通过不同的methods更新,每次method都会向ErrorCode、ErrorColVal和ErrorColName列添加新值。初始DataFrame如下:
+--------+-----------------+-----------+--------------------+ |sequence| ErrorCode|ErrorColVal| ErrorColName| +--------+-----------------+-----------+--------------------+ | 1|CHAR0008-CHAR0021|X-SAHU |X-LASTNAME | | 2|CHAR0009-CHAR0021|SAMBIT-SAHU|FIRSTNAME-LASTNAME | | 3|CHAR0009 |SANTOSH |FIRSTNAME | | 5|CHAR0009 |SANTOSH |FIRSTNAME | | 8|CHAR0009-CHAR0021|SAMBIT-SAHU|FIRSTNAME-LASTNAME | +--------+-----------------+-----------+--------------------+
部分methods会出现仅向ErrorCode添加1个值,但向ErrorColVal和ErrorColName添加多个值的情况。比如sequence = 2时,ErrorColName新增了SURNAME和TITLE,ErrorColVal新增了SINGH和DR,但ErrorCode仅新增CHAR0022,导致更新后列值数量不匹配:
+--------+--------------------------+--------------------+--------------------------------+ |sequence| ErrorCode| ErrorColVal| ErrorColName| +--------+--------------------------+--------------------+--------------------------------+ | 1|CHAR0008-CHAR0021-CHAR0022|X-SAHU-SINGH |X-LASTNAME-SURNAME-TITLE | | 2|CHAR0009-CHAR0021-CHAR0022|SAMBIT-SAHU-SINGH-DR|FIRSTNAME-LASTNAME-SURNAME-TITLE| | 3|CHAR0009 |SANTOSH |FIRSTNAME | | 5|CHAR0009 |SANTOSH |FIRSTNAME | | 8|CHAR0009-CHAR0021-CHAR0022|SAMBIT-SAHU-BORA |FIRSTNAME-LASTNAME-SURNAME | +--------+--------------------------+--------------------+--------------------------------+
这种不匹配会导致后续无法关联ErrorCode=CHAR0022对应的ErrorColVal和ErrorColName。需要通过重复指定错误码(存储在ErrCd字段,可通过lit("ErrCd")获取)来修正,确保三列值数量一致,最终目标DataFrame如下:
+--------+-----------------------------------+--------------------+--------------------------------+ |sequence| ErrorCode| ErrorColVal| ErrorColName| +--------+-----------------------------------+--------------------+--------------------------------+ | 1|CHAR0008-CHAR0021-CHAR0022 |X-SAHU-SINGH |X-LASTNAME-SURNAME-TITLE | | 2|CHAR0009-CHAR0021-CHAR0022-CHAR0022|SAMBIT-SAHU-SINGH-DR|FIRSTNAME-LASTNAME-SURNAME-TITLE| | 3|CHAR0009 |SANTOSH |FIRSTNAME | | 5|CHAR0009 |SANTOSH |FIRSTNAME | | 8|CHAR0009-CHAR0021-CHAR0022 |SAMBIT-SAHU-BORA |FIRSTNAME-LASTNAME-SURNAME | +--------+-----------------------------------+--------------------+--------------------------------+
实现方案
以下是基于PySpark的可复用方法,核心逻辑是计算三列的元素数量差,动态补充重复的错误码:
from pyspark.sql import functions as F from pyspark.sql.types import StringType def align_error_columns(df, err_cd_col="ErrCd"): # 计算各列按'-'分割后的元素数量 df = df.withColumn("code_count", F.size(F.split(F.col("ErrorCode"), "-"))) \ .withColumn("val_count", F.size(F.split(F.col("ErrorColVal"), "-"))) \ .withColumn("name_count", F.size(F.split(F.col("ErrorColName"), "-"))) # 取val_count和name_count的最大值作为目标数量(确保两者一致) df = df.withColumn("target_count", F.greatest(F.col("val_count"), F.col("name_count"))) # 计算需要补充的错误码数量 df = df.withColumn("need_add", F.when(F.col("target_count") > F.col("code_count"), F.col("target_count") - F.col("code_count")) .otherwise(0)) # 生成需要重复的错误码字符串,例如需要补充2个则生成"CHAR0022-CHAR0022" df = df.withColumn("add_codes", F.when(F.col("need_add") > 0, F.concat_ws("-", F.expr(f"array_repeat({err_cd_col}, need_add)"))) .otherwise("")) # 拼接原ErrorCode和补充的错误码,处理空补充字符串的情况 df = df.withColumn("ErrorCode", F.when(F.col("add_codes") != "", F.concat(F.col("ErrorCode"), F.lit("-"), F.col("add_codes"))) .otherwise(F.col("ErrorCode"))) # 清理临时列 return df.drop("code_count", "val_count", "name_count", "target_count", "need_add", "add_codes")
使用说明
- 确保DataFrame中包含
ErrCd字段(存储需要重复的错误码,如CHAR0022) - 调用方法时直接传入DataFrame:
# 假设不匹配的DataFrame名为df_unaligned df_aligned = align_error_columns(df_unaligned) df_aligned.show(truncate=False)
逻辑说明
- 通过
split和size计算三列的元素数量 - 以
ErrorColVal和ErrorColName的最大数量为基准,计算ErrorCode需要补充的数量 - 用
array_repeat生成对应数量的错误码数组,再用concat_ws拼接成字符串 - 将补充字符串拼接到原
ErrorCode后,确保三列元素数量完全匹配
内容的提问来源于stack exchange,提问作者SDS
相关产品推荐
相关产品推荐

