Polars流处理中实现数据类型推断
超大Parquet文件的类型推断与转换解决方案
针对你遇到的全字符串列Parquet文件类型转换问题,以下是几种内存友好的解决方案,避免直接cast导致的全null问题:
方案1:使用PySpark(适合超大规模文件,分布式处理)
PySpark的try_cast函数会在转换失败时返回null,结合when条件判断,可以实现"能转则转,转不了保留原字符串"的逻辑:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when, try_cast # 初始化Spark会话 spark = SparkSession.builder.appName("ParquetTypeInference").getOrCreate() # 读取原Parquet文件 df = spark.read.parquet("input_file.parquet") # 定义单列类型推断逻辑:优先转int,失败转float,均失败则保留原字符串 def infer_col_type(col_name): return ( when(try_cast(col(col_name), "int").isNotNull(), try_cast(col(col_name), "int")) .when(try_cast(col(col_name), "float").isNotNull(), try_cast(col(col_name), "float")) .otherwise(col(col_name)) ) # 应用逻辑到所有列 transformed_df = df.select([infer_col_type(c).alias(c) for c in df.columns]) # 导出为新Parquet文件 transformed_df.write.parquet("output_file.parquet", mode="overwrite")
方案2:使用Dask(单机处理大文件,分块计算)
Dask支持分块读取大文件,适合内存不足以加载全量数据的场景:
import dask.dataframe as dd import pandas as pd # 分块读取Parquet文件 ddf = dd.read_parquet("input_file.parquet") # 自定义列类型推断函数 def infer_dtype(series): # 尝试转换为整数(保留非空转换结果) int_converted = pd.to_numeric(series, errors='coerce', downcast='integer') if not int_converted.isna().all(): return int_converted # 尝试转换为浮点数 float_converted = pd.to_numeric(series, errors='coerce', downcast='float') if not float_converted.isna().all(): return float_converted # 均失败则保留原字符串类型 return series.astype('string') # 应用推断逻辑到所有列 transformed_ddf = ddf.apply(infer_dtype, axis=0, meta={col: 'object' for col in ddf.columns}) # 导出结果 transformed_ddf.to_parquet("output_file.parquet", overwrite=True)
方案3:使用Pandas分块处理(单机小内存场景)
如果只能用Pandas,通过分块读取避免内存溢出:
import pandas as pd import os # 分块读取Parquet文件(chunksize根据内存调整) chunk_iter = pd.read_parquet("input_file.parquet", chunksize=100_000) temp_files = [] for idx, chunk in enumerate(chunk_iter): # 逐列处理类型转换 for col in chunk.columns: # 尝试转整数 int_conv = pd.to_numeric(chunk[col], errors='coerce', downcast='integer') if not int_conv.isna().all(): chunk[col] = int_conv continue # 尝试转浮点数 float_conv = pd.to_numeric(chunk[col], errors='coerce', downcast='float') if not float_conv.isna().all(): chunk[col] = float_conv continue # 保留原字符串类型 chunk[col] = chunk[col].astype('string') # 保存临时块 temp_path = f"temp_chunk_{idx}.parquet" chunk.to_parquet(temp_path) temp_files.append(temp_path) # 合并临时块并导出最终文件 final_df = pd.concat([pd.read_parquet(f) for f in temp_files]) final_df.to_parquet("output_file.parquet", overwrite=True) # 清理临时文件 for f in temp_files: os.remove(f)
为什么之前的cast(strict=False)会全变null?
直接使用cast(col, "int", strict=False)时,只要列中有任何无法转换为int的字符串,该位置会被设为null;如果整列都存在无法转换的值,最终整列都会变成null。而上述方案通过"尝试转换+保留原值"的逻辑,只替换能成功转换的行,保留无法转换的原字符串,避免全列null的问题。
内容的提问来源于stack exchange,提问作者Heiaha
相关产品推荐
相关产品推荐

