Spec-Zone.ru › pandas 1

Масштабирование до больших наборов данных

pandas предоставляет структуры данных для аналитики в оперативной памяти, что делает использование pandas для анализа наборов данных, превышающих объём оперативной памяти, несколько сложным. Даже наборы данных, составляющие значительную часть памяти, становятся неуправляемыми, поскольку некоторые операции pandas требуют создания промежуточных копий.

Этот документ предоставляет несколько рекомендаций по масштабированию вашего анализа до больших наборов данных. Он дополняет Улучшение производительности, который фокусируется на ускорении анализа для наборов данных, помещающихся в оперативную память.

Но сначала стоит рассмотреть отказ от использования pandas. pandas не является правильным инструментом для всех ситуаций. Если вы работаете с очень большими наборами данных и такой инструмент, как PostgreSQL, подходит вам, то вы, вероятно, должны использовать его. Предполагая, что вы хотите или нуждаетесь в выразительности и мощности pandas, давайте продолжим.

Загрузка меньшего объёма данных

Предположим, что наш исходный набор данных на диске имеет множество столбцов:

                     id_0    name_0       x_0       y_0  id_1   name_1       x_1  ...  name_8       x_8       y_8  id_9   name_9       x_9       y_9
timestamp                                                                         ...
2000-01-01 00:00:00  1015   Michael -0.399453  0.095427   994    Frank -0.176842  ...     Dan -0.315310  0.713892  1025   Victor -0.135779  0.346801
2000-01-01 00:01:00   969  Patricia  0.650773 -0.874275  1003    Laura  0.459153  ...  Ursula  0.913244 -0.630308  1047    Wendy -0.886285  0.035852
2000-01-01 00:02:00  1016    Victor -0.721465 -0.584710  1046  Michael  0.524994  ...     Ray -0.656593  0.692568  1064   Yvonne  0.070426  0.432047
2000-01-01 00:03:00   939     Alice -0.746004 -0.908008   996   Ingrid -0.414523  ...   Jerry -0.958994  0.608210   978    Wendy  0.855949 -0.648988
2000-01-01 00:04:00  1017       Dan  0.919451 -0.803504  1048    Jerry -0.569235  ...   Frank -0.577022 -0.409088   994      Bob -0.270132  0.335176
...                   ...       ...       ...       ...   ...      ...       ...  ...     ...       ...       ...   ...      ...       ...       ...
2000-12-30 23:56:00   999       Tim  0.162578  0.512817   973    Kevin -0.403352  ...     Tim -0.380415  0.008097  1041  Charlie  0.191477 -0.599519
2000-12-30 23:57:00   970     Laura -0.433586 -0.600289   958   Oliver -0.966577  ...   Zelda  0.971274  0.402032  1038   Ursula  0.574016 -0.930992
2000-12-30 23:58:00  1065     Edith  0.232211 -0.454540   971      Tim  0.158484  ...   Alice -0.222079 -0.919274  1022      Dan  0.031345 -0.657755
2000-12-30 23:59:00  1019    Ingrid  0.322208 -0.615974   981   Hannah  0.607517  ...   Sarah -0.424440 -0.117274   990   George -0.375530  0.563312
2000-12-31 00:00:00   937    Ursula -0.906523  0.943178  1018    Alice -0.564513  ...   Jerry  0.236837  0.807650   985   Oliver  0.777642  0.783392

[525601 rows x 40 columns]

Это можно сгенерировать с помощью следующего фрагмента кода:

In [1]: import pandas as pd

In [2]: import numpy as np

In [3]: def make_timeseries(start="2000-01-01", end="2000-12-31", freq="1D", seed=None):
   ...:     index = pd.date_range(start=start, end=end, freq=freq, name="timestamp")
   ...:     n = len(index)
   ...:     state = np.random.RandomState(seed)
   ...:     columns = {
   ...:         "name": state.choice(["Alice", "Bob", "Charlie"], size=n),
   ...:         "id": state.poisson(1000, size=n),
   ...:         "x": state.rand(n) * 2 - 1,
   ...:         "y": state.rand(n) * 2 - 1,
   ...:     }
   ...:     df = pd.DataFrame(columns, index=index, columns=sorted(columns))
   ...:     if df.index[-1] == end:
   ...:         df = df.iloc[:-1]
   ...:     return df
   ...: 

In [4]: timeseries = [
   ...:     make_timeseries(freq="1T", seed=i).rename(columns=lambda x: f"{x}_{i}")
   ...:     for i in range(10)
   ...: ]
   ...: 

