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

验证DataFrame中id列对与value匹配性并定位异常行

简洁方案:检查ID对匹配及Value一致性(Pandas/PySpark)

示例数据

df_1(基准数据集)

# Pandas 版本
import pandas as pd
df_1 = pd.DataFrame({
    "id_1": [1, 1, 2, 3],
    "id_2": [10, 20, 30, 40],
    "value": [100, 200, 300, 400]
})

# PySpark 版本
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("match_check").getOrCreate()
df_1_spark = spark.createDataFrame([
    (1, 10, 100), (1, 20, 200), (2, 30, 300), (3, 40, 400)
], ["id_1", "id_2", "value"])

df_2(待校验数据集)

# Pandas 版本
df_2 = pd.DataFrame({
    "id_1": [1, 1, 2, 3, 4],
    "id_2": [10, 20, 30, 50, 60],
    "value": [100, 201, 300, 400, 500]
})

# PySpark 版本
df_2_spark = spark.createDataFrame([
    (1, 10, 100), (1, 20, 201), (2, 30, 300), (3, 50, 400), (4, 60, 500)
], ["id_1", "id_2", "value"])

Pandas 实现方案

核心逻辑:通过左关联匹配id_1+id_2组合,直接对比value值筛选异常行,无需复杂函数。

# 左关联保留df_2所有行,匹配df_1对应数据
merged = df_2.merge(df_1, on=["id_1", "id_2"], how="left", suffixes=("_df2", "_df1"))

# 筛选异常:ID对不存在 或 value不匹配
abnormal_rows = merged[(merged["value_df1"].isna()) | (merged["value_df2"] != merged["value_df1"])]

# 添加错误类型说明并整理输出格式
abnormal_rows["error_type"] = abnormal_rows.apply(
    lambda x: "ID对在df_1中不存在" if pd.isna(x["value_df1"]) else "Value值不匹配", axis=1
)
abnormal_rows = abnormal_rows[["id_1", "id_2", "value_df2", "value_df1", "error_type"]].rename(
    columns={"value_df2": "df2_value", "value_df1": "df1_value"}
)

print(abnormal_rows)

PySpark 实现方案

核心逻辑:用左关联替代复杂UDF,通过内置函数直接判断异常条件,性能更优。

from pyspark.sql.functions import col, when

# 左关联匹配ID对
merged_spark = df_2_spark.alias("df2").join(
    df_1_spark.alias("df1"),
    on=["id_1", "id_2"],
    how="left"
)

# 筛选异常行并标注错误类型
abnormal_rows_spark = merged_spark.filter(
    col("df1.value").isNull() | (col("df2.value") != col("df1.value"))
).withColumn(
    "error_type",
    when(col("df1.value").isNull(), "ID对在df_1中不存在")
    .otherwise("Value值不匹配")
).select(
    col("id_1"), col("id_2"),
    col("df2.value").alias("df2_value"),
    col("df1.value").alias("df1_value"),
    col("error_type")
)

abnormal_rows_spark.show()

期望输出

两种方案最终输出的异常行一致:

id_1id_2df2_valuedf1_valueerror_type
120201200Value值不匹配
350400nullID对在df_1中不存在
460500nullID对在df_1中不存在

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 04:20:33