TL;DR: мы сделали llm_markup — Spark‑приложение для массовой обработки данных через LLM API в существующем Airflow‑контуре. Источник данных, промпт, формат ответа, лимиты и целевая таблица задаются конфигурацией. Инструмент берет на себя батчинг, распределение запросов, контроль нагрузки на API, валидацию ответов и сопоставление результата с исходными строками.

Меня зовут Дима Иванов, я дата‑инженер в Lamoda Tech. Я работал над llm_markup и в этой статье расскажу, как мы пришли к его архитектуре и какие задачи пришлось решить, чтобы инструмент стабильно работал в продакшене. Разберем это на двух примерах: оценке качества поиска и классификации негативных отзывов.

Введение

LLM все чаще используют внутри дата‑пайплайнов: для классификации текстов, оценки соответствия объектов и извлечения структурированных данных. На небольшой выборке достаточно Python‑скрипта, но при обработке десятков и сотен тысяч строк нужны батчинг, параллелизм, контроль нагрузки, повторы и сопоставление ответов с исходными данными.

На практике мы впервые столкнулись со всем этим, когда стали масштабировать LLM‑разметку для оценки качества поиска. В fashion‑каталоге запросы варьируются от «летнее платье» до «adidas rhjccjdrb» — это «adidas кроссовки» в неправильной раскладке. Модели нужно восстановить смысл запроса и определить, насколько ему соответствует каждый товар. Позже та же инженерная механика понадобилась для другой бизнес‑задачи — классификации причин недовольства в отзывах. Два сценария показали, что общую часть пайплайна можно отделить от бизнес‑логики и переиспользовать.

Так появилось llm_markup — универсальное Spark‑приложение, в коде которого нет понятий поискового запроса, товара или отзыва. Если новый сценарий можно описать источником данных, промптом, форматом ответа и целевой таблицей, для него не требуется менять общий Spark‑код.

После перехода на llm_markup на сопоставимом объеме данных медианное время разметки в поисковом сценарии сократилось с ~10 часов до ~5 часов — почти вдвое. Дальше разберем, как распределение запросов между Spark‑партициями и контроль нагрузки на LLM API помогли получить это ускорение.

Архитектура: почему именно Spark

Упрощенная схема пайплайна выглядит так:

Упрощенная схема пайплайна llm_markup
Упрощенная схема пайплайна llm_markup

Ожидаемый вопрос: зачем Spark для LLM‑вызовов? Нельзя просто написать Python‑скрипт с asyncio и пулом соединений? Или поднять Ray/Celery отдельно?

Но в нашем случае Spark оказался наиболее подходящим способом встроить массовую LLM‑разметку в уже существующий пайплайн. Причин несколько:

  • Объемы. Десятки тысяч строк за один прогон требуют управляемого параллелизма, повторов и восстановления после сбоев.

  • Источники данных. Инструмент умеет читать данные из Hive, Feature Storage на HDFS и внешних баз по JDBC. Spark позволяет работать с ними в существующем контуре безопасности и вычислений.

  • Оркестрация. Airflow уже запускает Spark‑приложения в подготовленном окружении. Отдельный runtime для LLM‑батчей — еще один сервис, деплой и мониторинг.

  • Распределенное исполнение. Spark умеет параллельно выполнять задачи на executor‑ах и собирать результат в DataFrame.

Главная идея фреймворка: Spark можно использовать как распределенный движок для сетевого I/O, а не только для классических ETL‑вычислений. Партиции данных независимо ходят в LLM API, а результаты возвращаются в DataFrame. Spark берет на себя планирование и распределение работы. Общая конфигурация передается executor‑ам как broadcast‑переменная — доступная только для чтения копия, которую Spark рассылает на исполняющие узлы. Итоговая статистика собирается на driver‑е.

Но Spark не делает внешние HTTP‑вызовы exactly‑once. Если executor или Spark‑задача падает после отправки запроса, задача может быть переисполнена и LLM‑вызов повторится. Это нужно учитывать в стоимости, лимитах и проектировании побочных эффектов. Динамическое добавление executor‑ов тоже зависит от настройки конкретного Spark‑кластера, а не появляется автоматически из кода приложения.

Инфраструктура вокруг инструмента уже существовала: Airflow оркестрирует DAG‑и, данные лежат в хранилищах, BI‑дашборды показывают метрики. llm_markup добавляет между источником и потребителем звено с LLM‑обработкой. На выходе джоба может вернуть метки из заданного словаря или типизированные данные по декларативной схеме. Оба варианта используют общий механизм батчинга, распределения запросов и обработки ошибок.

