Інтеграція Apache Spark MLlib для обробки великих даних

Проектуємо та впроваджуємо системи штучного інтелекту: від прототипу до production-ready рішення. Наша команда поєднує експертизу в машинному навчанні, дата-інжинірингу та MLOps, щоб AI працював не в лабораторії, а в реальному бізнесі.
Показано 1 з 1Усі 1564 послуг
Інтеграція Apache Spark MLlib для обробки великих даних
Середній
~1-2 тижні
Часті запитання

Напрямки AI-розробки

Етапи розробки AI-рішення

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1359
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1251
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    957
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1188
  • image_logo-advance_0.webp
    Розробка логотипу компанії B2B Advance
    646
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    929

Ви завантажили 500 ГБ даних, запустили .fit() і через годину отримали OutOfMemoryError. Pandas та scikit-learn упираються в RAM однієї машини. Рішення — розподілене навчання на Spark MLlib. Ми займаємося розподіленим ML на Spark більше 5 років, виконали 20+ проєктів з датасетами до 10 ТБ. Нижче — практичний досвід, який допоможе уникнути типових помилок і швидко запустити навчання на кластері.

Проблеми масштабування та feature engineering

Spark MLlib розподіляє обчислення на кластер, обробляючи дані будь-якого обсягу. Типовий поріг — датасети >100 GB або >100 млн рядків, де pipeline на sklearn перестає працювати. Сотні ознак, категоріальні змінні з мільйонами унікальних значень — StringIndexer і OneHotEncoder в Spark справляються без вивантаження в пам'ять. CrossValidator на Spark запускає фолди паралельно, прискорюючи підбір гіперпараметрів у 3–5 разів. Вбудовані інструменти для роботи з пропусками (Imputer) та масштабування (StandardScaler) дозволяють будувати пайплайни без перемикання контексту.

Порівняння Spark MLlib з альтернативами

Порівняємо з підходами на pandas+sklearn та Dask. Spark виграє за рахунок нативної підтримки розподілених DataFrames, оптимізованих під shuffle, вбудованих алгоритмів (GBT, RandomForest, KMeans) та інтеграції з MLflow для трекінгу експериментів.

Інструмент Час навчання GBT (10 млн записів) Споживання пам'яті Масштабування
scikit-learn ~45 хв 32 GB+ (OOM) Ні
Dask+sklearn ~20 хв 16 GB Обмежене
Spark MLlib ~8 хв 8 GB на executor Горизонтальне

Spark MLlib на 80% швидший при вдвічі меншому споживанні ресурсів на executor.

Як налаштувати Spark MLlib для оптимальної продуктивності?

Ключові прийоми налаштування - Репартиціонування: `df.repartition(200)` перед fit — рівномірне навантаження на executor. - Кешування: `train_df.cache()` прискорює CV у 3–5 разів. - Налаштування shuffle partitions: `spark.sql.shuffle.partitions = 2 * total_cores`. - Паралелізм CV: параметр `parallelism=4` у CrossValidator запускає фолди паралельно.

Ці налаштування скорочують час CV з 6 до 1.5 годин на кластері з 10 executor.

Параметр Default Рекомендоване Ефект
spark.sql.shuffle.partitions 200 2x cores Уникнути skew
executor.memory 1g 4-8g Кеш датасету
spark.ml.param.maxParallelism 1 4-8 CV паралелізм
repartition перед fit ні 200-400 Рівномірне навантаження
caching train_df ні так 3-5x прискорення CV

Чому Spark MLlib швидший за sklearn на великих даних?

Spark MLlib використовує розподілені обчислення та оптимізовані алгоритми для роботи з даними, що не поміщаються в пам'ять. На відміну від sklearn, який завантажує все в RAM, Spark обробляє дані частинами на кластері. Це дозволяє досягти лінійної масштабованості при додаванні вузлів.

Реалізація пайплайну: від конфігурації до деплою

Конфігурація кластера та підготовка даних

Використовуємо PySpark 3.4+, MLflow 2.x, ONNX для інференсу. Конфігурація кластера підбирається під задачу: executor memory від 4 до 16 GB, кількість executor кратна числу партицій. Нижче — робочий пайплайн для бінарної класифікації з градієнтним бустингом.

