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

如何修复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


解决建议

问题根源

  1. 返回类型不一致:函数在l==1时返回字典,其他场景返回dict_items视图,Spark无法处理这种混合输出类型
  2. 序列化不兼容:dict_items是Python的视图对象,无法被Spark正确序列化为指定的ArrayType(StructType)结构
  3. 类型不匹配:部分时长值是float类型,与Schema中定义的IntegerType冲突

修复步骤

  1. 统一返回类型:无论哪种场景,都返回包含二元组的列表,每个二元组对应Struct的两个字段
  2. 转换时长为整数:确保所有时长值都是int类型,匹配Schema定义
  3. 简化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 19:05:07