mapInPandas: Spark как HTTP‑клиент

Механизм, который связывает Spark и LLM в обоих вариантах ответа, — mapInPandas. Это API, которое применяет Python‑функцию к pandas DataFrame каждой Spark‑партиции и возвращает новый DataFrame. Внутри функции можно сформировать промпт или payload, вызвать HTTP API и разобрать ответ.

Почему не обычный Python UDF? Скалярный UDF обрабатывает по одной строке и плохо подходит для задачи, где внутри партиции нужен собственный LLM‑клиент, ограничение частоты запросов и объединение нескольких строк в один запрос. pandas_udf тоже умеет работать с векторизованными порциями данных, но mapInPandas дает удобный iterator по pandas DataFrame и позволяет гибко формировать выходные строки.

На вход job принимает либо одну текстовую колонку, либо Jinja2-шаблон, который собирается из нескольких колонок DataFrame. Если тексты короткие, несколько строк можно отправить в одном LLM‑запросе. Это уменьшает долю токенов, которые тратятся на повторяющиеся инструкции. Батчинг не обязателен. Для длинных уникальных промптов можно использовать batch_size=1 и масштабировать обработку за счет Spark‑партиций.

Размер батча нельзя увеличивать бесконечно. Чем больше строк попадает в один запрос, тем реже повторяется общая инструкция, но тем длиннее становится ответ и тем выше риск, что модель пропустит или перепутает отдельные элементы. Поэтому batch_size подбирается для каждого сценария отдельно, а качество проверяется на контрольной выборке.

Общий каркас обработки одинаков для двух вариантов ответа:

Параллельная обработка Spark-партиций через mapInPandas
Параллельная обработка Spark‑партиций через mapInPandas
  1. Исходный DataFrame нарезается на батчи и репартиционируется.

  2. Каждая партиция батчей уходит в mapInPandas на executor‑е.

  3. Внутри партиции создается локальный LLM‑клиент, работают ограничение частоты запросов и повторы при ошибках.

  4. Общая конфигурация передается executor‑ам через broadcast‑переменные, а driver собирает статистику и диагностику.

  5. Spark формирует результирующий DataFrame и записывает его в Hive. Для типизированных данных можно дополнительно сохранить отдельную таблицу диагностики.

Дальше варианты расходятся: что именно отправляется в API и как сопоставляется ответ с исходными строками.

JSONL: классификация и метки

Режим по умолчанию (response_mode=jsonl) подходит для задач, где на вход — текст или Jinja2-шаблон, а на выход — одна или несколько меток из заранее известного словаря. Модель возвращает по одной JSON‑строке для каждого элемента батча, поэтому такой формат называется JSON Lines, или JSONL.

Перед разбиением на батчи каждая строка получает уникальный технический идентификатор global_index. Для этого Spark сначала создает монотонно возрастающий ID, а затем преобразует его в плотную нумерацию от нуля. Это не бизнес‑ключ и не способ сохранить исходный порядок строк. Индекс нужен только для сопоставления данных внутри текущего запуска.

Поток данных:

  1. Строки группируются в батчи: batchid = floor(_global_index / batch_size).

  2. DataFrame репартиционируется по _batch_id, строки батча агрегируются через groupBy + collect_list.

  3. Партиция уходит в mapInPandas. Внутри для каждого батча формируется промпт, вызывается get_answer(), ответ парсится и валидируется по списку допустимых меток.

  4. Spark делает LEFT JOIN результатов с исходным DataFrame по _global_index. Строки без LLM‑результата получают запасное значение согласно конфигурации.

Внутри одного LLM‑запроса используется короткий локальный индекс i. Модель возвращает его вместе с метками, после чего джоба восстанавливает _global_index соответствующей строки. В конце Spark делает LEFT JOIN с исходным DataFrame по этому идентификатору. Поэтому сопоставление не зависит от порядка, в котором модель вернула элементы батча.

Поле i — индекс строки внутри батча (0, 1, 2…), labels — массив текстовых меток. Инструкции по формату job добавляет в промпт автоматически. Пример ответа:

 {"i": 0, "labels": ["Проблемы с размером"]}
 {"i": 1, "labels": ["Проблемы с качеством", "Быстрое изнашивание"]}
 {"i": 2, "labels": ["Не определено"]}