from pyspark.sql import SparkSession
from pyspark.ml import Pipeline
from pyspark.ml.feature import (VectorAssembler, StringIndexer,
                                  StandardScaler, Imputer)
from pyspark.ml.classification import GBTClassifier, RandomForestClassifier
from pyspark.ml.evaluation import BinaryClassificationEvaluator
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder

spark = SparkSession.builder \
    .appName("ML Pipeline") \
    .config("spark.executor.memory", "8g") \
    .config("spark.executor.core", "4") \
    .config("spark.executor.instances", "10") \
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.ml.param.maxParallelism", "4") \
    .getOrCreate()

# Загрузка данных
df = spark.read.parquet("s3://data/training/*.parquet")
df = df.repartition(200)  # Оптимальное число партиций

# Feature engineering
numeric_cols = ['amount', 'age', 'days_since_last_tx', 'tx_count_30d']
categorical_cols = ['category', 'country', 'device_type']

# Imputer для числовых
imputer = Imputer(
    inputCols=numeric_cols,
    outputCols=[f"{c}_imputed" for c in numeric_cols],
    strategy="median"
)

# Кодирование категориальных
indexers = [
    StringIndexer(inputCol=col, outputCol=f"{col}_idx",
                  handleInvalid="keep")
    for col in categorical_cols
]

# Сборка вектора признаков
all_feature_cols = (
    [f"{c}_imputed" for c in numeric_cols] +
    [f"{c}_idx" for c in categorical_cols]
)

assembler = VectorAssembler(
    inputCols=all_feature_cols,
    outputCol="features_raw",
    handleInvalid="keep"
)

scaler = StandardScaler(
    inputCol="features_raw",
    outputCol="features",
    withMean=True,
    withStd=True
)

# Модель
gbt = GBTClassifier(
    labelCol="label",
    featuresCol="features",
    maxIter=100,
    maxDepth=5,
    stepSize=0.05,
    subsamplingRate=0.8,
    seed=42
)

# Pipeline
pipeline = Pipeline(stages=[
    imputer,
    *indexers,
    assembler,
    scaler,
    gbt
])

# Train/test split
train_df, test_df = df.randomSplit([0.8, 0.2], seed=42)

# Обучение
model = pipeline.fit(train_df)
predictions = model.transform(test_df)

# Оценка
evaluator = BinaryClassificationEvaluator(
    labelCol="label",
    rawPredictionCol="rawPrediction",
    metricName="areaUnderROC"
)
auc = evaluator.evaluate(predictions)
print(f"Test AUC: {auc:.4f}")

Гіперпараметрична оптимізація

# Cross-validation на кластере
param_grid = ParamGridBuilder() \
    .addGrid(gbt.maxDepth, [4, 6, 8]) \
    .addGrid(gbt.maxIter, [50, 100]) \
    .addGrid(gbt.stepSize, [0.05, 0.1]) \
    .build()

cv = CrossValidator(
    estimator=pipeline,
    estimatorParamMaps=param_grid,
    evaluator=evaluator,
    numFolds=3,
    parallelism=4,  # Параллельный запуск фолдов
    seed=42
)

cv_model = cv.fit(train_df)
best_model = cv_model.bestModel
print(f"Best params: {cv_model.bestModel.stages[-1].extractParamMap()}")

Feature Importance та інтерпретація

# Извлечение feature importance
gbt_model = best_model.stages[-1]
importance = gbt_model.featureImportances

# Маппинг на имена признаков
feature_names = all_feature_cols
importance_df = spark.createDataFrame(
    [(name, float(imp)) for name, imp in zip(feature_names, importance.toArray())],
    ["feature", "importance"]
).orderBy("importance", ascending=False)

importance_df.show(20)

# SHAP через pandas на выборке
sample_pandas = predictions.sample(fraction=0.01).toPandas()
# ... далее стандартный TreeExplainer

Збереження та деплой моделі

import mlflow
import mlflow.spark