In [5]: ts_wide = pd.concat(timeseries, axis=1)

In [6]: ts_wide.to_parquet("timeseries_wide.parquet")

Чтобы загрузить нужные столбцы, у нас есть два варианта. Вариант 1 загружает все данные, а затем фильтрует их до нужного набора.

In [7]: columns = ["id_0", "name_0", "x_0", "y_0"]

In [8]: pd.read_parquet("timeseries_wide.parquet")[columns]
Out[8]: 
                     id_0 name_0       x_0       y_0
timestamp                                           
2000-01-01 00:00:00   977  Alice -0.821225  0.906222
2000-01-01 00:01:00  1018    Bob -0.219182  0.350855
2000-01-01 00:02:00   927  Alice  0.660908 -0.798511
2000-01-01 00:03:00   997    Bob -0.852458  0.735260
2000-01-01 00:04:00   965    Bob  0.717283  0.393391
...                   ...    ...       ...       ...
2000-12-30 23:56:00  1037    Bob -0.814321  0.612836
2000-12-30 23:57:00   980    Bob  0.232195 -0.618828
2000-12-30 23:58:00   965  Alice -0.231131  0.026310
2000-12-30 23:59:00   984  Alice  0.942819  0.853128
2000-12-31 00:00:00  1003  Alice  0.201125 -0.136655

[525601 rows x 4 columns]

Вариант 2 загружает только запрошенные столбцы.

In [9]: pd.read_parquet("timeseries_wide.parquet", columns=columns)
Out[9]: 
                     id_0 name_0       x_0       y_0
timestamp                                           
2000-01-01 00:00:00   977  Alice -0.821225  0.906222
2000-01-01 00:01:00  1018    Bob -0.219182  0.350855
2000-01-01 00:02:00   927  Alice  0.660908 -0.798511
2000-01-01 00:03:00   997    Bob -0.852458  0.735260
2000-01-01 00:04:00   965    Bob  0.717283  0.393391
...                   ...    ...       ...       ...
2000-12-30 23:56:00  1037    Bob -0.814321  0.612836
2000-12-30 23:57:00   980    Bob  0.232195 -0.618828
2000-12-30 23:58:00   965  Alice -0.231131  0.026310
2000-12-30 23:59:00   984  Alice  0.942819  0.853128
2000-12-31 00:00:00  1003  Alice  0.201125 -0.136655

[525601 rows x 4 columns]

Если мы измерим использование памяти в обоих случаях, мы увидим, что указание columns в данном случае использует примерно в 10 раз меньше памяти.

С помощью pandas.read_csv() вы можете указать usecols для ограничения столбцов, считываемых в оперативную память. Не все форматы файлов, которые могут быть прочитаны pandas, предоставляют возможность чтения подмножества столбцов.

Использование эффективных типов данных

По умолчанию типы данных pandas не являются наиболее эффективными с точки зрения использования памяти. Это особенно верно для столбцов текстовых данных с относительно небольшим количеством уникальных значений (обычно называемых данными с «низкой кардинальностью»). Используя более эффективные типы данных, вы можете хранить более крупные наборы данных в оперативной памяти.

In [10]: ts = make_timeseries(freq="30S", seed=0)

In [11]: ts.to_parquet("timeseries.parquet")

In [12]: ts = pd.read_parquet("timeseries.parquet")

In [13]: ts
Out[13]: 
                       id     name         x         y
timestamp                                             
2000-01-01 00:00:00  1041    Alice  0.889987  0.281011
2000-01-01 00:00:30   988      Bob -0.455299  0.488153
2000-01-01 00:01:00  1018    Alice  0.096061  0.580473
2000-01-01 00:01:30   992      Bob  0.142482  0.041665
2000-01-01 00:02:00   960      Bob -0.036235  0.802159
...                   ...      ...       ...       ...
2000-12-30 23:58:00  1022    Alice  0.266191  0.875579
2000-12-30 23:58:30   974    Alice -0.009826  0.413686
2000-12-30 23:59:00  1028  Charlie  0.307108 -0.656789
2000-12-30 23:59:30  1002    Alice  0.202602  0.541335
2000-12-31 00:00:00   987    Alice  0.200832  0.615972

[1051201 rows x 4 columns]

Теперь давайте рассмотрим типы данных и использование памяти, чтобы определить, на чём следует сосредоточить внимание.

