АрхитектураРаспределённые системыИнженер по распределённым системам

В группе потребителей один участник внезапно перестал подтверждать сообщения. Как перераспределение партици...

В группе потребителей один участник внезапно перестал подтверждать сообщения. Как перераспределение партиций восстанавливает обработку?

Проходите собеседования с ИИ помощником Hintsage

Краткий ответ

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

Перераспределение восстанавливает доступность обработки, но само по себе не гарантирует отсутствие дублей или потерь. Надёжность зависит от момента фиксации смещения и от того, насколько корректно обработчик повторяет бизнес-операцию.

Исторический контекст

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

При изменении состава группы нужно определить единственного текущего владельца каждой партиции. Поэтому используется процедура rebalance: группа временно или поэтапно отзывает назначения и распределяет их заново.

Постановка проблемы

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

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

Подробное решение

Участники группы периодически посылают координатору heartbeat и сообщают о прогрессе обработки. При превышении допустимого интервала координатор считает участника недоступным, пересматривает состав группы и назначает его партиции оставшимся потребителям.

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

Новый владелец читает с позиции, сохранённой в хранилище смещений. Поэтому типичный результат при сбое — семантика at-least-once: сообщение может быть обработано повторно, но не должно исчезнуть из-за неподтверждённого прогресса.

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

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

Ребалансировка не должна рассматриваться как транзакция между чтением сообщения и бизнес-операцией. Для защиты от повторов обработчик обычно делает операцию идемпотентной, использует уникальный идентификатор события или атомарно связывает изменение состояния с фиксацией прогресса. Если бизнес-операция необратима, одной корректной передачи партиций недостаточно.

Ситуация из практики

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

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

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

Что кандидаты часто упускают

1. Гарантирует ли успешная ребалансировка отсутствие повторной обработки сообщений?

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

2. Почему нельзя считать партицию свободной сразу после первого пропущенного heartbeat?

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

3. Может ли кооперативная ребалансировка полностью исключить одновременную обработку одной партиции старым и новым потребителем?

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