0Pricing
Pandas & NumPy Academy · Урок

Введение в Dask DataFrames

Заменяйте pd.read_csv и pd.DataFrame эквивалентами Dask, вызывайте compute() для запуска вычислений и анализируйте графы задач.

«Введение в Dask DataFrames» — бесплатный урок Pandas & NumPy Academy на CoddyKit. Это урок 3 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Pandas & NumPy Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Pandas & NumPy Academy содержит 4 уроков всего.

Что такое Dask?

Dask — это библиотека параллельных вычислений для Python, расширяющая возможности NumPy и Pandas для наборов данных, превышающих объём RAM. Её модуль dask.dataframe предоставляет API DataFrame, почти идентичный API Pandas, но вместо немедленного выполнения операций Dask строит граф задач и выполняет его отложенно при вызове .compute(). Это позволяет Dask распараллеливать работу между несколькими ядрами или даже несколькими компьютерами, практически не меняя код.

Установка и импорт Dask

Dask устанавливается с помощью команды pip install dask[dataframe]. Стандартная конструкция импорта — import dask.dataframe as dd. Внутри Dask DataFrame разбивается на множество небольших DataFrames Pandas, каждый из которых обрабатывается независимо. Операции над Dask DataFrame создают отложенный граф задач — ничего не выполняется, пока не будет вызван .compute(). Такое разделение между описанием и выполнением вычислений — ключевая идея Dask.

import dask.dataframe as dd

# Read a large CSV — returns a Dask DataFrame immediately (no data loaded yet)
ddf = dd.read_csv('large_sales.csv')
print(type(ddf))     # dask.dataframe.core.DataFrame
print(ddf.columns.tolist())
print(ddf.dtypes)

Dask и Pandas: ключевое различие

В Pandas каждая операция выполняется немедленно и активно. В Dask операции возвращают другой объект Dask, представляющий отложенное вычисление. Только при вызове .compute() Dask фактически считывает данные и выполняет граф задач. Такая отложенность позволяет Dask оптимизировать план перед выполнением: например, объединять последовательные фильтры, чтобы не загружать данные несколько раз. Представьте это как рецепт: Dask записывает рецепт, а .compute() готовит блюдо.

import dask.dataframe as dd

ddf = dd.read_csv('sales.csv')

# This does NOT run yet — just builds the task graph
filtered = ddf[ddf['amount'] > 1000]
agg = filtered.groupby('region')['amount'].sum()

print(type(agg))  # dask.dataframe.core.Series

# NOW execute everything
result = agg.compute()
print(result)

Фрагменты: основная концепция

Dask DataFrame разделён на партиции, каждая из которых представляет собой обычный Pandas DataFrame. По умолчанию dd.read_csv создаёт одну партицию на файл (или одну на каждые 128 MB для больших файлов). Это можно настроить с помощью blocksize. Проверка ddf.npartitions показывает количество существующих партиций. Большее количество партиций обеспечивает более высокий уровень параллелизма, но увеличивает накладные расходы; меньшее количество снижает накладные расходы, но ограничивает параллелизм. Обычно оптимальным считается несколько сотен партиций.

import dask.dataframe as dd

ddf = dd.read_csv('data/*.csv')  # Read multiple CSV files at once
print('Number of partitions:', ddf.npartitions)

# Access a single partition as a Pandas DataFrame
first_partition = ddf.get_partition(0).compute()
print('Partition 0 shape:', first_partition.shape)

Знакомые операции Pandas в Dask

Большинство распространённых операций Pandas работают в Dask так же: .head(), .tail(), .describe(), булева индексация, .groupby(), .merge() и .assign() имеют эквиваленты в Dask. Главное отличие заключается в том, что для материализации результата необходимо вызвать .compute(). Операции, которые Pandas выполняет за миллисекунды, в Dask могут занимать секунды из-за накладных расходов графа задач — поэтому используйте Pandas для небольших объёмов данных, а Dask — когда данные не помещаются в RAM.

import dask.dataframe as dd

ddf = dd.read_csv('transactions.csv')

# Filtering — same syntax as Pandas
high_value = ddf[ddf['amount'] > 500]

# GroupBy aggregation
by_region = high_value.groupby('region')['amount'].mean()

# Execute
result = by_region.compute()
print(result.sort_values(ascending=False))

Чтение нескольких файлов с шаблонами Glob

Одна из самых полезных возможностей Dask — чтение нескольких файлов одновременно с помощью шаблонов glob. dd.read_csv('data/2024-*.csv') читает все подходящие файлы и создаёт одну партицию на файл. Это идеально подходит для данных, хранящихся в виде файлов, разделённых по месяцам или дням, — распространённого шаблона в озёрах данных. Dask автоматически согласует схемы. Это эквивалентно ручному перебору файлов и их объединению с помощью Pandas, но гораздо проще.

import dask.dataframe as dd

# Read all monthly files at once
ddf = dd.read_csv('sales/2024-*.csv',
                  dtype={'order_id': 'int32',
                         'amount': 'float32'})
print(f'Partitions: {ddf.npartitions}')  # One per file
print(f'Total rows (lazy): {len(ddf)}')  # This triggers a compute!

Метод visualize() для графов задач

Перед выполнением сложного конвейера Dask можно проверить граф задач, вызвав result.visualize(). Этот вызов создаёт схему PNG со всеми этапами вычислений. Это полезно для понимания того, что именно будет выполнять Dask, и для поиска причин неожиданной медленной работы. На графе показано, как партиции проходят через этапы фильтрации, группировки и агрегации, поэтому избыточные вычисления легко обнаружить. Требуется пакет graphviz.