In [14]: ts.dtypes
Out[14]: 
id        int64
name     object
x       float64
y       float64
dtype: object
In [15]: ts.memory_usage(deep=True)  # memory usage in bytes
Out[15]: 
Index     8409608
id        8409608
name     65176434
x         8409608
y         8409608
dtype: int64

Столбец name занимает намного больше памяти, чем любой другой. В нём всего несколько уникальных значений, поэтому это хороший кандидат для преобразования в pandas.Categorical. С помощью pandas.Categorical мы храним каждое уникальное имя один раз и используем эффективные с точки зрения места целые числа, чтобы знать, какое именно имя используется в каждой строке.

In [16]: ts2 = ts.copy()

In [17]: ts2["name"] = ts2["name"].astype("category")

In [18]: ts2.memory_usage(deep=True)
Out[18]: 
Index    8409608
id       8409608
name     1051495
x        8409608
y        8409608
dtype: int64

Мы можем продвинуться дальше и понизить числовые столбцы до их наименьших типов с помощью pandas.to_numeric().

In [19]: ts2["id"] = pd.to_numeric(ts2["id"], downcast="unsigned")

In [20]: ts2[["x", "y"]] = ts2[["x", "y"]].apply(pd.to_numeric, downcast="float")

In [21]: ts2.dtypes
Out[21]: 
id        uint16
name    category
x        float32
y        float32
dtype: object
In [22]: ts2.memory_usage(deep=True)
Out[22]: 
Index    8409608
id       2102402
name     1051495
x        4204804
y        4204804
dtype: int64
In [23]: reduction = ts2.memory_usage(deep=True).sum() / ts.memory_usage(deep=True).sum()

In [24]: print(f"{reduction:0.2f}")
0.20

В итоге мы сократили занимаемое этой областью данных место в оперативной памяти до 1/5 от первоначального размера.

См. Данные категорий для получения дополнительной информации о pandas.Categorical и типы данных для общего обзора всех типов данных pandas.

Использование чанков

Некоторые рабочие нагрузки могут быть реализованы с помощью чанков: разделение большой проблемы, такой как «преобразование этого каталога CSV в Parquet», на множество небольших проблем («преобразование этого файла CSV в файл Parquet. Теперь повторите это для каждого файла в этом каталоге»). До тех пор, пока каждый чанк помещается в оперативную память, вы можете работать с наборами данных, значительно превышающими объём оперативной памяти.

Примечание

Чанки хорошо подходят, когда операция, которую вы выполняете, требует нулевого или минимального взаимодействия между чанками. Для более сложных рабочих процессов лучше использовать другую библиотеку.

Предположим, что у нас есть ещё больший «логический набор данных» на диске, который представляет собой каталог файлов Parquet. Каждый файл в каталоге представляет собой разные годы всего набора данных.

In [25]: import pathlib

In [26]: N = 12

In [27]: starts = [f"20{i:>02d}-01-01" for i in range(N)]

In [28]: ends = [f"20{i:>02d}-12-13" for i in range(N)]

In [29]: pathlib.Path("data/timeseries").mkdir(exist_ok=True)

In [30]: for i, (start, end) in enumerate(zip(starts, ends)):
   ....:     ts = make_timeseries(start=start, end=end, freq="1T", seed=i)
   ....:     ts.to_parquet(f"data/timeseries/ts-{i:0>2d}.parquet")
   ....: 
data
└── timeseries
    ├── ts-00.parquet
    ├── ts-01.parquet
    ├── ts-02.parquet
    ├── ts-03.parquet
    ├── ts-04.parquet
    ├── ts-05.parquet
    ├── ts-06.parquet
    ├── ts-07.parquet
    ├── ts-08.parquet
    ├── ts-09.parquet
    ├── ts-10.parquet
    └── ts-11.parquet

Теперь мы реализуем внеоператорную pandas.Series.value_counts(). Пиковое использование памяти этой рабочей области равно самому большому чанку плюс небольшому ряду, хранящему счёт уникальных значений до этого момента. До тех пор, пока каждый отдельный файл помещается в оперативную память, это будет работать для произвольно больших наборов данных.

In [31]: %%time
   ....: files = pathlib.Path("data/timeseries/").glob("ts*.parquet")
   ....: counts = pd.Series(dtype=int)
   ....: for path in files:
   ....:     df = pd.read_parquet(path)
   ....:     counts = counts.add(df["name"].value_counts(), fill_value=0)
   ....: counts.astype(int)
   ....: 