# Логирование в MLflow
with mlflow.start_run():
    mlflow.log_param("max_depth", gbt.getMaxDepth())
    mlflow.log_param("max_iter", gbt.getMaxIter())
    mlflow.log_metric("auc", auc)

    # Сохранение Spark модели
    mlflow.spark.log_model(best_model, "spark_model")

    # Экспорт в ONNX для быстрого инференса
    from onnxmltools import convert_sparkml
    onnx_model = convert_sparkml(best_model, "GBT Model", test_df.limit(5))
    mlflow.onnx.log_model(onnx_model, "onnx_model")

# Загрузка для предсказаний
loaded_model = mlflow.spark.load_model("runs:/RUN_ID/spark_model")
batch_predictions = loaded_model.transform(new_data_df)

Що входить в роботу

При замовленні інтеграції Spark MLlib ви отримуєте:

  • Документація пайплайну та конфігурацій (включаючи обґрунтування вибору алгоритмів та параметрів).
  • Налаштований MLflow tracking server для відтворюваності експериментів.
  • Docker-образ для інференсу з ONNX Runtime.
  • Навчання команди замовника (1 день) — робота з PySpark, MLflow, оптимізація.
  • Підтримка 2 тижні після деплою — виправлення інцидентів, консультації.

Процес роботи, строки та гарантії

  1. Аналітика: вивчаємо дані, визначаємо ознаки, цільову змінну, метрики.
  2. Проектування: вибираємо алгоритм, конфігурацію кластера, пайплайн.
  3. Реалізація: пишемо код, налаштовуємо MLflow, запускаємо пробне навчання.
  4. Тестування: валідація на відкладеній вибірці, A/B тест у стейджингу.
  5. Деплой: упаковка в MLflow Model, сервінг через REST API або ONNX.

Строки під ключ — від 2 до 4 тижнів залежно від складності даних та алгоритму. Вартість розраховується індивідуально після аудиту вашого датасету та інфраструктури. Отримайте консультацію — оцінимо проєкт безкоштовно.

Багаторічний досвід у промисловому ML на Spark, десятки завершених проєктів з датасетами від 100 GB до 10 TB. Сертифіковані спеціалісти з Spark та MLflow. Гарантуємо відтворюваність результатів.

Замовте інтеграцію Spark MLlib — ми налаштуємо пайплайн під ваш датасет. Зв'яжіться для консультації — допоможемо вибрати оптимальну архітектуру та запустити перший пайплайн за тиждень.

Чому дата-інжиніринг визначає успіх ML-моделі

Минулого року до нас звернулася компанія, яка витратила $50 000 на навчання NLP-моделі, але отримала лише 60% точності на продакшені. Причина — data leakage через випадковий split часових даних. Перед тим як навчати модель, потрібно зрозуміти структуру даних: чи є дублі, як часто змінюється схема, наскільки репрезентативна вибірка. Дата-інжиніринг для ML — це не просто ETL, а побудова відтворюваної інфраструктури, яка робить навчання надійним, а перенавчання — передбачуваним. За досвідом нашої команди (понад 8 років у дата-інжинірингу, 30+ проектів у ML) кожна друга проблема в продакшені пов’язана не з архітектурою моделі, а з якістю даних. Замовте аудит ваших даних — оцінимо поточний пайплайн безкоштовно.

Як ETL-пайплайни для ML відрізняються від BI

ETL для аналітики та ETL для ML — різні завдання. В аналітиці важлива агрегація, у ML — індивідуальні записи з історією. В аналітиці train/val/test split не потрібен, у ML — критичний. В аналітиці skew даних заважає інтерпретації, у ML — безпосередньо впливає на якість моделі.

Інструменти. Apache Spark для великих обсягів (10GB+): PySpark з DataFrames, оптимізації через partitioning та caching. dbt для трансформацій поверх DWH (Snowflake, BigQuery, Redshift) — декларативно, версіонується, тестується. Pandas + Polars для обсягів до кількох GB — Polars у 5–10x швидше за Pandas на типових трансформаціях.

