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

PySpark applyInPandas分组后amount列出现null值问题排查

PySpark applyInPandas 运行后amount字段出现null的错误定位

你的代码存在4个直接导致异常的问题,按影响程度排序:

1. 核心触发原因:applyInPandas 按列位置匹配schema,而非列名

Spark 的applyInPandas不会根据列名对齐返回的Pandas DataFrame和定义的Schema,而是严格按照列的先后顺序做映射。你定义的Schema列顺序为:
id → time_stamp → account_id → IS_REPEATED → nature → amount → user_time → user_date
但你在check函数中返回的是未做列顺序调整的原始Pandas df,它的列顺序和Schema顺序完全不一致,且新增的IS_REPEATED列默认插在df末尾,直接导致所有列映射错位,映射到amount位置的列如果存在类型不匹配、空值的情况,就会出现你看到的amount为null的现象。

2. Pandas 赋值时索引自动对齐导致数据错位

你在check函数中给df_backup赋值nature_asc、nature_desc时,排序后调用了reset_index(drop=True),返回的Series索引是从0开始的连续整数,但df_backup的索引是传入时的原始行索引(非0开始连续值),Pandas 会按索引做自动对齐,对齐失败的位置直接填充NaN,后续计算IS_REPEATED和回写原df时会进一步放大行错位问题。

3. 分组键与需求不匹配

你注释说明需要按amount、nature分组,但实际groupby传入的分组键是["account_id", "amount"],分组逻辑本身不符合预期,会导致计算范围错误。

4. 基础语法错误

get_schema函数存在两个低级错误:

  • 变量名拼写错误:返回的shcema应为schema
  • 缩进不规范:虽然不直接触发本次null问题,但容易引发变量未定义的报错

修正后可运行代码

import pandas as pd
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

def xor_two_list(first,second):
    return 1 if first!=second else 0

def check(df: pd.DataFrame) -> pd.DataFrame:
    # 重置索引,避免原有索引干扰对齐
    df = df.reset_index(drop=True)
    # 按verified_time排序生成备份
    df_backup = df.sort_values('verified_time').reset_index(drop=True)
    # 直接计算升序、降序的nature序列,不需要提前赋值占位
    nature_asc = df_backup['nature'].sort_values(ascending=True).reset_index(drop=True)
    nature_desc = df_backup['nature'].sort_values(ascending=False).reset_index(drop=True)
    # 计算IS_REPEATED字段
    df_backup['IS_REPEATED'] = [xor_two_list(a,b) for a,b in zip(nature_asc, nature_desc)]
    # 严格按照schema定义的列顺序返回,彻底避免位置错位
    return df_backup[["id", "time_stamp", "account_id", "IS_REPEATED", "nature", "amount", "user_time", "user_date"]]

def get_schema(value):
    if value == 'trx':
        schema = StructType([
            StructField("id", StringType(), nullable=True),
            StructField("time_stamp", IntegerType(), nullable=True),
            StructField("account_id", StringType(), nullable=True),
            StructField("IS_REPEATED", IntegerType(), nullable=True),
            StructField("nature", IntegerType(), nullable=True),
            StructField("amount", IntegerType(), nullable=True),
            StructField("user_time", StringType(), nullable=True),
            StructField("user_date", StringType(), nullable=True)
        ])
    # 修正拼写错误
    return schema

# 分组键可根据实际需求替换,示例保留原代码的account_id+amount,若需按amount、nature分组直接替换参数即可
df_new = df_spark.groupby(["account_id", "amount"]).applyInPandas(check, schema=get_schema('trx'))

# 校验逻辑
df_new.select('account_id','amount','IS_REPEATED').where(col("amount").isNull()).show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.03 05:01:00