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
相关产品推荐
相关产品推荐