import dask.dataframe as dd

ddf = dd.read_csv('orders.csv')
pipeline = (
    ddf[ddf['status'] == 'completed']
    .groupby('product_id')['revenue']
    .sum()
)

# Visualise the task graph (saves to PNG)
# pipeline.visualize('task_graph.png')

# Check number of tasks in the graph
print('Number of tasks:', len(pipeline.__dask_graph__()))

Применение пользовательских функций с помощью map_partitions

Если необходимо применить пользовательскую функцию Pandas к Dask DataFrame, используйте ddf.map_partitions(func). Эта функция применяет func к каждой партиции независимо и возвращает новый Dask DataFrame. Функция получает обычный Pandas DataFrame и должна вернуть такой же объект. Это эквивалент Dask для df.apply() и способ интегрировать Dask с кодом, который понимает только Pandas.

import dask.dataframe as dd
import pandas as pd

def normalise_chunk(df):
    df = df.copy()
    df['amount_norm'] = (df['amount'] - df['amount'].mean()) / df['amount'].std()
    return df

ddf = dd.read_csv('data.csv')
normalised = ddf.map_partitions(normalise_chunk)
result = normalised[['order_id', 'amount_norm']].compute()
print(result.head())

Варианты планировщика Dask

В Dask есть несколько планировщиков, управляющих выполнением задач. Планировщик 'synchronous' выполняет задачи последовательно в текущем потоке (это полезно при отладке). Планировщик 'threads' использует пул потоков (подходит для задач, ограниченных операциями ввода-вывода). Планировщик 'processes' запускает несколько процессов для задач, ограниченных производительностью CPU (обходит GIL Python). Распределённый кластер Dask позволяет выполнять задачи на нескольких машинах. Укажите планировщик с помощью compute(scheduler='threads').

import dask.dataframe as dd

ddf = dd.read_csv('data.csv')
agg = ddf.groupby('category')['sales'].sum()

# Choose scheduler based on workload
result_sync = agg.compute(scheduler='synchronous')  # sequential, easy to debug
result_threads = agg.compute(scheduler='threads')   # parallel I/O
result_processes = agg.compute(scheduler='processes')  # parallel CPU

Преобразование между Dask и Pandas

Часто большой набор данных обрабатывают с помощью Dask, а затем переносят агрегированный результат в Pandas для окончательного анализа или визуализации. Используйте .compute(), чтобы преобразовать Dask DataFrame в Pandas. В обратном направлении вызов dd.from_pandas(df, npartitions=4) преобразует Pandas DataFrame в Dask DataFrame. Это удобно для тестирования кода Dask на небольших объёмах данных перед масштабированием до полного набора данных.

import pandas as pd
import dask.dataframe as dd

# Start with a small Pandas DF for testing
df_small = pd.DataFrame({'a': range(100), 'b': range(100, 200)})

# Convert to Dask for development/testing
ddf = dd.from_pandas(df_small, npartitions=4)
result = ddf.groupby('a')['b'].sum().compute()
print(type(result))  # pandas.Series
print(result.head())

Когда использовать Dask, Pandas и SQL

Dask подходит не всегда. Используйте Pandas, когда данные помещаются в RAM (меньше нескольких GB): благодаря меньшим накладным расходам этот инструмент проще и быстрее. Используйте Dask, когда объём данных превышает RAM, но Вам нужен синтаксис, похожий на Pandas, и параллельная обработка на одной машине. Используйте SQL или базу данных, когда данные находятся в реляционной базе данных и агрегации можно передать механизму базы данных. Используйте Spark или BigQuery, когда требуется распределённая обработка на нескольких машинах для данных петабайтного масштаба.

Быстрая проверка

Проверьте, насколько хорошо Вы усвоили концепции анализа данных из этого урока.

Повторение урока

В этом уроке Вы узнали, что Dask DataFrames — это лениво вычисляемые коллекции партиций Pandas со знакомым API, вызов .compute() запускает фактическое выполнение графа задач, а map_partitions позволяет применить любую пользовательскую функцию Pandas ко всем партициям. Далее мы рассмотрим формат Parquet — быстрый столбцовый вариант CSV для хранения больших наборов данных.

Часто задаваемые вопросы

Урок «Введение в Dask DataFrames» бесплатный?

Да — полный текст урока «Введение в Dask DataFrames» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Pandas & NumPy Academy, подпишись на CoddyKit PRO. Курс Pandas & NumPy Academy содержит 4 уроков всего.

Чему я научусь в уроке «Введение в Dask DataFrames»?

Заменяйте pd.read_csv и pd.DataFrame эквивалентами Dask, вызывайте compute() для запуска вычислений и анализируйте графы задач. Ты практикуешь Pandas & NumPy Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать Pandas & NumPy Academy?

Предыдущий опыт не требуется. Pandas & NumPy Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 3 из 4.

Сколько времени занимает урок «Введение в Dask DataFrames»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке Pandas & NumPy Academy?

Да. Каждый урок Pandas & NumPy Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Потоковое чтение CSV с chunksize
  2. Постепенное агрегирование по частям
  3. Введение в Dask DataFrames
  4. Parquet: быстрое столбцовое хранилище
← Назад к Pandas & NumPy Academy