【问题标题】:Aggregate a Dask dataframe and produce a dataframe of aggregates聚合 Dask 数据帧并生成聚合数据帧
【发布时间】:2018-03-04 15:54:16
【问题描述】:

我有一个如下所示的 Dask 数据框:

url     referrer    session_id ts                  customer
url1    ref1        xxx        2017-09-15 00:00:00 a.com
url2    ref2        yyy        2017-09-15 00:00:00 a.com
url2    ref3        yyy        2017-09-15 00:00:00 a.com
url1    ref1        xxx        2017-09-15 01:00:00 a.com
url2    ref2        yyy        2017-09-15 01:00:00 a.com

我想对 url 和时间戳上的数据进行分组,聚合列值并生成一个看起来像这样的数据框:

customer url    ts                  page_views visitors referrers
a.com    url1   2017-09-15 00:00:00 1          1        [ref1]
a.com    url2   2017-09-15 00:00:00 2          2        [ref2, ref3]

在 Spark SQL 中,我可以这样做:

select 
    customer,
    url,
    ts,
    count(*) as page_views,
    count(distinct(session_id)) as visitors,
    collect_list(referrer) as referrers
from df
group by customer, url, ts

有什么方法可以使用 Dask 数据帧来实现吗?我试过了,但是只能单独计算聚合列,如下:

# group on timestamp (rounded) and url
grouped = df.groupby(['ts', 'url'])

# calculate page views (count rows in each group)
page_views = grouped.size()

# collect a list of referrer strings per group
referrers = grouped['referrer'].apply(list, meta=('referrers', 'f8'))

# count unique visitors (session ids)
visitors = grouped['session_id'].count()

但我似乎找不到生成我需要的组合数据框的好方法。

【问题讨论】:

  • 在 Pandas 中有没有很好的方法来做到这一点?这种方式适用于 dask.dataframe 吗?

标签: group-by aggregation dask


【解决方案1】:

以下确实有效:

gb = df.groupby(['customer', 'url', 'ts'])
gb.apply(lambda d: pd.DataFrame({'views': len(d), 
     'visitiors': d.session_id.count(), 
     'referrers': [d.referer.tolist()]})).reset_index()

(假设访问者按照上面的 sql 应该是唯一的) 您可能希望定义输出的meta。

【讨论】:

  • 不错!如果我从我的数据中构造一个pd.DataFrame,它会强制所有数据进入一台机器上的内存吗?现在它只是一个玩具示例,但真正的工作是处理千兆字节的分布式数据。
  • 它似乎与您的数据完全一样;你应该尝试提供一个元参数dask.pydata.org/en/latest/…
  • 你说得对,它完全按照我在本示例中指定的数据处理数据。它不适用于从分区镶木地板读取的数据稍大的示例。我想弄清楚那个到底有什么问题——我会用我的数据样本在 dask 中提出一个问题。 Stackoverflow 看起来不是一个好地方。谢谢!
  • 在这个例子中你如何提供一个正确的meta?我尝试使用meta = {'views': int, 'visitors': int, 'referrers': object},但没有成功。
【解决方案2】:

这是@j-bennet 打开的link to the github issue,它提供了一个额外的选项。基于这个问题,我们实现了如下聚合:
custom_agg = dd.Aggregation( 'custom_agg', lambda s: s.apply(set), lambda s: s.apply(lambda chunks: list(set(itertools.chain.from_iterable(chunks)))), ).
为了与计数结合,代码如下
dfgp = df.groupby(['ID1','ID2']) df2 = dfgp.assign(cnt=dfgp.size()).agg(custom_agg).reset_index()

【讨论】:

    猜你喜欢
    • 2017-05-12
    • 1970-01-01
    • 1970-01-01
    • 2018-09-04
    • 2022-11-18
    • 1970-01-01
    • 2016-10-09
    • 2016-05-24
    • 2015-08-12
    相关资源
    最近更新 更多