使用Polars处理Azure Blob存储中60GB Parquet数据时遭遇PyArrow字节读取错误的排查求助
Polars处理Azure Blob存储中60GB Parquet数据时遭遇PyArrow字节读取错误的排查求助
大家好,我遇到了一个棘手的问题,想请各位帮忙分析下:
我正在用Polars结合PyArrow读取Azure Blob存储里的一批Parquet文件,总大小大概60GB,运行环境是Standard_E8_v3计算实例(8核、64GB内存、200GB磁盘)。数据读取完成后,我尝试对数据做分组聚合操作,但在执行collect()的时候触发了错误,我完全看不懂错误信息想表达的意思,希望能搞清楚这几个问题:
- 是不是数据量太大,超出了当前机器的处理能力?
- 我的代码是不是哪里写得有问题?
- 是不是数据本身存在异常,需要做预处理?
特别说明:解决方案必须是基于Polars的,麻烦各位帮忙定位下问题,非常感谢!
我的代码如下:
import pyarrow.dataset as ds from azureml.fsspec import AzureMachineLearningFileSystem import polars as pl from azureml.core import Workspace ws = Workspace.from_config() # Azure Machine Learning workspace details: subscription = ws.subscription_id resource_group = ws.resource_group workspace = ws.name datastore_name = 'datastore_name' path_on_datastore = 'path_to_data' # long-form Datastore uri format: uri = f'azureml://subscriptions/{subscription}/resourcegroups/{resource_group}/workspaces/{workspace}/datastores/{datastore_name}' aml_fs = AzureMachineLearningFileSystem(uri) files = aml_fs.glob() myds=ds.dataset(path_on_datastore, filesystem=aml_fs, format="parquet") df = ( pl.scan_pyarrow_dataset(myds) .select([ 'COLUMN_LIST' ]) #.with_columns(pl.col('turnovervalue').cast(pl.Float64, strict=False)) .filter((pl.col('col1')>0)&(pl.col('col2') >= 2022)) ) grouped = (df.lazy() .groupby(['colA','colB']) .agg( [ pl.n_unique('colC').alias('Blah'), pl.sum('colD').alias("BlahBlah"), pl.n_unique('colE').alias('BlahBlahBlah'), (pl.col('colF') == "C").count().alias('BlahBlahBlahBlah') ] ) ).collect()
编辑补充:
我检查了Polars DataFrame的schema,结果看起来是合理的,所以我觉得Azure的连接应该没问题。于是我往上排查,发现问题可能出在PyArrow这边——当我尝试查看从Azure获取的PyArrow数据集的head时,出现了字节读取相关的错误。看起来是PyArrow推断的数据类型和实际读取到的数据不匹配,但我不知道整个表中第4位置的数据是什么,也完全不知道该怎么去排查这个点。现在我打算调整问题的标签和标题,希望能找到熟悉这类问题的人来帮忙。
备注:内容来源于stack exchange,提问作者Hillygoose
相关产品推荐
相关产品推荐

