Программирование и разработка программного обеспечения

Как построить конвейер Kafka, сохраняющий порядок сессий, используя Go

Joshua Oluikpe описывает практическую архитектуру обработки сообщений Kafka в строгом порядке внутри каждой сессии с возможностью параллельной работы независимых сессий. Архитектура использует согласованное хеширование, отдельную горутину Go для каждой сессии, повторные попытки в том же пути обработки и управление соседними точками фиксации, чтобы избежать потери сообщений после сбоя.

2026-10-07
4 мин. чтения
1 просмотров
certi.news Editorial Team
Как построить конвейер Kafka, сохраняющий порядок сессий, используя Go

Недостаточно полагаться на порядок Kafka внутри раздела (partition), когда тысячи сессий используют один и тот же раздел. В системах чатов и последовательных задач исправляющее сообщение вроде «сделай пунктом назначения Париж» может прийти после сообщения «забронируй поездку в Лондон», но его выполнение перед первым сообщением нарушит контекст. Решение, которое Joshua Oluikpe описывает в InfoQ, переносит часть ответственности за порядок на прикладной слой, написанный на Go.

Порядок внутри сессии и параллелизм между сессиями

Архитектура опирается на два уровня воркеров. Уровень распределения получает записи Kafka и направляет их по идентификатору сессии с помощью согласованного хеширования, так что сообщения одной сессии попадают к одному и тому же воркеру, а разные сессии можно распределять между несколькими воркерами. Затем для каждой активной сессии создаётся отдельная goroutine, которая обрабатывает по одному сообщению за раз в порядке поступления.

Благодаря этому медленная сессия не становится препятствием для других сессий, как это происходит в плоском пуле воркеров, когда повторная попытка для одной сессии полностью занимает воркер. Goroutine создаются по мере необходимости и удаляются после периода бездействия, поэтому потребление ресурсов зависит от числа активных сессий, а не от общего числа возможных сессий.

Повторные попытки без пропуска сообщений

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

Неисправимые ошибки, такие как некорректные данные, сразу направляются в очередь недоставленных сообщений (DLQ). Такое разделение временных и окончательных ошибок устраняет необходимость в отдельном автомате состояний для управления повторными попытками, но также означает, что сообщения определённой сессии могут накапливаться в период задержки.

Безопасная фиксация и восстановление после сбоя

Когда несколько сессий обрабатывают свои сообщения параллельно, нельзя фиксировать только наибольший завершённый offset. Если offset 104 завершён раньше 102 и система фиксирует 104, сбой потребителя может окончательно привести к пропуску 102. Поэтому система отслеживает для каждого раздела находящиеся в обработке и завершённые offsets и перемещает точку фиксации только до наибольшего соседнего завершённого offset.

В результате перезапуск может начаться с более раннего offset и повторно обработать некоторые сообщения. Архитектура использует постоянный идентификатор события, помогающий устранять дубликаты, но не заявляет о достижении обработки «ровно один раз»: внешние эффекты по-прежнему зависят от способности зависимых служб идемпотентно обрабатывать повторы. В редких случаях застрявший пробел можно пропустить после истечения заданного времени ожидания, чтобы восстановить продвижение раздела, зафиксировав это как операционный сбой и направив сообщение в DLQ. Это явный компромисс между непрерывностью работы и полнотой обработки.

Что требуется для реальной эксплуатации?

Одной базовой логики недостаточно для полной гарантии порядка. Во время перераспределения разделов потребитель прекращает принимать новые записи, ждёт завершения текущей работы, а затем фиксирует последний соседний watermark перед передачей владения. Когда каналы воркеров заполняются, система временно прекращает извлечение сообщений вместо их отбрасывания, сохраняя уже извлечённые записи в хранилищах, выделенных для каждого раздела, чтобы более новые сообщения не опережали их.

Системе также требуется индивидуальное обнаружение застрявших offsets, сквозная трассировка через асинхронные границы Kafka, сохранение сигналов отставания, необходимых для автоматического масштабирования во время остановки извлечения, а также завершение обработки и запись в DLQ перед закрытием. Продукт также применяет согласованное хеширование на стороне отправки, чтобы обеспечивать порядок сообщений при асинхронном производстве.

Результаты и ограничения

В синтетическом тесте с 50 000 сообщениями и 10 сессиями, где временные ошибки были искусственно внесены в 10% сообщений, система обработала 14 027 сообщений в секунду без нарушений порядка, а медианная задержка составила 11,7 миллисекунды. В нагрузочном тесте с 1 000 сессиями фактическая скорость отправки достигла 48 805 сообщений в секунду при целевом показателе 50 000, без ошибок отправки, однако хвостовая задержка превысила бюджет в 200 миллисекунд из-за ограничения скорости сообщений в API Gateway, а не из-за механизма упорядочивания сессий.

Автор сообщает, что система обработала более 40 миллионов сообщений в рабочей среде без заметных нарушений порядка, направив в DLQ всего 460 сообщений, то есть около 0,002%. Однако эта цифра основана на операционном мониторинге, а не на полномасштабной проверке последовательности каждого сообщения; кроме того, в более высоком нагрузочном тесте порядок для каждой сессии независимо не проверялся. Поэтому практическая ценность архитектуры заключается в превращении порядка из предположения на уровне Kafka в свойство, обеспечиваемое приложением, при сохранении необходимости в явных тестах последовательности и гарантиях идемпотентности зависимых служб.

Источник новости
InfoQ - Architecture Articles
Открыть первоисточник ↗
c
Автор

certi.news Editorial Team

В той же категории

Вам также может понравиться

Все новости