使用dask.read_parquet读取Parquet文件时字段类型不兼容报错
解决Dask读取Parquet时的类型冲突问题
方法1:指定统一Schema强制转换
Dask支持通过schema参数手动指定字段类型,强制将冲突的dt字段统一为字符串类型:
import dask.dataframe as dd # 直接定义schema,将dt指定为字符串类型 schema = {"dt": str} mme = dd.read_parquet(mme_path, storage_options=S3_CREDENTIALS, schema=schema)
如果不确定其他字段的类型,可以先用Pandas读取单个分片生成基准meta:
import pandas as pd # 读取单个Parquet分片获取样本结构 sample_df = pd.read_parquet(f"{mme_path}/part-00000.parquet", storage_options=S3_CREDENTIALS) # 修改dt字段类型为字符串 sample_df["dt"] = sample_df["dt"].astype(str) # 生成Dask可用的meta对象 meta = dd.utils.make_meta(sample_df) mme = dd.read_parquet(mme_path, storage_options=S3_CREDENTIALS, schema=meta)
方法2:禁用类型合并校验(谨慎使用)
通过arrow_options关闭Arrow的严格类型合并校验,让Dask兼容不同分片的类型差异(类似Pandas的处理逻辑):
mme = dd.read_parquet( mme_path, storage_options=S3_CREDENTIALS, arrow_options={"merge_schema": False} )
注意:该方法可能引发后续计算的类型异常,仅在确认dt字段内容可兼容转换时使用。
方法3:手动读取分片再合并
如果前两种方法无效,可以遍历所有Parquet分片,用Pandas读取后统一类型,再转为Dask DataFrame:
import s3fs from dask.dataframe import from_delayed, delayed import dask.dataframe as dd import pandas as pd # 列出S3路径下的所有Parquet文件 fs = s3fs.S3FileSystem(**S3_CREDENTIALS) file_paths = [f"s3://{path}" for path in fs.glob(f"{mme_path}/*.parquet")] # 定义单文件读取并转换类型的函数 def process_file(path): df = pd.read_parquet(path, storage_options=S3_CREDENTIALS) df["dt"] = df["dt"].astype(str) return df # 生成延迟任务并合并为Dask DataFrame delayed_tasks = [delayed(process_file)(path) for path in file_paths] mme = from_delayed(delayed_tasks)
内容的提问来源于stack exchange,提问作者dkjackyu
相关产品推荐
相关产品推荐

