如何修复PySpark中ValueError: Unexpected tuple with StructType错误
问题描述
编写了用于计算两个时间戳之间每小时时长分布的PySpark UDF,独立运行函数输出符合预期,但应用到DataFrame时抛出ValueError: Unexpected tuple '2021-11-01:20' with StructType错误。
函数代码
def compute_hourly_viewing(start:datetime, end:datetime): hours_list = getHoursList(start,end) hourViewingMap={} l = len(hours_list) if l==1: hourViewingMap[createKeyHoursViewing(hours_list[0])]=(end-start).total_seconds() return hourViewingMap if l>2: i =1 while i < l-1: hourViewingMap[createKeyHoursViewing(hours_list[i])]=3600 i+=1 hourViewingMap[createKeyHoursViewing(hours_list[0])] = 3600 - getSecondsinHour(start) hourViewingMap[createKeyHoursViewing(hours_list[l-1])] = getSecondsinHour(end) return (hourViewingMap.items())
独立测试结果
d1 = datetime.datetime(2010, 1, 22, 13, 29, 40) d2= datetime.datetime(2010, 1, 22, 16, 31, 23) compute_hourly_viewing(d1,d2) Output: dict_items([('2010-01-22:14', 3600), ('2010-01-22:15', 3600), ('2010-01-22:13', 1820), ('2010-01-22:16', 1883)])
UDF定义与报错代码
schema = ArrayType(StructType( ([ StructField('date_hour' , StringType() , False), StructField('duration' , IntegerType() , False) ]) )) computeHourlyViewingUdf = udf(lambda x,y:compute_hourly_viewing(x,y),schema) df_step3 = df_step2.withColumn("hour_secondCnt", explode_outer(computeHourlyViewingUdf(col("customer_program_start_time"),col("customer_program_end_time"))))
错误信息:
An exception was thrown from the Python worker. Please see the stack trace below.
'ValueError: Unexpected tuple '2021-11-01:20' with StructType'. Full traceback below:
ValueError: Unexpected tuple '2021-11-01:20' with StructType
解决建议
问题根源
- 返回类型不一致:函数在
l==1时返回字典,其他场景返回dict_items视图,Spark无法处理这种混合输出类型 - 序列化不兼容:
dict_items是Python的视图对象,无法被Spark正确序列化为指定的ArrayType(StructType)结构 - 类型不匹配:部分时长值是
float类型,与Schema中定义的IntegerType冲突
修复步骤
- 统一返回类型:无论哪种场景,都返回包含二元组的列表,每个二元组对应Struct的两个字段
- 转换时长为整数:确保所有时长值都是
int类型,匹配Schema定义 - 简化UDF定义:直接使用函数作为UDF参数,无需额外Lambda包装
修改后的函数与UDF代码
def compute_hourly_viewing(start: datetime, end: datetime): hours_list = getHoursList(start, end) hourViewingMap = {} l = len(hours_list) if l == 1: key = createKeyHoursViewing(hours_list[0]) # 转换为整数,匹配Schema的IntegerType hourViewingMap[key] = int((end - start).total_seconds()) else: if l > 2: i = 1 while i < l - 1: key = createKeyHoursViewing(hours_list[i]) hourViewingMap[key] = 3600 i += 1 # 处理起始小时 start_key = createKeyHoursViewing(hours_list[0]) hourViewingMap[start_key] = int(3600 - getSecondsinHour(start)) # 处理结束小时 end_key = createKeyHoursViewing(hours_list[-1]) hourViewingMap[end_key] = int(getSecondsinHour(end)) # 统一返回列表形式的二元组,确保Spark能正确序列化 return list(hourViewingMap.items()) # 定义UDF时直接传入函数,无需Lambda computeHourlyViewingUdf = udf(compute_hourly_viewing, schema) # 应用UDF df_step3 = df_step2.withColumn( "hour_secondCnt", explode_outer(computeHourlyViewingUdf(col("customer_program_start_time"), col("customer_program_end_time"))) )
额外验证点
- 确认
getHoursList函数返回的是正确的小时时间戳列表 - 确认
createKeyHoursViewing函数返回的是符合预期的字符串格式(如YYYY-MM-DD:HH) - 确认
getSecondsinHour函数返回的是整数类型的秒数
内容的提问来源于stack exchange,提问作者Sri
相关产品推荐
相关产品推荐

