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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 17:50:37