Универсальный Spark-пайплайн: как Lamoda масштабировала LLM-разметку для продакшена
Компания Lamoda Tech создала инструмент llm_markup — универсальное Spark-приложение, позволяющее интегрировать вызовы больших языковых моделей в существующие конвейеры обработки данных на базе Apache Spark и Apache Airflow. Решение решает проблемы батчинга, контроля нагрузки на API и обработки ошибок, сократив время разметки данных почти вдвое.
# Как превратить LLM-разметку в универсальный Spark-пайплайн
Масштабирование задач, требующих взаимодействия с большими языковыми моделями (LLM), часто сталкивается с инженерными трудностями: от управления нагрузкой на внешний API до восстановления данных после сбоев. В статье дата-инженер Lamoda Tech Дима Иванов рассказывает о разработке llm_markup — инструмента, который абстрагирует сложность работы с LLM под капотом надежного распределенного движка Apache Spark.
Ранее масштабирование LLM-разметки для оценки качества поиска в каталоге fashion требовало значительных временных ресурсов. При переходе на новый инструмент медианное время обработки сократилось с 10 часов до 5 часов. Теперь одна и та же инженерная механика используется как для восстановления смысла искаженных поисковых запросов, так и для классификации причин недовольства в отзывах клиентов.
Архитектурное решение: Spark как HTTP-клиент
Ключевым выбором стала интеграция непосредственно с Apache Spark. Вместо создания отдельного сервиса на Python или использования Ray/Celery, инженеры Lamoda решили превратить Spark в распределенный клиент для сетевого I/O. Это позволило использовать существующую инфраструктуру, включая Airflow для оркестрации и HDFS для хранения данных.
Сердце инструмента — API mapInPandas. В отличие от стандартных векторизованных функций, mapInPandas позволяет применять пользовательский Python-код к отдельным партициям данных (подробным блокам строк) внутри каждого исполняющего узла (executor). Это дает необходимую гибкость для:
* Формирования промптов, учитывающих конкретный блок данных. * Ограничения частоты запросов для каждой партиции локально. * Управления параллелизмом и обработки ошибок без блокировки всего кластера.
Основная идея заключается в том, что каждая партиция данных независимо взаимодействует с LLM, а финальный результат собирается в единый DataFrame. Конфигурация промпта, форматов ответов и лимитов передается всем узлам через broadcast-переменные, что делает код приложения универсальным и не зависящим от конкретной бизнес-задачи.
Режимы работы: от классификации к структурированным данным
Фреймворк поддерживает два основных режима ответа от модели, каждый из которых имеет свою логику сопоставления результатов с исходными строками.
JSONL-режим: классификация и метки
По умолчанию используется формат JSON Lines (JSONL). Он оптимален для задач классификации, где нужно получить набор меток для каждой строки текста. Модель возвращает JSON-объект, содержащий локальный индекс строки внутри батча и массив меток.
* Уникальная идентификация: Перед разбиением на батчи каждая строка получает глобальный идентификатор. Внутри батча строки получают локальные индексы (0, 1, 2...), которые возвращаются вместе с ответом модели. * Сопоставление: После получения ответов Spark выполняет LEFT JOIN, связывая локальные индексы с глобальными ID исходных данных. Это гарантирует корректное распределение меток даже при изменении порядка элементов.
Structured-режим: типизированное извлечение
Для более сложных задач, требующих вложенных структур или нескольких записей на одну строку (например, извлечение аспектов из отзыва), используется режим structured.
В этом случае перед запуском компилируется контракт ответа (используя библиотеку Pydantic), который определяет схему ожидаемого JSON. Модель должна возвращать строго типизированные объекты. Инструмент автоматически валидирует ответы, проверяя соответствие схеме. Если валидация проходит, данные преобразуются в Spark-таблицу с заданными типами данных.
Управление нагрузкой и обработка ошибок
При массовых запросах к LLM-апи критически важным становится контроль скорости (Rate Limiting). Непредсказуемое время ответа и сетевые задержки могут приводить к превышению лимитов запросов в секунду (RPS).
Решение Lamoda использует динамическое распределение нагрузки:
1. Расчет безопасного RPS: Целевая скорость умножается на коэффициент запаса (например, 0.8). 2. Jitter (рандомизация): Количество рабочих партиций рассчитывается исходя из безопасного RPS. Каждая партиция получает случайную задержку перед стартом, что предотвращает резкий всплеск запросов в первый момент времени. 3. Локальное ограничение: Внутри каждой партиции реализован таймер, следящий за интервалом между вызовами API.
Обработка ошибок реализована через два режима:
* Partition (по партиции): Если в пределах одной Spark-партиции накапливается определенное количество ошибок (параметр n_errors), задача завершается. Остальные партиции продолжают работу. Это позволяет кластеру восстановиться и переисполнить неудачную задачу, если ошибка была временной. * Wave (по волнам): Глобальная проверка для всего задания. Батчи обрабатываются последовательными группами. Если в одной из волн слишком много сбоев, все последующие батчи пропускаются, чтобы избежать бесконечного потребления ресурсов.
Особый акцент сделан на умном повторе (Smart Retry). Инструмент различает ситуации, когда модель вернула «Не определено» (корректный ответ), и случаи, когда ответ не был распарсирован или таймаут произошел. Повторный запрос отправляется только во втором случае, если ранее была подставлена резервная метка.
Заключение
Инструмент llm_markup демонстрирует, как можно эффективно объединить возможности больших языковых моделей с надежностью и масштабируемостью традиционных Big Data-платформ. Ключевым преимуществом является декларативная конфигурация: новые сценарии обрабатываются путем задания источника данных, промпта и схемы ответа, без необходимости переписывать код пайплайна. Этот подход значительно ускоряет внедрение AI-решений в корпоративные процессы.