Масштабирование до больших наборов данных
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'>
Все эти примеры с 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