В JSONL‑режиме label_mapping связывает текстовые метки с выходными значениями. Отдельно настраиваются повтор строк с ошибочным ответом и обработка невалидных меток. Если LLM вернула несколько меток для одной строки, job может развернуть их через explode.

Structured: типизированное извлечение

Второй режим (response_mode=structured) — для задач, где JSONL‑меток уже недостаточно: нужны поля с типами, вложенные объекты, связь «один ко многим». Например, из отзыва «подошел маломерит, через неделю порвался шов» нужно получить несколько записей:

[
   {"aspect": "размер", "confidence": 0.9},
   {"aspect": "качество шва", "confidence": 0.85}
 ]

Здесь тоже батчи и mapInPandas, но другой контракт и другой вызов API:

  1. structured_contract описывает, какие колонки отправлять модели, какую схему ответа ждать и как собрать выходную Spark‑таблицу. Контракт валидируется на driver‑е до первого LLM‑запроса.

  2. Строкам батча назначаются локальные ID 1, 2, 3 — в модель не уходят чувствительные бизнес‑идентификаторы. После ответа ID однозначно сопоставляются с исходными строками.

  3. Контракт компилируется в Pydantic‑модель и передается в API через response_format. Внутри mapInPandas вызывается get_structured_answer(), ответ дополнительно проверяется через parse().

  4. При необходимости job пишет таблицу диагностики: для проблемной строки сохраняются тип проблемы и детали (batch_failed, validation_substring_failed, unknown_local_id и др.).

Режимы не смешиваются: параметры JSONL и structured задаются раздельно. При этом оба варианта используют общие механизмы ограничения нагрузки и обработки ошибок.

Фреймворк работает с API, совместимым с OpenAI. Подключение настраивается через base_url, ключ и заголовки.

Конфигурация вместо кода

Ниже упрощенный иллюстративный пример — не продакшен‑конфиг, а схема параметров для классификации отзывов из корпоративного DWH по JDBC. В примере также включен enable_smart_retry, который повторяет отдельные строки после некорректного ответа. Для остановки при систематических сбоях выбран error_mode=partition: лимит n_errors проверяется отдельно в каждой Spark‑партиции. Альтернативный режим и различия между ними подробно разберем ниже.

comments_llm {
  type = "llm_markup"
  sql = "llm/comments/select_for_labeling.sql"
  connection_id = "dwh"
  prompt_file = "comments/classification_prompt.txt"
  llm_conn_id = "llm_api"
  output_db = "analytics"
  output_table = "comments_llm_labeled"

  # Входные данные и формат результата
 processing_config {
    text_column = "comment_text"
    result_column = "llm_label"
    batch_size = 25
    replace_invalid_labels_with_default = false
  }

  # Ограничение нагрузки, повторы и остановка при ошибках
  rate_limit_config {
    requests_per_second = 5
    safety_factor = 0.8
    max_concurrent_partitions = 10
    max_retries = 3
    enable_smart_retry = true
    error_mode = "partition"
    n_errors = 10
  }

  # Преобразование LLM-меток в значения целевой таблицы
  label_mapping {
    "Проблемы с размером" = 1
    "Проблемы с качеством" = 2
    "Не определено" = 0
    "default" = 0
  }
}

В шаблонном режиме (Jinja2 в prompt_file) параметр text_column не нужен — промпт собирается из колонок DataFrame на executor‑е. Маппинг меток работает без учета регистра: "релевантно" и "РЕЛЕВАНТНО" эквивалентны.

Конфигурация не избавляет от инженерных требований: нужны права на Airflow, доступ к источнику и Hive‑таблице, соединение с LLM, понимание схемы данных и контроль качества разметки. Но новый сценарий не требует писать новый Spark‑код.

Jinja2 в промпте и Airflow: как развести два уровня шаблонизации

Для отзыва достаточно передать модели одну текстовую колонку. С поиском сложнее. Чтобы оценить пару «запрос → товар», модели нужны название, бренд, категория, цвет, сезон и описание. Поэтому user‑промпт мы сделали Jinja2-шаблоном:

Запрос: {{ query }}
Товар: {{ name }}, {{ brand }}
Описание: {{ description }}

На executor‑е такой шаблон отдельно рендерится для каждой строки DataFrame. NULL превращается в пустую строку, а значения полей ограничиваются по длине, чтобы промпт не разрастался бесконтрольно. Порог нужно проверять для каждого нового сценария, иначе вместе с лишним текстом можно обрезать и полезный контекст. Общие правила оценки при этом остаются в системном промпте.