CPU times: user 893 ms, sys: 91.5 ms, total: 984 ms
Wall time: 960 ms
Out[31]: 
Alice      1994645
Bob        1993692
Charlie    1994875
dtype: int64

Некоторые читатели, такие как pandas.read_csv(), предлагают параметры для управления chunksize при чтении одного файла.

Ручное использование чанков — это приемлемый вариант для рабочих процессов, не требующих слишком сложных операций. Некоторые операции, такие как pandas.DataFrame.groupby(), гораздо сложнее выполнять по частям. В этих случаях вам может быть лучше перейти к другой библиотеке, которая реализует эти алгоритмы вне оперативной памяти за вас.

Использование других библиотек

pandas — всего лишь одна библиотека, предлагающая API DataFrame. Благодаря своей популярности, API pandas стал своего рода стандартом, который реализуют другие библиотеки. Документация pandas содержит список библиотек, реализующих API DataFrame на странице нашей экосистемы.

Например, библиотека параллельного вычисления Dask имеет dask.dataframe, API, подобный pandas, для работы с наборами данных, превышающими объём памяти, в параллельном режиме. Dask может использовать несколько потоков или процессов на одном компьютере или кластере компьютеров для параллельной обработки данных.

Мы импортируем dask.dataframe и заметим, что API похож на pandas. Мы можем использовать функцию read_parquet Dask, но предоставить строку шаблона файлов для чтения.

In [32]: import dask.dataframe as dd

In [33]: ddf = dd.read_parquet("data/timeseries/ts*.parquet", engine="pyarrow")

In [34]: ddf
Out[34]: 
Dask DataFrame Structure:
                   id    name        x        y
npartitions=12                                 
                int64  object  float64  float64
                  ...     ...      ...      ...
...               ...     ...      ...      ...
                  ...     ...      ...      ...
                  ...     ...      ...      ...
Dask Name: read-parquet, 1 graph layer

Рассмотрев объект ddf, мы видим несколько моментов:

  • Есть знакомые атрибуты, такие как .columns и .dtypes

  • Есть знакомые методы, такие как .groupby, .sum, и т. д.

  • Есть новые атрибуты, такие как .npartitions и .divisions

Разбиения и разделения — это способ, которым Dask параллелизует вычисления. DataFrame Dask состоит из многих pandas pandas.DataFrame. Один вызов метода DataFrame Dask приводит к выполнению многих вызовов методов pandas, и Dask знает, как координировать все для получения результата.

In [35]: ddf.columns
Out[35]: Index(['id', 'name', 'x', 'y'], dtype='object')

In [36]: ddf.dtypes
Out[36]: 
id        int64
name     object
x       float64
y       float64
dtype: object

In [37]: ddf.npartitions
Out[37]: 12

Одно важное отличие: API dask.dataframe отложенный. Если вы посмотрите на представленный выше вывод, вы заметите, что значения фактически не выводятся; отображаются только имена столбцов и типы данных. Это потому, что Dask ещё не прочитал данные. Вместо немедленного выполнения, операции формируют граф задач.

In [38]: ddf
Out[38]: 
Dask DataFrame Structure:
                   id    name        x        y
npartitions=12                                 
                int64  object  float64  float64
                  ...     ...      ...      ...
...               ...     ...      ...      ...
                  ...     ...      ...      ...
                  ...     ...      ...      ...
Dask Name: read-parquet, 1 graph layer

In [39]: ddf["name"]
Out[39]: 
Dask Series Structure:
npartitions=12
    object
       ...
     ...  
       ...
       ...
Name: name, dtype: object
Dask Name: getitem, 2 graph layers

In [40]: ddf["name"].value_counts()
Out[40]: 
Dask Series Structure:
npartitions=1
    int64
      ...
Name: name, dtype: int64
Dask Name: value-counts-agg, 4 graph layers

Каждый из этих вызовов мгновенен, потому что результат ещё не вычисляется. Мы просто составляем список вычислений, которые нужно выполнить, когда кто-то нуждается в результате. Dask знает, что тип возвращаемого значения pandas.Series.value_counts — pandas pandas.Series с определённым типом данных и именем. Таким образом, версия Dask возвращает Dask Series с тем же типом данных и тем же именем.

Чтобы получить фактический результат, вы можете вызвать .compute().

In [41]: %time ddf["name"].value_counts().compute()
CPU times: user 914 ms, sys: 28.6 ms, total: 943 ms
Wall time: 931 ms
Out[41]: 
Charlie    1994875
Alice      1994645
Bob        1993692
Name: name, dtype: int64

