【发布时间】:2018-03-26 02:36:59
【问题描述】:
我有一个 hive 格式和快速压缩的 parquet 文件。它适合内存,pandas.info 提供以下数据。
parquet 文件中每组的行数只有 100K
>>> df.info()
<class 'pandas.core.frame.DataFrame'>
Index: 21547746 entries, YyO+tlZtAXYXoZhNr3Vg3+dfVQvrBVGO8j1mfqe4ZHc= to oE4y2wK5E7OR8zyrCHeW02uTeI6wTwT4QTApEVBNEdM=
Data columns (total 8 columns):
payment_method_id int16
payment_plan_days int16
plan_list_price int16
actual_amount_paid int16
is_auto_renew bool
transaction_date datetime64[ns]
membership_expire_date datetime64[ns]
is_cancel bool
dtypes: bool(2), datetime64[ns](2), int16(4)
memory usage: 698.7+ MB
现在,用 dask 进行一些简单的计算,我得到以下时间
使用线程
>>>time.asctime();ddf.actual_amount_paid.mean().compute();time.asctime()
'Fri Oct 13 23:44:50 2017'
141.98732048354384
'Fri Oct 13 23:44:59 2017'
使用分布式(本地集群)
>>> c=Client()
>>> time.asctime();ddf.actual_amount_paid.mean().compute();time.asctime()
'Fri Oct 13 23:47:04 2017'
141.98732048354384
'Fri Oct 13 23:47:15 2017'
>>>
没关系,每次大约 9 秒。
现在使用多处理,惊喜来了……
>>> time.asctime();ddf.actual_amount_paid.mean().compute(get=dask.multiprocessing.get);time.asctime()
'Fri Oct 13 23:50:43 2017'
141.98732048354384
'Fri Oct 13 23:57:49 2017'
>>>
我希望多处理和分布式/本地集群处于同一数量级,但线程可能存在一些差异(无论好坏)
但是,多处理需要 47 倍以上的时间来对 in16 列求一个简单的平均值?
我的环境只是一个带有所需模块的全新 conda 安装。没有精心挑选任何东西。
为什么会有这种差异??我无法管理 dask/distributed 以使其具有可预测的行为,以便能够根据问题的性质在不同的调度程序之间做出明智的选择。
这只是一个玩具示例,但我无法找到符合我期望的示例(至少我对阅读文档的理解)。
有什么我应该记住的吗?还是我完全没有抓住重点?
谢谢
JC
【问题讨论】:
标签: python pandas distributed dask fastparquet