Здесь возникает конфликт двух шаблонизаторов: Airflow тоже использует Jinja2 и пытается обработать {{ query }} еще до запуска Spark, хотя такой переменной в контексте DAG‑а нет. Поэтому handler временно заменяет фигурные скобки нейтральными плейсхолдерами, а Spark‑приложение восстанавливает их уже перед рендерингом строк. В итоге каждый из двух Jinja2 получает только свои переменные.

Ограничение частоты запросов: не перегрузить API

У LLM API есть ограничение на количество запросов в секунду, или RPS. Когда десятки Spark‑задач одновременно отправляют запросы без координации, лимит будет превышен: начнутся 429, таймауты и повторные попытки.

Подход: задается целевой RPS и коэффициент запаса, например 0.8. Из них вычисляются число рабочих Spark‑партиций и задержка между запросами внутри каждой:

safe_rps = target_rps * safety_factor
num_partitions = min(max(1, int(safe_rps)), max_concurrent_partitions, total_batches)
delay_per_partition = num_partitions / safe_rps

Каждая обрабатывающая партиция перед стартом получает случайную задержку от 0 до delay_per_partition. Это разброс старта (jitter): без него первая секунда дает всплеск запросов, с ним старты распределяются во времени.

Между запросами задача выдерживает паузу, проверяя, сколько прошло с прошлого вызова. Если LLM отвечает дольше рассчитанной задержки, дополнительное ожидание не требуется. Ограничение уже соблюдается естественным образом.

Пример: целевой RPS = 5, коэффициент запаса = 0.8, значит безопасный RPS = 4. Четыре рабочие партиции делают примерно по одному запросу в секунду, а суммарная целевая нагрузка составляет около 4 RPS.

Ограничение частоты можно отключить (enable_rate_limiting = false), если пропускная способность внутреннего сервиса подтверждена нагрузочными тестами. Тогда фактический RPS определяется временем ответа модели и числом Spark‑задач, которые выполняются одновременно. max_concurrent_partitions ограничивает количество рабочих партиций, но не равен числу одновременных запросов. Фактическую конкурентность дополнительно ограничивают доступные Spark task slots на executor‑ах.

Когда LLM ошибается

LLM API — не база данных. Бывают сетевые таймауты, 5xx, пустой content, ответ не в ожидаемом формате и отказ модели отвечать. На объемах в тысячи батчей это не «если», а «когда». Ниже — механизмы, общие для обоих режимов, и детали, специфичные для JSONL.

Повторы и два режима контроля ошибок

Каждый запрос повторяется до max_retries раз с экспоненциальной задержкой 2^attempt секунд между попытками (1, 2, 4…). Но повтор — это только половина решения. Нужно еще понимать, когда остановить задание при систематической проблеме.

Для контроля используется параметр n_errors — допустимое число неуспешных батчей. Проверять этот лимит можно в двух режимах.

Режим partition — по умолчанию и основной для стабильных продакшен‑пайплайнов. Лимит применяется локально к каждой Spark‑партиции. После превышения лимита обработка этой партиции завершается ошибкой. Spark может переисполнить задачу, поэтому HTTP‑вызовы потенциально повторятся. Это компромисс между скоростью восстановления и стоимостью внешнего API.

Режим wave — глобальный лимит на все задание. Батчи разбиваются на последовательные волны заданного размера. Driver запускает волну, дожидается материализации результата, считает неуспешные батчи. При превышении n_errors текущая диагностика сохраняется, последующие волны помечаются как пропущенные, после чего задание завершается ошибкой. Это быстрее останавливает систематические сбои, но добавляет материализацию между волнами. Имеет смысл при первых прогонах нового промпта или модели.

Для стабильных пайплайнов обычно подходит partition. Для новых или экспериментальных — wave. В structured‑режиме те же режимы работают на уровне батчей. Ошибки валидации и сопоставления дополнительно попадают в таблицу диагностики.

Умный повтор (JSONL)

Представьте батч из 25 строк. LLM вернула ответ, но для трех строк задание подставило запасное значение. Причины разные: модель могла честно не найти категорию или ответ мог сломаться при генерации и разборе.

Умный повтор (enable_smart_retry, по умолчанию включен) отправляет отдельно только строки, для которых запасное значение появилось из‑за ошибки разбора или отсутствующего ответа. Если модель теперь дает валидный ответ, результат обновляется. Если запасное значение повторяется, задание сохраняет его.

Критический нюанс: система различает «модель вернула значение по умолчанию как предсказание» и «значение по умолчанию подставлено как запасной вариант из‑за ошибки». Повтор делается только для второго случая. Если модель честно вернула допустимую метку Не определено или пустой labels=[] при replace_invalid_labels_with_default=false, переспрашивать ее не нужно.

Первая версия фреймворка перезапрашивала все строки со значением по умолчанию без разбора — это давало лишние запросы. Разделение на «предсказанное значение по умолчанию» и «запасное значение после ошибки» эту проблему сняло.

Пустой список меток и null вместо default (JSONL)

В некоторых сценариях допустим результат «подходящей категории нет»:

{"i": 0, "labels": []}

Это не ошибка разбора. Параметр replace_invalid_labels_with_default=false позволяет сохранить такой пустой результат как null в выходной колонке вместо автоматической подстановки default_label. Та же настройка помогает отличать корректное отсутствие метки от запасного значения, вызванного поврежденным или невалидным ответом.

Как тестируем llm_markup

Тесты используют локальный HTTP‑сервер на pytest, который имитирует OpenAI‑compatible API рядом с локальным Spark. Сервер возвращает заготовленные ответы, JSONL и ошибки. Этого хватает, чтобы покрыть батчинг, rate limiting, режимы контроля ошибок и structured‑контракты без обращения к реальной модели.

Как это работает на практике: оценка качества поиска

Соберем все вместе на оценке качества поиска, с которой началась история llm_markup.

Подробнее о методике оценки поиска с помощью метрик, LLM и обратной связи пользователей мы рассказывали в отдельной статье.

Перед разметкой другие задачи пайплайна готовят исходный набор данных и складывают его в Feature Storage.

Сама LLM‑разметка появилась еще до общего инструмента. Для нее было написано отдельное Spark‑приложение, которое читало подготовленный фичасет, собирало данные на driver и последовательно отправляло в LLM пары «запрос → товар». Решение работало, но с ростом выборки росла и одна длинная очередь запросов.

В llm_markup источник данных и логика оценки остались прежними, а выполнение изменилось. Теперь данные не собираются на driver, а распределяются между Spark‑партициями, каждая из которых независимо обращается к LLM API.

Получается набор признаков с парами «запрос → товар»: поисковый запрос, SKU, ранг товара, название, бренд, категория, атрибуты, гендер, цвет, сезон.

Затем запускается LLM‑разметка в режиме JSONL. Spark читает этот набор и для каждой пары «запрос → товар» рендерит промпт из Jinja2-шаблона:

Запрос: летнее платье
Товар:  Название: Платье макси  
Категория: Платья макси  
Бренд: пример бренда  
Гендер: women  
Цвет: ['синий']  
Сезон: ['мульти']  
Описание: Летнее платье из легкой ткани...

Отдельный системный промпт описывает правила оценки релевантности: учитывает гендер, сезонность, синонимы, опечатки, раскладку клавиатуры и брендовые связи. User‑шаблон содержит только контекст «запрос + товар». Перед продакшен‑запуском такие правила проверяются на контрольной человеческой выборке.

LLM возвращает метки релевантно / частично релевантно / нерелевантно, затем они маппятся в 2 / 1 / 0. Отдельное значение default: -1 означает «разметить не удалось» — его не смешивают с «нерелевантно».

В поисковом сценарии на каждую пару «запрос → товар» приходится большой уникальный контекст, поэтому несколько строк не объединяются в один запрос и используется batch_size=1. Количество обращений к модели не уменьшилось, но теперь они выполняются параллельно в Spark‑партициях, а не ждут друг друга в одной очереди.

После перехода на llm_markup на сопоставимом объеме данных медианное время работы LLM‑разметки снизилось с 9 часов 59 минут до 5 часов 16 минут — почти вдвое.

Сама LLM при этом отвечать быстрее не стала. Выигрыш появился за счет большей конкурентности, а значит, и большей одновременной нагрузки на API. Ее ограничивают число рабочих партиций и доступные task slots на executor‑ах, поэтому уровень параллелизма подбирается под возможности LLM‑сервиса.

Результат пишется в целевую Hive‑таблицу с оценками релевантности для товаров в выдаче. Дальше SQL считает DCG, NDCG и precision. Запросы с метрикой ниже порога помечаются как проблемные, а BI‑дашборд помогает сравнивать варианты ранжирования и находить регрессии.

Классификация негативных отзывов: тот же фреймворк, другая задача

Второй продакшен‑сценарий — классификация негативных отзывов, тоже в режиме JSONL. Источник — Oracle DWH через JDBC, откуда выбираются отзывы с рейтингом 1–3.

Конфигурация другая: JDBC‑источник вместо Feature Storage, простая текстовая колонка вместо Jinja2-шаблона и batch_size=25 для коротких текстов. Ограничение частоты запросов задается через общий механизм, для контроля ошибок выбран режим partition. Результат пишется в Hive. Следующая SQL‑таска маппит LLM‑метки в категории проблем: качество, запах, размер, брак, несоответствие описанию и другие.

В одном из запусков 234 отзыва удалось собрать в 10 обращений к LLM вместо 234 отдельных. Длинная инструкция с описанием категорий повторилась 10 раз, а не для каждой строки.

Батч при этом не превращается в мешок ответов. Каждый отзыв получает локальный индекс, который модель возвращает вместе с метками. После валидации llm_markup связывает результат с исходной строкой по глобальному индексу, не полагаясь на порядок ответа. В этом запуске все 234 отзыва были сопоставлены без ошибок и запасных значений.

Сопоставление ответов LLM со строками внутри батча
Сопоставление ответов LLM со строками внутри батча

Тот же Docker‑образ и общий код. Разница — SQL, HOCON‑конфиг, промпт и целевая таблица. Команде не нужно писать или выкатывать новое Spark‑приложение для каждого сценария.

Ограничения, которые нельзя скрывать

  • Повторы Spark‑задач могут повторить внешний LLM‑вызов. Это не exactly‑once обработка.

  • LLM‑разметку нужно регулярно сверять с человеческой контрольной выборкой: технически валидный ответ не гарантирует правильную оценку.

  • Стоимость и задержка зависят от размера промпта, количества токенов, времени ответа API и количества повторов.

  • wave уменьшает масштаб потерь при системной ошибке, но не заменяет мониторинг API и качество промпта.

  • Структурированный ответ делает результат формально валидируемым, но не устраняет содержательные ошибки модели.

Подводя итоги

Spark как распределенный LLM‑клиент — идея, которая работает, если учитывать специфику внешнего API. mapInPandas, партиционирование и передача общей конфигурации дают основу для параллельного сетевого I/O. Сверху нужны ограничение частоты запросов, повторы, контроль ошибок, рендеринг промптов и диагностика.

Два режима ответа — JSONL для меток и structured для типизированного извлечения. Оба через один и тот же Spark‑каркас, с разными контрактами на входе и выходе.

Конфигурация вместо нового Spark‑кода — главная ценность инструмента. Новый сценарий часто сводится к SQL, HOCON‑конфигу, промпту и заранее созданной целевой таблице. Само Spark‑приложение при этом менять не требуется.

Умный повтор и честный запасной результат — важная эксплуатационная деталь. Разделение между предсказанной меткой по умолчанию и подстановкой после ошибки уменьшает лишние запросы. Возможность сохранить пустой результат без default полезна для задач, где отсутствие категории — нормальный сигнал.

Режим волн появился из практической боли: при систематических проблемах нельзя часами продолжать тратить ресурсы. Волны позволяют ограничить масштаб обработки и сохранить диагностику до аварийного завершения.

В результате у нас получился универсальный инструмент для массовой LLM‑разметки в существующей Spark‑инфраструктуре, с контролем нагрузки, обработкой ошибок и проверкой ответов. Если решаете похожую задачу, имеет смысл присмотреться к Spark и mapInPandas: каркас распределенного исполнения уже есть, а специфику сценария можно вынести во входные данные, промпт и контракт результата.

Комментарии (1)


  1. Innesiya
    28.08.2026 09:07

    Расчет задержки из target_rps и safety_factor — это открытый контур: он держит лимит только пока ваш джоб единственный потребитель модели. У нас два параллельных пайплайна на одну точку входа дружно собирали 429, хотя каждый по своим настройкам был в рамках, и внятно заработало только общим токен-бакетом во внешнем хранилище. Отдельно смущает, что в режиме partition переисполнение таски гонит всю партицию в API заново, то есть сам факт упирания в лимит увеличивает нагрузку. target_rps у вас проставляется руками в каждом конфиге и за суммой следит человек, или есть общая квота на llm_conn_id, которую одновременные джобы делят между собой?