Skip to content

[benchmark][Python] feat(uptime): Add ability to use queues to manage parallelism - #18

Open
hsander wants to merge 1 commit into
benchmark/sentry/pr-95633-basefrom
benchmark/sentry/pr-95633
Open

hsander wants to merge 1 commit into
benchmark/sentry/pr-95633-basefrom
benchmark/sentry/pr-95633

Conversation

@hsander

@hsander hsander commented Jul 1, 2026

Copy link
Copy Markdown
Owner

Benchmark fixture for sentry PR #95633.\n\nGolden comments: benchmark/golden_comments/sentry.json — PR 95633.

except queue.ShutDown:
break

try:

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 Ошибка обработки не предотвращает фиксацию offset — сообщение теряется безвозвратно

В OrderedQueueWorker.run() вызов result_processor обёрнут в try/except/finally. При возбуждении исключения оно лишь логируется, а блок finally безусловно вызывает self.offset_tracker.complete_offset(...). Это удаляет offset из набора outstanding, после чего get_committable_offsets() продвинет коммит past проваленного сообщения. Необходимо вызывать complete_offset только при успешной обработке, либо переводить partition в error-state и блокировать коммит.

self.work_queue.qsize(),
tags={
"identifier": self.identifier,
},

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 Использование queue.ShutDown и Queue.shutdown() — API, доступного только в Python 3.13+

Код использует исключение queue.ShutDown и метод queue.Queue.shutdown(immediate=False), появившиеся в Python 3.13. Если runtime использует Python 3.12 или ранее (что типично для Sentry), будет выброшен AttributeError при любом вызове shutdown, а обработчик except queue.ShutDown также сломается, так как ShutDown не определён. Это затрагивает все пути завершения. Требуется либо версионная проверка, либо альтернативная реализация shutdown через sentinel-объект в очереди.


def shutdown(self) -> None:
"""Gracefully shutdown all workers."""
for worker in self.workers:

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 Гонка при завершении: флаг worker.shutdown устанавливается до q.shutdown(), что может привести к потере сообщений в очереди

В FixedQueuePool.shutdown() для всех воркеров сначала устанавливается worker.shutdown = True, и только затем вызывается q.shutdown(immediate=False). Воркер, завершивший обработку элемента и вернувшийся к проверке while not self.shutdown, увидит флаг и выйдет, не обработав оставшиеся элементы в очереди. Поскольку каждый воркер — единственный потребитель своей очереди, эти элементы будут потеряны. Следует менять порядок: сначала q.shutdown(), затем устанавливать флаг, либо использовать механизм drain-очереди.

if offset in all_offsets and offset not in outstanding:
highest_committable = offset
else:
break

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Множество outstanding в OffsetTracker никогда не очищается — неограниченный рост памяти

В методе mark_committed фильтруется all_offsets (удаляются offset <= committed), но outstanding не очищается. Завершённые и закоммиченные offset остаются в outstanding навсегда. При длительной работе потребителя множество растёт без ограничений, расходуя память и замедляя проверки offset not in outstanding. Необходимо удалять из outstanding все offset <= committed в mark_committed.


