- c - 在位数组中找到第一个零
- linux - Unix 显示有关匹配两种模式之一的文件的信息
- 正则表达式替换多个文件
- linux - 隐藏来自 xtrace 的命令
我正在尝试创建一个 dask.dataframe来自一堆大型 CSV 文件(目前有 12 个文件,每个文件有 8-10 百万行和 50 列)。它们中的一些可能会一起放入我的系统内存中,但它们肯定不会同时全部放入,因此使用 dask 而不是常规 pandas。
由于读取每个 csv 文件涉及一些额外的工作(添加包含文件路径中的数据的列),我尝试从延迟对象列表创建 dask.dataframe,类似于 to this example .
这是我的代码:
import dask.dataframe as dd
from dask.delayed import delayed
import os
import pandas as pd
def read_file_to_dataframe(file_path):
df = pd.read_csv(file_path)
df['some_extra_column'] = 'some_extra_value'
return df
if __name__ == '__main__':
path = '/path/to/my/files'
delayed_collection = list()
for rootdir, subdirs, files in os.walk(path):
for filename in files:
if filename.endswith('.csv'):
file_path = os.path.join(rootdir, filename)
delayed_reader = delayed(read_file_to_dataframe)(file_path)
delayed_collection.append(delayed_reader)
df = dd.from_delayed(delayed_collection)
print(df.compute())
启动此脚本(Python 3.4,dask 0.12.0)时,它会运行几分钟,而我的系统内存会不断填满。当它被完全使用时,一切都开始滞后,它会再运行几分钟,然后它会因 killed
或 MemoryError
而崩溃。
我认为 dask.dataframe 的全部意义在于能够对跨越磁盘上多个文件的大于内存的数据帧进行操作,那么我在这里做错了什么?
编辑: 据我所知,使用 df = dd.read_csv(path + '/*.csv')
读取文件似乎工作正常。但是,这不允许我使用文件路径中的附加数据来更改每个数据帧。
编辑#2:按照 MRocklin 的回答,我尝试使用 dask 的 read_bytes() method 读取我的数据。以及使用 single-threaded scheduler以及两者结合使用。尽管如此,即使在具有 8GB 内存的笔记本电脑上以单线程模式读取 100MB 的 block ,我的进程迟早会被杀死。不过,在一堆形状相似的小文件(每个大约 1MB)上运行下面所述的代码效果很好。知道我在这里做错了什么吗?
import dask
from dask.bytes import read_bytes
import dask.dataframe as dd
from dask.delayed import delayed
from io import BytesIO
import pandas as pd
def create_df_from_bytesio(bytesio):
df = pd.read_csv(bytesio)
return df
def create_bytesio_from_bytes(block):
bytesio = BytesIO(block)
return bytesio
path = '/path/to/my/files/*.csv'
sample, blocks = read_bytes(path, delimiter=b'\n', blocksize=1024*1024*100)
delayed_collection = list()
for datafile in blocks:
for block in datafile:
bytesio = delayed(create_bytesio_from_bytes)(block)
df = delayed(create_df_from_bytesio)(bytesio)
delayed_collection.append(df)
dask_df = dd.from_delayed(delayed_collection)
print(dask_df.compute(get=dask.async.get_sync))
最佳答案
如果您的每个文件都很大,那么在 Dask 有机会变聪明之前,对 read_file_to_dataframe
的几个并发调用可能会淹没内存。
Dask 尝试通过按顺序运行函数来在低内存中运行,以便快速删除中间结果。然而,如果只有几个函数的结果可以填满内存,那么 Dask 可能永远没有机会删除东西。例如,如果您的每个函数都生成了一个 2GB 的数据帧,并且您同时运行了 8 个线程,那么在 Dask 的调度策略生效之前,您的函数可能会生成 16GB 的数据。
read_csv 起作用的原因是它将大的 CSV 文件分成许多约 100MB 的字节 block (参见 blocksize=
关键字参数)。您也可以这样做,尽管这很棘手,因为您需要始终在端线处中断。
dask.bytes.read_bytes
函数可以帮到你。它可以将单个路径转换为 delayed
对象列表,每个对象对应于该文件的一个字节范围,该字节范围在分隔符上干净地开始和停止。然后,您可以将这些字节放入 io.BytesIO
(标准库)并调用 pandas.read_csv
。请注意,您还必须处理 header 等。该函数的文档字符串非常广泛,应该会提供更多帮助。
在上面的示例中,如果我们没有来自并行性的 8 倍乘数,一切都会很好。我怀疑如果你一次只运行一个函数,那么事情可能会在没有达到你的内存限制的情况下流水线。您可以使用以下行将 dask 设置为仅使用单个线程
dask.set_options(get=dask.async.get_sync)
注意:For Dask 版本 >= 0.15,您需要改用 dask.local.get_sync
。
如果你制作一个 dask.dataframe 然后立即计算它
ddf = dd.read_csv(...)
df = ddf.compute()
您正在将所有数据加载到 Pandas 数据框中,这最终会耗尽内存。相反,最好在 Dask 数据帧上操作并且只计算小的结果。
# result = df.compute() # large result fills memory
result = df.groupby(...).column.mean().compute() # small result
CSV 是一种普遍实用的格式,但也有一些缺陷。您可能会考虑使用 HDF5 或 Parquet 等数据格式。
关于python - 从延迟收集创建大型 dask.dataframe 时被杀死/内存错误,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/41266078/
假设我有 3 个 DataFrame。其中一个 DataFrame 的列名不在其他两个中。 using DataFrames df1 = DataFrame([['a', 'b', 'c'], [1,
假设我有 3 个 DataFrame。其中一个 DataFrame 的列名不在其他两个中。 using DataFrames df1 = DataFrame([['a', 'b', 'c'], [1,
我有一个 largeDataFrame(多列和数十亿行)和一个 smallDataFrame(单列和 10,000 行)。 只要 largeDataFrame 中的 some_identifier 列
我有一个函数,可以在其中规范化 DataFrame 的前 N 列。我想返回规范化的 DataFrame,但不要管原来的。然而,该函数似乎也会对传递的 DataFrame 进行变异! using D
我想在 Scala 中使用指定架构在 DataFrame 上创建。我尝试过使用 JSON 读取(我的意思是读取空文件),但我认为这不是最佳实践。 最佳答案 假设您想要一个具有以下架构的数据框: roo
我正在尝试从数据框中删除一些列,并且不希望返回修改后的数据框并将其重新分配给旧数据框。相反,我希望该函数只修改数据框。这是我尝试过的,但它似乎并没有做我所除外的事情。我的印象是参数是作为引用传递的,而
我有一个包含大约 60000 个数据的庞大数据集。我会首先使用一些标准对整个数据集进行分组,接下来我要做的是将整个数据集分成标准内的许多小数据集,并自动对每个小数据集运行一个函数以获取参数对于每个小数
我遇到了以下问题,并有一个想法来解决它,但没有成功:我有一个月内每个交易日的 DAX 看涨期权和看跌期权数据。经过转换和一些计算后,我有以下 DataFrame: DaxOpt 。现在的目标是消除没有
我正在尝试做一些我认为应该是单行的事情,但我正在努力把它做好。 我有一个大数据框,我们称之为lg,还有一个小数据框,我们称之为sm。每个数据帧都有一个 start 和一个 end 列,以及多个其他列所
我有一个像这样的系列数据帧的数据帧: state1 state2 state3 ... sym1 sym
我有一个大约有 9k 行和 57 列的数据框,这是“df”。 我需要一个新的数据框:'df_final'- 对于“df”的每一行,我必须将每一行复制“x”次,并将每一行中的日期逐一增加,也就是“x”次
假设有一个 csv 文件如下: # data.csv 0,1,2,3,4 a,3.0,3.0,3.0,3.0,3.0 b,3.0,3.0,3.0,3.0,3.0 c,3.0,3.0,3.0,3.0,3
我只想知道是否有人对以下问题有更优雅的解决方案: 我有两个 Pandas DataFrame: import pandas as pd df1 = pd.DataFrame([[1, 2, 3], [
我有一个 pyspark 数据框,我需要将其转换为 python 字典。 下面的代码是可重现的: from pyspark.sql import Row rdd = sc.parallelize([R
我有一个 DataFrame,我想在 @chain 的帮助下对其进行处理。如何存储中间结果? using DataFrames, Chain df = DataFrame(a = [1,1,2,2,2
我有一个包含 3 列的 DataFrame,名为 :x :y 和 :z,它们是 Float64 类型。 :x 和 "y 在 (0,1) 上是 iid uniform 并且 z 是 x 和 y 的总和。
这个问题在这里已经有了答案: pyspark dataframe filter or include based on list (3 个答案) 关闭 2 年前。 只是想知道是否有任何有效的方法来过
我刚找到这个包FreqTables ,它允许人们轻松地从 DataFrames 构建频率表(我正在使用 DataFrames.jl)。 以下代码行返回一个频率表: df = CSV.read("exa
是否有一种快速的方法可以为 sort 指定自定义订单?/sort!在 Julia DataFrames 上? julia> using DataFrames julia> srand(1); juli
在 Python Pandas 和 R 中,可以轻松去除重复的列 - 只需加载数据、分配列名,然后选择那些不重复的列。 使用 Julia Dataframes 处理此类数据的最佳实践是什么?此处不允许
我是一名优秀的程序员,十分优秀!