Открыть сервис

DataFrame API

DataFrame API — это программный интерфейс приложения (API) для работы с датафреймами (DataFrame), структурированными наборами данных, организованными в виде таблицы с именованными столбцами и строками. DataFrame API предоставляет набор функций и методов для создания, преобразования, фильтрации, агрегации и анализа данных, абстрагируя пользователя от низкоуровневых деталей хранения и распределённой обработки. Он является ключевым компонентом многих современных систем обработки данных, включая Apache Spark, Pandas (Python), R (dplyr) и Julia (DataFrames.jl), и ориентирован на удобство, производительность и масштабируемость.

История

Концепция датафрейма восходит к языку программирования R и его пакету data.frame, появившемуся в начале 1990-х годов. В 2008 году Уэс Маккинни создал библиотеку Pandas для Python, которая популяризировала датафреймы в экосистеме Python. В 2013 году в рамках проекта Apache Spark был разработан Spark DataFrame API, который объединил идеи реляционных баз данных (SQL) и распределённых вычислений. В 2015 году вышла версия Spark 1.3, где DataFrame API стал стабильным. В последующие годы аналогичные API появились в Julia, Kotlin (Kotlin DataFrame), а также в виде надстроек над SQL-движками (например, DuckDB).

Основные характеристики

DataFrame API отличается от традиционных SQL-запросов и низкоуровневых RDD (Resilient Distributed Dataset) по нескольким параметрам:

  • Декларативность: Пользователь описывает, что нужно сделать (например, отфильтровать строки), а не как это реализовать на уровне распределённых операций. Оптимизатор запросов (например, Catalyst в Spark) автоматически строит эффективный план выполнения.
  • Типизация: Столбцы имеют определённые типы данных (целые, строки, даты, массивы, структуры). API предоставляет строгую проверку типов на этапе компиляции или выполнения, что снижает количество ошибок.
  • Ленивые вычисления: В системах распределённой обработки (Spark, Dask) операции преобразования (например, select, filter) не выполняются немедленно, а формируют логический план. Вычисления запускаются только при вызове «действия» (action) — например, show(), count(), collect().
  • Оптимизация: Встроенные оптимизаторы (Catalyst для Spark, CBO для SQL) переписывают запросы, переставляют фильтры, объединяют проекции и используют колоночное хранение для ускорения.
  • Интеграция с SQL: DataFrame API часто позволяет выполнять SQL-запросы непосредственно к датафрейму, а также регистрировать датафреймы как временные таблицы.

Классификация

DataFrame API можно классифицировать по среде выполнения и языку:

По среде выполнения

  • Однопоточные (in-memory): Pandas, R data.frame, Julia DataFrames.jl. Работают в оперативной памяти одного узла, подходят для малых и средних объёмов данных (до десятков гигабайт).
  • Распределённые (cluster): Apache Spark, Dask, Koalas (теперь Pandas API on Spark), Apache Flink (Table API). Обрабатывают данные, разбитые на партиции, на кластерах из сотен узлов.
  • Гибридные: DuckDB, Polars. Могут работать как в памяти одного узла, так и с использованием колоночных форматов и векторных инструкций, обеспечивая высокую производительность на одном узле.

По языку программирования

  • Python: Pandas, PySpark, Polars, Dask DataFrame.
  • Scala/Java: Spark DataFrame API, Kotlin DataFrame.
  • R: dplyr, data.table, sparklyr.
  • Julia: DataFrames.jl, JuliaDB.
  • SQL: Встроенные функции в PostgreSQL (JSON/JSONB), DuckDB, ClickHouse.

Устройство и компоненты

Основные сущности

  • DataFrame: Таблица с именованными столбцами. Каждый столбец — это серия (Series) или колонка (Column) с однотипными данными.
  • Column: Объект, представляющий столбец. Может быть результатом выражения (например, col("age") + 1).
  • Row: Одна запись (строка) в датафрейме. В распределённых системах строки распределены по партициям.
  • Schema (Схема): Определение структуры датафрейма — имена столбцов и их типы данных. Может быть задана явно или выведена автоматически при чтении данных.

Типы операций

API делится на два класса:

  1. Трансформации (Transformations): Возвращают новый датафрейм, не изменяя исходный. Примеры:
  • select() — выбор столбцов.
  • filter() / where()фильтрация строк по условию.
  • groupBy() — группировка по столбцам.
  • join()объединение двух датафреймов.
  • withColumn() — добавление или замена столбца.
  • orderBy() / sort()сортировка.
  • drop()удаление столбцов.
  • distinct() — удаление дубликатов.
  1. Действия (Actions): Запускают вычисления и возвращают результат (не датафрейм). Примеры:
  • show() — вывод первых N строк.
  • count() — количество строк.
  • collect() — сбор всех данных на драйвер.
  • write() — запись в файл (Parquet, CSV, JSON, JDBC).
  • head() / take()выборка первых строк.

