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.
解决步骤
处理字典查询的None值
动态生成的readingType可能不在_readingType字典的键集合中,导致get()返回None。修改查询逻辑,添加默认值避免空指针:# 替换原查询代码,指定默认元组避免None readingTypeMapping = _readingType.get(readingType, ("unknown_code", "unknown_desc"))修正Schema与返回值类型不匹配
原代码中readingTypeMapping是元组类型,但Schema定义为IntegerType/StringType,类型不匹配导致异常或错误输出。根据需求调整:- 若需要返回元组中的单个值(如
rc或rd),修改Schema对应字段类型为StringType,并在append时传入对应值:
同时Schema中保持# 比如返回rc作为readingTypeMapping "readingTypeMapping": rc,StructField("readingTypeMapping", StringType(), False)。 - 若需要返回完整元组,需将Schema中
readingTypeMapping定义为StructType:
并确保append时传入对应结构:StructField("readingTypeMapping", StructType([ StructField("rc", StringType(), False), StructField("rd", StringType(), False) ]), False),"readingTypeMapping": {"rc": rc, "rd": rd},
- 若需要返回元组中的单个值(如
校验动态生成的readingType值
添加临时打印语句,排查动态生成的readingType是否存在字典中:readingType = int(strIn[j + 2:j + 6], 16) # 临时打印调试 if readingType not in _readingType: print(f"Missing readingType: {readingType}")
内容的提问来源于stack exchange,提问作者Sai Pavan
相关产品推荐
相关产品推荐

