使用PySpark读取字节文件转ASCII时遇AttributeError问题排查
问题与解决方案
问题场景
有一个以#*为分隔符的字节信息文件,内容示例如下:
b'\x00\x00V\x97'#*b'2%'#*b'\x00\x00'#*b'\xc5'#*b'\t'#*b'\xc0'
使用PySpark读取该文件并按不同列规则转换为ASCII时,抛出错误:
AttributeError: 'str' object has no attribute 'decode'
原读取与转换代码如下:
原读取代码
from pyspark.sql.types import StructType, StructField, StringType from pyspark.sql.functions import udf, col customSchema = StructType( [ StructField("Col1", StringType(), True), StructField("Col2", StringType(), True), StructField("Col3", StringType(), True), ] ) df = ( spark.read.format("csv") .option("inferSchema", "true") .option("header", "true") .schema(customSchema) .option("sep", "#*") .load("/FileStore/tables/EbcdicTextData.txt") ) df.show() # 输出示例: # |b'\x00\x01d0'|b'I\x08'|b'\x00\x00'|
原转换函数与UDF
def unpack_ch_or_zd(bytes: bytearray) -> str: ascii_text = bytes.decode("cp037").replace("\x00", "").rstrip() return ascii_text if ascii_text.isascii() else "Non-ASCII" def unpack_pd_or_pd_plus(bytes) -> str: ascii_text = ( "" if bytes.hex()[-1:] != "d" and bytes.hex()[-1:] != "b" else "-" ) + bytes.hex()[:-1] return ascii_text if ascii_text.isascii() else "Non-ASCII" def unpack_pd_or_pd_plus_dec(bytes, decimal: int) -> str: ascii_text = ( "" if bytes.hex()[-1:] != "d" and bytes.hex()[-1:] != "b" else "-" ) + bytes.hex()[:-1] ascii_text = ( ascii_text[:-decimal] + "." + ascii_text[-decimal:] if ascii_text.isascii() else "Non-ASCII" ) return ascii_text def unpack_bi_or_biplus_no_dec(bytes) -> str: a = str(int("0x" + bytes.hex(), 0)) return a if a.isascii() else "Non-ASCII" unpack_ch_or_zd_UDF = udf(lambda x: unpack_ch_or_zd(x), StringType()) unpack_pd_or_pd_plus_UDF = udf(lambda x: unpack_pd_or_pd_plus(x), StringType()) unpack_pd_or_pd_plus_dec_UDF = udf(lambda x: unpack_pd_or_pd_plus_dec(x), StringType())
原列处理逻辑
for row in layout_df.collect(): column_name = row["Field_name"] conversion_type = row["Python_Data_type"] decimal = row.get("decimal", 0) # 补充原代码缺失的decimal获取逻辑 if conversion_type.lower() == "ch" or conversion_type.lower() == "zd": df = df.withColumn(column_name, unpack_ch_or_zd_UDF(col(column_name))) elif ( conversion_type.lower() == "pd" or conversion_type.lower() == "pd+" ) and decimal == 0: df = df.withColumn(column_name, unpack_pd_or_pd_plus_UDF(col(column_name))) elif ( conversion_type.lower() == "pd" or conversion_type.lower() == "pd+" ) and decimal > 0: df = df.withColumn(column_name, unpack_pd_or_pd_plus_dec_UDF(col(column_name)))
错误原因
读取后的列值是字符串形式的bytes字面量(如"b'\x00\x01d0'"),而非真实的bytes对象。原转换函数直接调用decode()方法,而str类型没有该方法,因此抛出错误。
解决办法
步骤1:新增字符串转bytes的辅助函数
将字符串形式的bytes字面量解析为真实的bytes对象,同时处理解析失败的情况:
import ast def str_to_bytes(byte_str: str) -> bytes: try: # 安全解析字符串为bytes对象,避免eval的安全风险 return ast.literal_eval(byte_str) except (SyntaxError, ValueError): # 解析失败时返回空bytes return b""
步骤2:修改转换函数,先解析字符串为bytes
调整所有转换函数,先调用str_to_bytes处理输入:
def unpack_ch_or_zd(byte_str: str) -> str: bytes_data = str_to_bytes(byte_str) ascii_text = bytes_data.decode("cp037").replace("\x00", "").rstrip() return ascii_text if ascii_text.isascii() else "Non-ASCII" def unpack_pd_or_pd_plus(byte_str: str) -> str: bytes_data = str_to_bytes(byte_str) if not bytes_data: return "Non-ASCII" hex_str = bytes_data.hex() sign = "-" if hex_str[-1:] in ("d", "b") else "" ascii_text = sign + hex_str[:-1] return ascii_text if ascii_text.isascii() else "Non-ASCII" def unpack_pd_or_pd_plus_dec(byte_str: str, decimal: int) -> str: bytes_data = str_to_bytes(byte_str) if not bytes_data: return "Non-ASCII" hex_str = bytes_data.hex() sign = "-" if hex_str[-1:] in ("d", "b") else "" num_str = sign + hex_str[:-1] if not num_str.isascii() or decimal < 0 or len(num_str) < decimal: return "Non-ASCII" # 处理小数位 if decimal == 0: return num_str return num_str[:-decimal] + "." + num_str[-decimal:] def unpack_bi_or_biplus_no_dec(byte_str: str) -> str: bytes_data = str_to_bytes(byte_str) if not bytes_data: return "Non-ASCII" hex_str = bytes_data.hex() try: num = int(f"0x{hex_str}", 0) num_str = str(num) return num_str if num_str.isascii() else "Non-ASCII" except ValueError: return "Non-ASCII"
步骤3:更新UDF定义
去掉冗余lambda,直接绑定修改后的函数;针对带decimal参数的场景,用闭包创建动态UDF:
unpack_ch_or_zd_UDF = udf(unpack_ch_or_zd, StringType()) unpack_pd_or_pd_plus_UDF = udf(unpack_pd_or_pd_plus, StringType()) # 为带小数位的PD类型创建动态UDF def create_unpack_pd_dec_udf(decimal: int): def func(byte_str: str): return unpack_pd_or_pd_plus_dec(byte_str, decimal) return udf(func, StringType())
步骤4:调整列处理逻辑
适配带参数的UDF:
for row in layout_df.collect(): column_name = row["Field_name"] conversion_type = row["Python_Data_type"] decimal = row.get("decimal", 0) if conversion_type.lower() in ("ch", "zd"): df = df.withColumn(column_name, unpack_ch_or_zd_UDF(col(column_name))) elif conversion_type.lower() in ("pd", "pd+"): if decimal == 0: df = df.withColumn(column_name, unpack_pd_or_pd_plus_UDF(col(column_name))) else: # 创建对应小数位的UDF并应用 pd_dec_udf = create_unpack_pd_dec_udf(decimal) df = df.withColumn(column_name, pd_dec_udf(col(column_name)))
关键说明
- 使用
ast.literal_eval安全解析字符串形式的bytes字面量,避免eval的安全风险 - 增加了解析失败和非法输入的容错处理,返回统一的"Non-ASCII"占位符
- 针对带小数位的PD类型,通过闭包创建动态UDF,灵活传递decimal参数
内容的提问来源于stack exchange,提问作者phani437
相关产品推荐
相关产品推荐