Оптимизация запросов

В распределённых системах (Spark) DataFrame API использует оптимизатор Catalyst, который:

  • Преобразует логический план в физический.
  • Применяет правила переписывания (например, проталкивание фильтров, объединение проекций).
  • Использует статистику для выбора стратегии соединения (broadcast join, sort-merge join).
  • Генерирует байт-код Java для выполнения.

Применение

DataFrame API широко используется в следующих областях:

  • Обработка больших данных: ETL-процессы (извлечение, преобразование, загрузка) в Apache Spark, где датафреймы читаются из HDFS, S3, Kafka, а затем записываются в хранилища.
  • Анализ данных и статистика: Pandas и R — стандартные инструменты для исследовательского анализа, построения отчётов, сводных таблиц.
  • Машинное обучение: DataFrame API служит интерфейсом для передачи данных в библиотеки ML (Spark MLlib, scikit-learn, XGBoost). В Spark ML Pipeline датафреймы используются для хранения признаков и меток.
  • Работа с временными рядами: Фильтрация по датам, оконные функции, агрегация по периодам.
  • Интеграция с SQL: Многие системы позволяют выполнять SQL-запросы к датафреймам, что удобно для аналитиков, не владеющих программированием.

Примеры

Пример 1: Pandas (Python)

```python import pandas as pd

df = pd.DataFrame({'name': ['Alice', 'Bob'], 'age': [25, 30]}) df_filtered = df[df['age'] > 25] print(df_filtered) ```

Пример 2: PySpark (Spark DataFrame API)

```python from pyspark.sql import SparkSession from pyspark.sql.functions import col

spark = SparkSession.builder.getOrCreate() df = spark.read.csv("data.csv", header=True, inferSchema=True) df_filtered = df.filter(col("age") > 25) df_filtered.show() ```

Пример 3: R (dplyr)

``r library(dplyr) df <- data.frame(name = c("Alice", "Bob"), age = c(25, 30)) df_filtered <- df %>% filter(age > 25) print(df_filtered) ``

Сравнение с альтернативами

ХарактеристикаDataFrame API (Spark)SQLRDD (Spark)
Уровень абстракцииВысокийВысокийНизкий
ТипизацияСтрогая, схемаСлабая (типы столбцов)Нет (Java/Scala объекты)
ОптимизацияCatalystОптимизатор SQLРучная
ПроизводительностьВысокая (колоночное хранение, кодогенерация)ВысокаяСредняя (сериализация)
Удобство для аналитиковВысокоеВысокоеНизкое

Критика

  • Сложность отладки: В распределённых системах ленивые вычисления затрудняют поиск ошибок — сообщение об ошибке может появиться только при выполнении действия, далеко от места трансформации.
  • Ограничения на пользовательские функции: В Spark использование пользовательских функций (UDF) на Python приводит к сериализации и десериализации данных, что снижает производительность.
  • Потребление памяти: В Pandas все данные загружаются в память одного узла, что ограничивает размер обрабатываемого набора данных.
  • Неявное поведение: Некоторые операции (например, join с несовпадающими ключами) могут приводить к неожиданным результатам, если не указан тип соединения.

Интересные факты

  • Spark DataFrame API изначально был вдохновлён структурой данных DataFrame в R и Pandas, но адаптирован для распределённых вычислений.
  • В 2019 году компания Databricks представила проект Koalas (теперь Pandas API on Spark), который предоставляет знакомый синтаксис Pandas поверх Spark, чтобы упростить миграцию.
  • Библиотека Polars, написанная на Rust, использует DataFrame API с нулевой копией (zero-copy) и векторными инструкциями, что делает её одной из самых быстрых однопоточных реализаций.

Источники

  • Apache Spark Documentation: DataFrame API
  • Pandas Documentation: DataFrame
  • R Documentation: data.frame
  • Polars Documentation: DataFrame API
  • «Learning Spark» by Jules Damji, Brooke Wenig, Tathagata Das, Denny Lee (O'Reilly, 2020)
  • «Python for Data Analysis» by Wes McKinney (O'Reilly, 2017)

BFOmetr — база данных и аналитика по компаниям России.

На главную BFOmetr →