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")
额外校验步骤
跑代码前可以做几个验证:
- 对比Python和PySpark对同一年份的天数计算结果,确保X值正确
- 检查
Info字段长度,确认每行length(Info) == X*2 - 展开后检查AM/PM列是否为空:
df_expanded.filter(F.col("AM").isNull() | F.col("PM").isNull()).count(),结果应为0
内容的提问来源于stack exchange,提问作者Eli
相关产品推荐
相关产品推荐