В этот момент вы получите то же, что и с pandas, в данном случае конкретный pandas pandas.Series со счётом каждого name.

Вызов .compute вызывает полное выполнение графа задач. Это включает чтение данных, выбор столбцов и выполнение value_counts. Выполнение происходит параллельно, где это возможно, и Dask пытается сохранить общий объём используемой памяти небольшим. Вы можете работать с наборами данных, значительно превышающими объём памяти, при условии, что каждая часть (обычный pandas pandas.DataFrame) помещается в память.

По умолчанию, операции dask.dataframe используют пул потоков для параллельного выполнения операций. Мы также можем подключиться к кластеру для распределения работы по нескольким машинам. В этом случае мы подключимся к локальному «кластеру», состоящему из нескольких процессов на этом одном компьютере.

>>> from dask.distributed import Client, LocalCluster

>>> cluster = LocalCluster()
>>> client = Client(cluster)
>>> client
<Client: 'tcp://127.0.0.1:53349' processes=4 threads=8, memory=17.18 GB>

После создания этого client все вычисления Dask будут выполняться на кластере (в данном случае — это просто процессы).

Dask реализует наиболее часто используемые части API pandas. Например, мы можем выполнить обычное агрегирование по группам.

In [42]: %time ddf.groupby("name")[["x", "y"]].mean().compute().head()
CPU times: user 2.04 s, sys: 119 ms, total: 2.16 s
Wall time: 1.99 s
Out[42]: 
                x         y
name                       
Alice   -0.000224 -0.000194
Bob     -0.000746  0.000349
Charlie  0.000604  0.000250

Группировка и агрегирование выполняются вне памяти и параллельно.

Когда Dask знает divisions набора данных, возможны определённые оптимизации. При чтении наборов данных parquet, созданных dask, разделения будут определяться автоматически. В данном случае, так как мы создали файлы parquet вручную, нам необходимо вручную указать разделения.

In [43]: N = 12

In [44]: starts = [f"20{i:>02d}-01-01" for i in range(N)]

In [45]: ends = [f"20{i:>02d}-12-13" for i in range(N)]

In [46]: divisions = tuple(pd.to_datetime(starts)) + (pd.Timestamp(ends[-1]),)

In [47]: ddf.divisions = divisions

In [48]: ddf
Out[48]: 
Dask DataFrame Structure:
                   id    name        x        y
npartitions=12                                 
2000-01-01      int64  object  float64  float64
2001-01-01        ...     ...      ...      ...
...               ...     ...      ...      ...
2011-01-01        ...     ...      ...      ...
2011-12-13        ...     ...      ...      ...
Dask Name: read-parquet, 1 graph layer

Теперь мы можем выполнять такие действия, как быстрый произвольный доступ с помощью .loc.

In [49]: ddf.loc["2002-01-01 12:01":"2002-01-01 12:05"].compute()
Out[49]: 
                       id     name         x         y
timestamp                                             
2002-01-01 12:01:00   971      Bob -0.659481  0.556184
2002-01-01 12:02:00  1015  Charlie  0.120131 -0.609522
2002-01-01 12:03:00   991      Bob -0.357816  0.811362
2002-01-01 12:04:00   984    Alice -0.608760  0.034187
2002-01-01 12:05:00   998  Charlie  0.551662 -0.461972

Dask знает, что нужно искать только в третьем разделе для выбора значений в 2002 году. Ему не нужно просматривать никакие другие данные.

Многие рабочие процессы включают в себя большое количество данных и их обработку таким образом, чтобы уменьшить объём до размера, который помещается в память. В этом случае мы перевыведем частоту на ежедневную и возьмём среднее значение. После вычисления среднего значения мы знаем, что результаты поместятся в память, поэтому мы можем безопасно вызвать compute без исчерпания памяти. В этот момент это обычный объект pandas.

In [50]: ddf[["x", "y"]].resample("1D").mean().cumsum().compute().plot()
Out[50]: <AxesSubplot: xlabel='timestamp'>
../_images/dask_resample.png

Все эти примеры с Dask были выполнены с использованием нескольких процессов на одном компьютере. Dask может быть развёрнут на кластере для масштабирования до ещё больших наборов данных.

Вы найдете больше примеров с Dask на https://examples.dask.org.

© 2008–2022, AQR Capital Management, LLC, Lambda Foundry, Inc. and PyData Development Team
Licensed under the 3-clause BSD License.
https://pandas.pydata.org/pandas-docs/version/1.5.0/user_guide/scale.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API