【问题标题】:Dask filter on DataFrame index using a negated range使用否定范围对 DataFrame 索引进行 Dask 过滤器
【发布时间】:2019-02-20 09:52:36
【问题描述】:

我的用例是我每天处理约 100MB。我将 Pandas DataFrame 用作单个文件,但这失败了,因为 pandas 倾向于强制依赖于不同日期的数据的 dtypes。我尝试使用 Dask DataFrame 阅读这些内容,但由于架构不同而失败。具有描述性的列名和 717 列的异常消息是无法管理的(100KB 的固定长度的密集二进制字符串)。

所以我尝试使用 Dask 写出一个巨型镶木地板,并希望它能够解决 pandas dtype skullduggery。有时我需要在我已经拥有的全部天数数据的中间重新处理一两天。

到目前为止,我设法想出了这个,它非常丑陋,我不禁想到有更好的方法。似乎没有办法在 read_parquet 中使用过滤器,因为我们根据索引进行过滤。似乎没有一种方法可以否定索引值的范围。索引只是一个日期,没有小时等。df 是当天的数据价值,mdf 是我的 mega-df,里面有一年的数据

            mdf = dd.read_parquet(self.local_location + self.megafile, engine='pyarrow')
            inx = df.index.unique()
            start1 = '2016-01-01'
            end1 = pd.to_datetime(inx.values.min()).strftime('%Y-%m-%d')
            start2 = pd.to_datetime(inx.values.max()).strftime('%Y-%m-%d')
            end2 = '2029-01-01'
            mdf1 = mdf[start1:end1]
            mdf2 = mdf[start2:end2]
            if len(mdf1) > 0:
                df_usage1 = 1 + mdf1.memory_usage(deep=True).sum().compute() // 100000001

                if len(mdf2) > 0:
                    df_usage2 = 1 + mdf1.memory_usage(deep=True).sum().compute() // 100000001
                    mdf1 = mdf1.append(mdf2, npartitions=df_usage2)
            else:
                if len(mdf2) > 0:
                    df_usage2 = 1 + mdf2.memory_usage(deep=True).sum().compute() // 100000001
                    mdf1 = dd.from_pandas(df).append(mdf2, npartitions=df_usage2)

这也会在

处引发异常
mdf1 = mdf1.append(df, npartitions=df_usage1)

{ValueError}Exactly one of npartitions and chunksize must be specified.

这很有趣,因为这正是我正在做的事情。

df_usage2 在这种情况下 = 2

寻求更好的替代方法,并可能解释附加的实际错误。

【问题讨论】:

  • 对于mega-dfskullduggery 的行话大声笑

标签: python-3.x pandas dask


【解决方案1】:

我建议不要提供npartitions= 关键字

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-07-08
    • 1970-01-01
    • 2018-05-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多