Temporal splits. Для ML важливо, що split за часом, а не випадковий. Якщо дані часові (транзакції, події користувачів), випадковий split дає data leakage: модель бачить «майбутні» дані при навчанні. Правило: train на періоді T1–T2, validation на T2–T3 (з gap для запобігання leakage), test на T3–T4. Неправильний split може коштувати 10–15% якості моделі на валідації. Temporal split best practices (scikit-learn docs)

Інкрементальні пайплайни. Модель перенавчається щотижня на нових даних. Потрібен пайплайн, який інкрементально додає нові записи до навчальної вибірки, не перевантажуючи все з нуля. Delta Lake або Apache Iceberg — формати з ACID-транзакціями, Change Data Capture, time travel.

Як уникнути training-serving skew за допомогою Feature Store

Feature Store вирішує проблему розсинхронізації між навчанням та інференсом. Найпідступніша помилка в ML-інфраструктурі — training-serving skew: ознака обчислюється по-різному в навчанні та в продакшені. Модель вчиться на «правильних» даних, а інференс отримує інші.

Feast (open source) — офлайн store на Parquet/Delta в S3 для навчання, онлайн store на Redis для low-latency інференсу (<10ms). Feature definitions як Python-код:

from feast import FeatureView, Field
from feast.types import Float32, Int64

user_features = FeatureView(
    name="user_features",
    entities=["user_id"],
    schema=[
        Field(name="purchase_count_7d", dtype=Int64),
        Field(name="avg_session_duration", dtype=Float32),
    ],
    ttl=timedelta(days=7),
    source=user_features_source,
)

Один definition використовується всюди — немає розбіжностей.

Потокові ознаки. Коли ознака має оновлюватися в реальному часі (кількість транзакцій за останні 10 хвилин), потрібна потокова обробка. Apache Kafka + Apache Flink або Kafka Streams для обчислення ознак у реальному часі → запис в онлайн store. Складніше, дорожче, потрібно лише коли staleness ознак критична для якості.

Розмітка даних: як не витратити бюджет даремно

Розмітка — найтрудомісткіша та недооцінювана частина ML-проекту. Погано розмічені дані не виправить жодна архітектура.

Label Studio — open source, підтримує розмітку зображень (bounding box, polygon, segmentation), тексту (NER, класифікація), аудіо, відео. Піднімається за 10 хвилин через Docker. Для невеликих команд — перший вибір.

Оцінка якості розмітки. Inter-annotator agreement — наскільки згодні розмітники між собою. Cohen's Kappa > 0.8 — добре, 0.6–0.8 — прийнятно, < 0.6 — завдання неоднозначне або інструкція погана. Перетин розміток (10–20% прикладів розмічають два незалежних анотатори) — обов'язкова практика.

Active learning. Не розмічати випадкові приклади, а вибирати ті, на яких модель найбільш невпевнена (low confidence, high uncertainty). Дозволяє досягти тієї ж якості при 50–70% обсягу розмітки. Modals, Prodigy, Label Studio підтримують active learning workflows. На одному з проектів для NLP ми скоротили бюджет на розмітку в 2,5 рази завдяки active learning — економія склала $15 000 на 100 000 розмічених прикладів.

Синтетичні дані. Коли реальних даних мало або отримати їх дорого. Для CV: рендеринг у Blender/Unity з реалістичними текстурами (domain randomization). Для NLP: parafrase через LLM, backtranslation. Ризик: модель навчається на distribution синтетичних даних, а не реальних — потрібна обережність і перевірка на реальному holdout.

Якість даних: валідація та моніторинг

Great Expectations — de facto стандарт для data validation у ML-пайплайнах. Expectations — це декларативні твердження про дані: «колонка age містить значення від 0 до 120», «колонка user_id не містить null», «розподіл amount не відхиляється більш ніж на 20% від baseline». Запускається в пайплайні, при провалі — блокує проходження.

Pandera — Pythonic alternative для pandas/polars DataFrames. Schema-based validation з type hints:

import pandera as pa

schema = pa.DataFrameSchema({
    "user_id": pa.Column(int, nullable=False),
    "score": pa.Column(float, pa.Check.between(0, 1)),
    "label": pa.Column(str, pa.Check.isin(["positive", "negative", "neutral"])),
})

