Polars中列转Decimal类型报错求助(附Spark实现示例)
Polars实现平面文件读取与类型转换解决方案
问题说明
读取平面文件并为DataFrame分配列名,将指定列转换为Decimal类型时出现报错,该功能已在Spark中实现,需适配Polars完成相同操作。
示例数据
data = b""" D|120277|Patricia|167.2|26/12/1982 D|757773|Charles|167.66|08/04/2019 D|248498|Katrina|168.68|20/11/2016 D|325561|Christina|170.86|05/10/1998 D|697464|Joshua|171.41|07/09/1970 D|244750|Zachary|169.43|08/12/2014 """.strip()
原Polars代码及报错
原代码
import polars as pl cols_dict = {'column_1': 'rtype', 'column_2': 'EMP_ID', 'column_3': 'EMP_NAME', 'column_4': 'SALARY', 'column_5': 'DOB'} df = pl.read_csv(data, separator='|', has_header=False) df = df.select(pl.all().name.map(lambda col_name: cols_dict.get(col_name))) df = df.with_columns( pl.col('EMP_ID').cast(pl.Decimal(scale=6, precision=0)))
报错信息
InvalidOperationError: conversion from
i64todecimal[1,6]failed in column 'EMP_ID' for 5 out of 5 values: [120277, 757773, … 697464]
Spark实现参考
Spark代码
rdd = sparkContext.textFile(core_data_file).filter(lambda x: x[0] == "D").map( lambda x: x.split('|')) sparkDf= rdd.toDF(schema=["rtype"] + list(cols_dict.values())) sparkDf= sparkDf.withColumn(col_name, coalesce( col(col_name).cast(DecimalType(data_length, data_scale)), lit(0))) sparkDf.show()
Spark输出结果
+-----+------+---------+------+-------------------+ |rtype|EMP_ID|EMP_NAME |SALARY|DOB | +-----+------+---------+------+-------------------+ |D |120277|Patricia |167.20|1982-12-26 00:00:00| |D |757773|Charles |167.66|2019-04-08 00:00:00| |D |248498|Katrina |168.68|2016-11-20 00:00:00| |D |325561|Christina|170.86|1998-10-05 00:00:00| |D |697464|Joshua |171.41|1970-09-07 00:00:00| |D |244750|Zachary |169.43|2014-12-08 00:00:00| +-----+------+---------+------+-------------------+
Polars正确实现
错误原因
Polars中pl.Decimal的参数逻辑需注意:
pl.Decimal(precision, scale)里,precision是总位数(整数位+小数位),scale是小数位数。- 原代码错误将
precision设为0、scale设为6,导致精度不足以容纳6位整数,触发转换失败。
修正后的Polars代码
import polars as pl cols_dict = {'column_1': 'rtype', 'column_2': 'EMP_ID', 'column_3': 'EMP_NAME', 'column_4': 'SALARY', 'column_5': 'DOB'} # 读取文件时直接指定列名与初始类型,提升效率 df = pl.read_csv( data, separator='|', has_header=False, new_columns=list(cols_dict.values()), dtypes={ 'EMP_ID': pl.Int64, 'SALARY': pl.Float64, 'DOB': pl.String } ) # 按需求转换列类型 df = df.with_columns([ # EMP_ID为6位整数,设精度6、小数位0 pl.col('EMP_ID').cast(pl.Decimal(precision=6, scale=0)), # SALARY保留2位小数,总位数设为5(3位整数+2位小数) pl.col('SALARY').cast(pl.Decimal(precision=5, scale=2)), # DOB按指定格式转日期时间类型 pl.col('DOB').str.to_date(format="%d/%m/%Y").cast(pl.Datetime) ]) # 查看结果 print(df)
输出结果
shape: (6, 5) ┌──────┬────────────┬───────────┬────────────┬─────────────────────┐ │ rtype┆ EMP_ID ┆ EMP_NAME ┆ SALARY ┆ DOB │ │ --- ┆ --- ┆ --- ┆ --- ┆ --- │ │ str ┆ decimal[6,0]┆ str ┆ decimal[5,2]┆ datetime[μs] │ ╞══════╪════════════╪═══════════╪════════════╪═════════════════════╡ │ D ┆ 120277 ┆ Patricia ┆ 167.20 ┆ 1982-12-26 00:00:00 │ │ D ┆ 757773 ┆ Charles ┆ 167.66 ┆ 2019-04-08 00:00:00 │ │ D ┆ 248498 ┆ Katrina ┆ 168.68 ┆ 2016-11-20 00:00:00 │ │ D ┆ 325561 ┆ Christina ┆ 170.86 ┆ 1998-10-05 00:00:00 │ │ D ┆ 697464 ┆ Joshua ┆ 171.41 ┆ 1970-09-07 00:00:00 │ │ D ┆ 244750 ┆ Zachary ┆ 169.43 ┆ 2014-12-08 00:00:00 │ └──────┴────────────┴───────────┴────────────┴─────────────────────┘
内容的提问来源于stack exchange,提问作者Smruti Prakash Mohanty
相关产品推荐
相关产品推荐

