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

PySpark UDF读取字典异常:动态取值失败与类型匹配问题

PySpark UDF字典查询与Schema类型匹配问题解决

问题现象

  • 定义了包含300个键的_readingType字典,静态传入测试值(如readingType=83)时,_readingType.get(readingType)能正常返回预期结果;但动态解析字符串生成readingType时,查询字典返回None,触发TypeError: 'NoneType' object is not subscriptable错误。
  • Schema定义中,若readingTypeMapping设为IntegerType,输出结果中该字段固定为0;设为StringType则触发空指针异常。

相关代码

解析函数定义

from datetime import datetime
from pyspark.sql.types import *
from pyspark.sql.functions import udf, col, explode

def parse_data(strIn):
    result = []
    j = 0
    _readingQuality = {0: ("", "")}
    
    _readingType = {0: ("", "")}  # 实际包含300个键的字典

    while j < len(strIn) - 1:
        ts = int(strIn[j:j + 8], 16)
        timeStamp = str(datetime.utcfromtimestamp(ts))
        numEvents = int(strIn[j + 8], 16)
        numReadings = int(strIn[j + 9], 16)
        j += 10
        for i in range(numReadings):
            binary = bin(int(strIn[j:j + 2], 16))[2:].zfill(8)
            timeStampPresent = bool(int(binary[0], 2))
            readingQualitiesPresent = bool(int(binary[1], 2))
            pendingPowerOfTen = int(binary[2:5], 2)
            if pendingPowerOfTen == 7:
                pendingPowerOfTen = 9
            readingsValueSizeInBytes = int(binary[5:], 2) + 1
            readingType = int(strIn[j + 2:j + 6], 16)
            readingTypeMapping = _readingType.get(readingType)
            rc = readingTypeMapping[0]  # 此处readingTypeMapping为None时触发报错
            rd = readingTypeMapping[1]
            if timeStampPresent:
                ts = int(strIn[j + 6:j + 14], 16)
                timeStamp = str(datetime.utcfromtimestamp(ts))
                j = j + 14
            else:
                timeStamp = ''
                j = j + 6

            # 处理读取质量
            if readingQualitiesPresent:
                q = []
                c = []
                d = []
                finished = False
                while not finished:
                    binary = bin(int(strIn[j:j + 2], 16))[2:].zfill(8)
                    num = int(binary[1:], 2)
                    q.append(num)
                    t = _readingQuality.get(num, (None, None))
                    c.append(t[0])
                    d.append(t[1])
                    if binary[0] == '0':
                        finished = True
                    else:
                        j = j + 2
                j = j + 2

            # 读取值
            readingValue = strIn[j:j + readingsValueSizeInBytes * 2]
            j = j + readingsValueSizeInBytes * 2
            
            result.append({
                "i": str(i),
                "timeStamp": timeStamp,
                "numEvents": numEvents,
                "numReadings": numReadings,
                "readingType": readingType,
                "readingTypeMapping": readingTypeMapping,
                "rc": rc,
                "rd": rd,
                "readingValue": readingValue
            })
    return result

UDF与Schema定义

parse_data_udf = udf(parse_data, ArrayType(StructType([
    StructField("i", StringType(), False),
    StructField("timeStamp", StringType(), False),
    StructField("numEvents", IntegerType(), False),
    StructField("numReadings", IntegerType(), False),
    StructField("readingType", IntegerType(), False),
    StructField("readingTypeMapping", IntegerType(), False),
    StructField("rc", StringType(), False),
    StructField("rd", StringType(), False),  # 修正原代码重复定义rc的问题
    StructField("readingValue", StringType(), False)
])))

# 应用UDF并展开结果
df_parsed = df.withColumn("parsed", parse_data_udf(col("PL")))
df_exploded = df_parsed.select("PL", explode(col("parsed")).alias("parsed"))
display(df_exploded)

错误信息

PythonException: 'TypeError: 'NoneType' object is not subscriptable', from , line 386.

解决步骤

  1. 处理字典查询的None值
    动态生成的readingType可能不在_readingType字典的键集合中,导致get()返回None。修改查询逻辑,添加默认值避免空指针:

    # 替换原查询代码,指定默认元组避免None
    readingTypeMapping = _readingType.get(readingType, ("unknown_code", "unknown_desc"))
    
  2. 修正Schema与返回值类型不匹配
    原代码中readingTypeMapping是元组类型,但Schema定义为IntegerType/StringType,类型不匹配导致异常或错误输出。根据需求调整:

    • 若需要返回元组中的单个值(如rc或rd),修改Schema对应字段类型为StringType,并在append时传入对应值:
      # 比如返回rc作为readingTypeMapping
      "readingTypeMapping": rc,
      
      同时Schema中保持StructField("readingTypeMapping", StringType(), False)。
    • 若需要返回完整元组,需将Schema中readingTypeMapping定义为StructType:
      StructField("readingTypeMapping", StructType([
          StructField("rc", StringType(), False),
          StructField("rd", StringType(), False)
      ]), False),
      
      并确保append时传入对应结构:
      "readingTypeMapping": {"rc": rc, "rd": rd},
      
  3. 校验动态生成的readingType值
    添加临时打印语句,排查动态生成的readingType是否存在字典中:

    readingType = int(strIn[j + 2:j + 6], 16)
    # 临时打印调试
    if readingType not in _readingType:
        print(f"Missing readingType: {readingType}")
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 00:45:52