For each partition, finds the highest contiguous offset that has been processed.
"""
committable = {}

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Линейный скан range(start, max_offset + 1) в get_committable_offsets при разреженных offset

Метод get_committable_offsets итерирует range(start, max_offset + 1) для каждой партиции. Если offset разрежены (например, 0 и 1 000 000 с разрывом), цикл выполнится до 1M итераций. Вызов происходит каждую секунду в commit-цикле и под блокировкой партиции, что может вызвать существенную нагрузку на CPU и задержки. Следует использовать структуру данных с O(1) проверкой принадлежности (например, set или отсортированный список завершённых offset).

if num_processes is None:
num_processes = multiprocessing.cpu_count()
self.multiprocessing_pool = MultiprocessingPool(num_processes)
if mode == "thread-queue-parallel":

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Разделяемый queue_pool порождает несколько commit-потоков при ребалансе партиций, что ведёт к двойным коммитам

FixedQueuePool (и его единственный OffsetTracker) создаётся один раз в ResultsStrategyFactory.__init__. Каждый вызов create_thread_queue_parallel_worker создаёт новый SimpleQueueProcessingStrategy, конструктор которого безусловно запускает новый commit_thread, работающий с разделяемым queue_pool.offset_tracker. При ребалансе партиций arroyo вызывает create_with_partitions повторно, порождая дополнительные commit-потоки. Это может привести к гонкам при коммите и двойной фиксации offset. Следует гарантировать единственный commit-поток на queue_pool.

offset_tracker=self.offset_tracker,
)
worker.start()
self.workers.append(worker)

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Очереди создаются без maxsize — отсутствует backpressure, заявленный в докстринге

Докстринг класса утверждает «Natural backpressure when queues fill up», но queue.Queue() создаётся без параметра maxsize, то есть очередь не ограничена. Если потребление сообщений опережает обработку, память будет расти без ограничений. Следует задать maxsize или убрать заявление о backpressure из документации.

worker.join(timeout=5.0)


class SimpleQueueProcessingStrategy(ProcessingStrategy[KafkaPayload], Generic[T]):

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Докстринг SimpleQueueProcessingStrategy заявляет коммит только после успешной обработки, что не соответствует реализации

Докстринг класса перечисляет гарантию «Only commits offsets after successful processing». Однако, как отмечено в c2, offset фиксируется в finally независимо от результата result_processor. Докстринг вводит в заблуждение будущих разработчиков. Следует либо исправить реализацию, либо обновить документацию.


committable = self.queue_pool.offset_tracker.get_committable_offsets()

if committable:

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Метрика offsets_committed инкрементируется до успешного выполнения commit_function

В _commit_loop вызов metrics.incr('remote_subscriptions.queue_pool.offsets_committed', ...) происходит до self.commit_function(committable). Если commit_function возбуждает исключение (перехватывается широким except), mark_committed пропускается, но метрика уже учла несостоявшийся коммит. Повторные сбои коммита будут раздувать метрику. Следует инкрементировать метрику после успешного вызова commit_function.

@hsander

hsander commented Jul 2, 2026

Copy link
Copy Markdown
Owner Author

🔍 Автоматическое ревью кода

Краткий итог

Синтез ревью: 16 уникальных групп замечаний по 2 исходн. выполнениям; 10 для публикации (критических: 0, высоких: 3, средних: 3, низких: 4, мелких: 0).

Замечания

🟠 Важно

  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:135] Ошибка обработки не предотвращает фиксацию offset — сообщение теряется безвозвратно — В OrderedQueueWorker.run() вызов result_processor обёрнут в try/except/finally. При возбуждении исключения оно лишь логируется, а блок finally безусловно вызывает self.offset_tracker.complete_offset(...). Это удаляет offset из набора outstanding, после чего get_committable_offsets() продвинет коммит past проваленного сообщения. Необходимо вызывать complete_offset только при успешной обработке, либо переводить partition в error-state и блокировать коммит.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:155] Использование queue.ShutDown и Queue.shutdown() — API, доступного только в Python 3.13+ — Код использует исключение queue.ShutDown и метод queue.Queue.shutdown(immediate=False), появившиеся в Python 3.13. Если runtime использует Python 3.12 или ранее (что типично для Sentry), будет выброшен AttributeError при любом вызове shutdown, а обработчик except queue.ShutDown также сломается, так как ShutDown не определён. Это затрагивает все пути завершения. Требуется либо версионная проверка, либо альтернативная реализация shutdown через sentinel-объект в очереди.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:233] Гонка при завершении: флаг worker.shutdown устанавливается до q.shutdown(), что может привести к потере сообщений в очереди — В FixedQueuePool.shutdown() для всех воркеров сначала устанавливается worker.shutdown = True, и только затем вызывается q.shutdown(immediate=False). Воркер, завершивший обработку элемента и вернувшийся к проверке while not self.shutdown, увидит флаг и выйдет, не обработав оставшиеся элементы в очереди. Поскольку каждый воркер — единственный потребитель своей очереди, эти элементы будут потеряны. Следует менять порядок: сначала q.shutdown(), затем устанавливать флаг, либо использовать механизм drain-очереди.

💡 Дополнительно

  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:93] Множество outstanding в OffsetTracker никогда не очищается — неограниченный рост памяти — В методе mark_committed фильтруется all_offsets (удаляются offset <= committed), но outstanding не очищается. Завершённые и закоммиченные offset остаются в outstanding навсегда. При длительной работе потребителя множество растёт без ограничений, расходуя память и замедляя проверки offset not in outstanding. Необходимо удалять из outstanding все offset <= committed в mark_committed.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:73] Линейный скан range(start, max_offset + 1) в get_committable_offsets при разреженных offset — Метод get_committable_offsets итерирует range(start, max_offset + 1) для каждой партиции. Если offset разрежены (например, 0 и 1 000 000 с разрывом), цикл выполнится до 1M итераций. Вызов происходит каждую секунду в commit-цикле и под блокировкой партиции, что может вызвать существенную нагрузку на CPU и задержки. Следует использовать структуру данных с O(1) проверкой принадлежности (например, set или отсортированный список завершённых offset).
  • [src/python/src/sentry/remote_subscriptions/consumers/result_consumer.py:40] Разделяемый queue_pool порождает несколько commit-потоков при ребалансе партиций, что ведёт к двойным коммитам — FixedQueuePool (и его единственный OffsetTracker) создаётся один раз в ResultsStrategyFactory.__init__. Каждый вызов create_thread_queue_parallel_worker создаёт новый SimpleQueueProcessingStrategy, конструктор которого безусловно запускает новый commit_thread, работающий с разделяемым queue_pool.offset_tracker. При ребалансе партиций arroyo вызывает create_with_partitions повторно, порождая дополнительные commit-потоки. Это может привести к гонкам при коммите и двойной фиксации offset. Следует гарантировать единственный commit-поток на queue_pool.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:196] Очереди создаются без maxsize — отсутствует backpressure, заявленный в докстринге — Докстринг класса утверждает «Natural backpressure when queues fill up», но queue.Queue() создаётся без параметра maxsize, то есть очередь не ограничена. Если потребление сообщений опережает обработку, память будет расти без ограничений. Следует задать maxsize или убрать заявление о backpressure из документации.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:246] Докстринг SimpleQueueProcessingStrategy заявляет коммит только после успешной обработки, что не соответствует реализации — Докстринг класса перечисляет гарантию «Only commits offsets after successful processing». Однако, как отмечено в c2, offset фиксируется в finally независимо от результата result_processor. Докстринг вводит в заблуждение будущих разработчиков. Следует либо исправить реализацию, либо обновить документацию.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:280] Метрика offsets_committed инкрементируется до успешного выполнения commit_function — В _commit_loop вызов metrics.incr('remote_subscriptions.queue_pool.offsets_committed', ...) происходит до self.commit_function(committable). Если commit_function возбуждает исключение (перехватывается широким except), mark_committed пропускается, но метрика уже учла несостоявшийся коммит. Повторные сбои коммита будут раздувать метрику. Следует инкрементировать метрику после успешного вызова commit_function.
  • [src/python/src/sentry/remote_subscriptions/consumers/queue_consumer.py:339] Метод join() безусловно вызывает close(), что приводит к двойному shutdown — Жизненный цикл arroyo ProcessingStrategy обычно вызывает close(), затем join(). Поскольку join() вызывает self.close() повторно, происходит второй shutdown_event.set(), второй commit_thread.join() и второй queue_pool.shutdown(). Повторный queue_pool.shutdown() попытается вызвать q.shutdown() на уже остановленных очередях и worker.join() на уже присоединённых воркерах, что может вызвать исключения. Следует добавить флаг защиты от повторного вызова close().

Автоматическое ревью от Review Engine v2.0

self.shutdown_event.set()
self.commit_thread.join(timeout=5.0)
self.queue_pool.shutdown()

Copy link
Copy Markdown
Owner Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔵 Метод join() безусловно вызывает close(), что приводит к двойному shutdown

Жизненный цикл arroyo ProcessingStrategy обычно вызывает close(), затем join(). Поскольку join() вызывает self.close() повторно, происходит второй shutdown_event.set(), второй commit_thread.join() и второй queue_pool.shutdown(). Повторный queue_pool.shutdown() попытается вызвать q.shutdown() на уже остановленных очередях и worker.join() на уже присоединённых воркерах, что может вызвать исключения. Следует добавить флаг защиты от повторного вызова close().

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant