Последние годы я работал разработчиком в области машинного обучения, компьютерного зрения и 3D-реконструкции. Писал и видел множество вариантов организации вычислений и обработки потока данных. Но в них мне не хватало системы и организованности, прозрачности и гибкости. Поэтому я набрался решимости и воплотил самые мои смелые — и, как оказалось, вполне рабочие — идеи в новом фреймворке, который я назвал ICO.
Предлагаю смотреть на эту работу не просто как на «очередной фрейморк-велосипед», а как на своего рода инженерное исследование, в чем-то творческое, — попытку переосмыслить и систематизировать знакомые многим профессионалам задачи.
Фреймворк состоит из нескольких подсистем:
вычислительное ядро: база для организации потока вычислений и обработки данных
среда выполнения: мониторинг, события, обмен сообщениями, мульти-процессинг
прикладные реализации: алгоритмы обучения, аналог PyTorch DataLoader
профилирование / мониторинг производительности: в разработке
В данном посте я затрону вычислительное ядро и немного среду выполнения. Более подробное описание потребует дополнительной статьи.
ICO — Input, Context, Output
Атомарный элемент ICO — это оператор производящий вычисления или модификацию данных. Он имеет сигнатуру, в которой указывается тип входных и выходных данных: I → O
В расширенном варианте оператор принимает еще и контекст, становясь, как правило оператором обучения, — ведь контекст хранит и передает информацию, полученную из входных данных: I, C → O
Композиция операторов в поток вычисления происходит через оператор “|”
pipeline = load_data | augment | train
Таким образом, на уровне оператора задается строгая типизация, а при композиции операторов статический анализатор может проверить соответствие типов входа и выхода.
Анализаторы кода, такие как Pylance и mypy, предоставляют возможность провести валидацию типов в цепочке еще до выполнения кода.
Формируя цепочку выполнения мы «под капотом» создаем дерево операторов, что вместе с наличием сигнатуры у каждого оператора дает на еще одно важное свойство — интроспекцию и автоматическое описание плана вычислений.
Любой пайплайн может «рассказать» о себе с помощью утилиты describe() еще до начала своего выполнения. В результате будет показан план выполнения, отрендеренный своим рендерером с помощью модуля Rich Console.
Это две ключевые, но не единственные особенности ICO. Давайте посмотрим немного «учебный», но живой пример.
Пример: вычисление числа Фибоначчи.
from ico import IcoProcess, operator # Мы создаем состояние для хранения двух последних чисел из последовательности State = tuple[int, int] # Объявляем оператор, используя декоратор. # Шаг Фибоначчи - это оператор который модифицирует состояние @operator() def fib_step(state: State) -> State: a, b = state return (b, a + b) # Последний оператор в цепочке — получение результата из состояния @operator() def take_first(state: State) -> int: return state[0] # Создаем оператор-процесс, повторяющий заданный оператор восемь итераций fib8 = IcoProcess(fib_step, num_iterations=8) # Собираем «пайплайн» — последовательность из двух шагов pipeline = fib8 | take_first
Теперь посмотрим план выполнения
pipeline.describe()

Поток операторов читается сверху вниз, каждая строка — это оператор и его сигнатура с типами входных и выходных данных.
Оператор может иметь свой способ отображения, и для процесса используется группировка — внутри отображаются операторы тела процесса, а сам процесс обозначается ключевыми словами «iterate in» и «emit».
Таким образом, еще до запуска можно увидеть последовательность операторов и как данные меняются при прохождении через пайплайн.
Для выполнения вычислений и получения результата нам надо вызывать пайплайн, передав в него входные данные.
# Запускаем пайплайн, передавая входные данные - начальную последовательность result = pipeline((0, 1)) # И вот наш результат! print(f"{result=}") # 21
Потоки данных и ленивые вычисления
В реальных задачах мы часто имеем дело с потоками данных, таких как набор кадров из видео, бачи в эпохе обучения нейронной сети и т.д. В ICO для таких задач используется интерфейс Iterator[T]. Помимо однозначности описания интерфейса, это позволяет осуществлять загрузку данных по необходимости, используя итераторы.
Посмотрим опять игрушечный пример.
# Объявим оператор который работает с одним числом @operator() def scale_by_10(x: int) -> int: return x * 10 # Обернем оператор в поток, что поменяет его тип на Iterator[int] scale_stream = scale_by_10.stream() scale_stream.describe()

Видим, исходных оператор int → int теперь стал обернут в потоковый оператор Iterator[int] → Iterator[int]
Такой поток мы можем использовать с другими операторами, ожидающими Iterator на вход. Важно заметить, что выполнение оператора будет происходит после запроса следующего элемента у итератора, т.е. по принципу ленивой загрузки.
data = [1, 2, 3, 4, 5] for scaled in scale_stream(iter(data)): print(f"{scaled=}")
Элемент среды выполнения - прогресс
При долгих вычислениях важно понимать, что сейчас происходит и сколько работы уже выполнено. В ICO мониторинг прогресса относится не к вычислительному ядру, а к среде выполнения. Оператор сообщает о прогрессе через события, а runtime отвечает за их обработку и отображение. Благодаря этому тот же механизм работает и для вычислений в отдельных процессах.
Описание дизайна и особенностей работы с рантайм подсистемой потребует отдельной статьи (которую я бы с радостью написал). Здесь я бы хотел показать как это выглядит на простом примере.
В любой пайплайн, который работает с потоком данных, можно встроить оператор мониторинга, который будет посылать событие инструменту отображения прогресса.
# Создаем оператор/ноду прогресса progress = IcoProgress(name="Overall progress", total=100) # интегрируем ее в пайплайн pipeline = source | progress | processing | train

При выполнении пайплайна прогресс будет отображаться в реальном времени. При этом возможно наличие нескольких нод прогресса, в данном примере это два воркера.
Многопроцессорное выполнение
Для работы в реальных условиях нам часто требуется многозадачность. И для меня было важно сделать ее частью фреймворка, не нарушая его целостности. Так получился MPAgent - он имеет интерфейс оператора и выполняет пайплайн в отдельном процессе. С его помощью даже строится аналог PyTorch DataLoader, не нарушая целостности ICO.
В примере ниже, асинхронный поток использует пул воркеров, который состоит из нескольких MPAgent. Каждый агент выполняется в отдельном процессе и использует фабричный метод для создания пайплайна. Он получает данные, выполняет пайплайн и передает результат обратно.
Пример асинхронного пула воркеров.
workers = IcoAsyncStream( lambda: MPAgent(heavy_computation), pool_size=cpu_count() ) # встраиваем пул воркеров в пайплайн pipeline = source | workers | train
Многопроцессорное выполнение встраивается в тот же вычислительный пайплайн и не требует изменения интерфейса операторов.
Заключение
Я проделал путь почти в сотню коммитов, пытаясь создать удобную и стройную систему — модель вычислений, которая решала бы определённый, хорошо знакомый мне класс задач. Как минимум, это было интересное исследование и в каком-то смысле даже открытие для меня — что так тоже можно сделать.
Получился ли из этого действительно удобный инструмент? На этот вопрос мне как раз хотелось бы получить ответ от сообщества.
В этом посте я обозначил одни из основных особенностей ICO, но это далеко не все. Буду рад ответить на вопросы и раскрыть темы в следующих постах. Это можно было бы превратить в цикл статей, где я могу описать среду выполнения, архитектуру или конкретные примеры использования.
Сам исходный ICO код доступен на GitHub.
Для желающих ознакомиться более подробно, есть следующий набор материалов:
📖 Примеры
Примеры представлены в виде Jupyter-ноутбуков и могут быть запущены прямо в Google Colab — без дополнительной настройки.
Основы ICO
Базовое введение в ICO — основные строительные блоки и ключевые концепции. Открыть в Colab
Введение в ICO Runtime — отслеживание прогресса, вывод информации и архитектура среды выполнения Открыть в Colab
Многопроцессорная обработка
Примеры с многопроцессорной обработкой нельзя запустить в Jupyter или Google Colab. Для их запуска необходимо установить фреймворк локально и выполнить скрипты из терминала.
Инструкции по настройке см в разделе Установка.
Пример многопроцессорной обработки — базовый пример распределённых вычислительных потоков
Параллельный многопроцессорный пул — вычислительные потоки с параллельным пулом воркеров
Машинное обучение
Линейная регрессия — разработка ML-пайплайна на базе ICO Открыть в Colab
Классификация CIFAR-10 с валидацией — полный CV-пайплайн, заменяющий PyTorch DataLoader Открыть в Colab
Классификация CIFAR-10 с пулом воркеров — полный CV-пайплайн с параллельной обработкой данных