Data freshness. Модель очікує дані за останні N днів. ETL впав, дані не оновилися — модель використовує застарілі ознаки. Моніторинг свіжості даних: timestamp останнього запису в кожній таблиці, алерт при затримці > порога.

Дедуплікація. Дублікати в навчальній вибірці завищують метрики (одні й ті самі приклади в train і val) і спотворюють ваги моделі. MinHash LSH для наближеної дедуплікації великих датасетів. Для точної — хеш за нормалізованим контентом.

Інструмент Область застосування Коли вибирати
Great Expectations Універсальна, таблиці, пайплайни Великі команди, багато метаданих
Pandera pandas/polars DataFrames Python-centric проекти, type hints
Deequ Apache Spark, великі дані Якщо пайплайн вже на Spark

Сховища та формати

Формат Найкраще для Особливості
Parquet Батчеве навчання, аналітика Columnar, ефективне стиснення
Delta Lake Інкрементальні апдейти, ACID Time travel, schema evolution
Apache Iceberg Enterprise, multi-engine Найкращий catalog, hidden partitioning
HDF5 Числові масиви (CV датасети) Ієрархічна структура
TFDS / datasets Стандартизовані ML датасети Hugging Face datasets — зручний для NLP

Для більшості ML-проектів на старті: Parquet в S3 + DVC для версіонування. Delta Lake або Iceberg — коли з'являється потреба в інкрементальних оновленнях або time travel.

Типові помилки при побудові пайплайнів

  • Пропуск перевірки свіжості даних. Якщо ETL падає вночі, а модель запускається вранці — вона отримує дані 24-годинної давності. Рішення: алерт при затримці > 30 хвилин.
  • Відсутність версіонування даних. Не можна відтворити експеримент, бо дані змінилися. DVC або Delta Lake time travel виправляють це.
  • Забувають про schema evolution. Нове поле з’являється, а пайплайн падає. Автоматичне виявлення змін схеми через Great Expectations.

Active learning дозволяє скоротити бюджет на розмітку до 50–70%. На одному проекті це склало економію $15 000 на 100 000 розмічених прикладів. Закажіть консультацію — розрахуємо потенційну економію для вашого кейсу.

Що входить у проект з дата-інжинірингу для ML

Ми надаємо повний цикл:

  • Аудит існуючих даних та пайплайнів (1 тиждень).
  • Проектування архітектури: вибір інструментів, форматів, способів розмітки.
  • Реалізація ETL/ELT пайплайну з валідацією та моніторингом.
  • Документація коду та процесів (model card, data card).
  • Навчання вашої команди роботі з пайплайном.
  • SLA на супровід та підтримку.

Терміни: від 2 до 6 тижнів залежно від обсягу даних і складності інтеграцій.

Як ми будуємо пайплайн: покроково

  1. Аудит існуючих даних. Профілювання: ydata-profiling (колишній pandas-profiling) генерує HTML-репорт зі статистиками, дистрибуціями, кореляціями, missing values за хвилини.
  2. Проектування пайплайну. Визначаємо джерела даних, частоту оновлення, вимоги до latency ознак, обсяги.
  3. Реалізація та тестування. Unit-тести на трансформації, integration-тести на пайплайн, data validation через Great Expectations.
  4. Деплой та моніторинг. Алерти на freshness, quality checks, аномалії в обсягах даних.

Чому варто довірити це нам

Ми займаємося дата-інжинірингом та ML з понад 8-річним досвідом. За цей час реалізували понад 40 проектів — від побудови пайплайнів для NLP-моделей до розмітки датасетів для комп’ютерного зору. Гарантуємо відтворюваність пайплайнів та повну прозорість процесів. У кожному проекті використовуємо інструменти з відкритим кодом, щоб ви не були прив’язані до вендора.

Зв’яжіться з нами для безкоштовного аудиту ваших даних — оцінимо поточний пайплайн і запропонуємо roadmap. Замовте побудову ML-пайплайну під ключ.