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

Spark Worker抛出AssertionError:PySpark数据处理代码排查

PySpark AssertionError排查+业务需求落地方案

先明确业务要求

  • 给input_table每一行,在output_table生成X条数据——X是该行StartDate对应年份的总天数
  • Info字段必须凑够Y个字符(Y=X*2),长度不足时用#补位
  • output_table的AM、PM列,按顺序从Info中逐个取字符填充

为啥Spark Worker会抛AssertionError?

Python代码正常运行,转PySpark就炸,大概率是踩了这几个坑:

1. 年份天数计算逻辑不一致

Python用datetime算闰年天数的逻辑,和PySpark日期函数可能存在偏差(比如时区解析、日期格式兼容性问题),导致X值错误,后续Info长度、字符分配全乱套,触发Spark内部断言校验。

  • 别自己写闰年判断,直接用PySpark原生函数计算全年天数:
    df_with_x = input_df.withColumn(
        "year_start", F.date_trunc("year", F.to_date("StartDate"))
    ).withColumn(
        "year_end", F.last_day(F.to_date(F.concat(F.year(F.to_date("StartDate")), "-12-31")))
    ).withColumn(
        "X", F.datediff(F.col("year_end"), F.col("year_start")) + 1
    )
    

2. Info字段长度未达标Y=X*2

如果生成的Info长度不等于X*2,后续取字符填充AM/PM时会触发索引越界,直接炸断言。

  • 用PySpark的lpad函数统一补位,注意length函数是按字符数统计(和Python字节统计逻辑不同):
    df_with_info = df_with_x.withColumn(
        "Info", F.lpad(F.coalesce(F.col("Info"), F.lit("")), F.col("X")*2, "#")
    )
    
    补完可以抽样验证:df_with_info.select("X", F.length("Info")).show(),确保每行length(Info)严格等于X*2。

3. 行展开时索引计算错误

生成X行数据时,若索引范围或字符位置计算错误,会触发Spark内部的断言校验。

  • 用sequence生成1到X的序列,explode展开后,按day_idx*2-1和day_idx*2取AM/PM(注意PySpark的substring是从1开始计数的):
    df_expanded = df_with_info.withColumn(
        "day_idx", F.explode(F.sequence(F.lit(1), F.col("X")))
    ).withColumn(
        "AM", F.substring(F.col("Info"), F.col("day_idx")*2 - 1, 1)
    ).withColumn(
        "PM", F.substring(F.col("Info"), F.col("day_idx")*2, 1)
    )
    

完整可运行代码

from pyspark.sql import functions as F

# 读取原表
input_df = spark.table("input_table")

# 计算年份天数X
df_with_x = input_df.withColumn(
    "year_start", F.date_trunc("year", F.to_date("StartDate"))
).withColumn(
    "year_end", F.last_day(F.to_date(F.concat(F.year(F.to_date("StartDate")), "-12-31")))
).withColumn(
    "X", F.datediff(F.col("year_end"), F.col("year_start")) + 1
)

# 生成符合长度要求的Info字段
df_with_info = df_with_x.withColumn(
    "Info", F.lpad(F.coalesce(F.col("Info"), F.lit("")), F.col("X")*2, "#")
)

# 展开为X行并填充AM、PM列
df_expanded = df_with_info.withColumn(
    "day_idx", F.explode(F.sequence(F.lit(1), F.col("X")))
).withColumn(
    "AM", F.substring(F.col("Info"), F.col("day_idx")*2 - 1, 1)
).withColumn(
    "PM", F.substring(F.col("Info"), F.col("day_idx")*2, 1)
).drop("year_start", "year_end", "X", "day_idx")  # 清理临时字段

# 追加到目标表
df_expanded.write.mode("append").saveAsTable("output_table")

额外校验步骤

跑代码前可以做几个验证:

  1. 对比Python和PySpark对同一年份的天数计算结果,确保X值正确
  2. 检查Info字段长度,确认每行length(Info) == X*2
  3. 展开后检查AM/PM列是否为空:df_expanded.filter(F.col("AM").isNull() | F.col("PM").isNull()).count(),结果应为0

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 19:02:16