Repository navigation
Conversation
| except queue.ShutDown: | ||
| break | ||
|
|
||
| try: |
There was a problem hiding this comment.
🟠 Ошибка обработки не предотвращает фиксацию 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, | ||
| }, |
There was a problem hiding this comment.
🟠 Использование 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: |
There was a problem hiding this comment.
🟠 Гонка при завершении: флаг 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 |
There was a problem hiding this comment.
🟡 Множество 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 = {} |
There was a problem hiding this comment.
🟡 Линейный скан 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": |
There was a problem hiding this comment.
🟡 Разделяемый 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) |
There was a problem hiding this comment.
🔵 Очереди создаются без maxsize — отсутствует backpressure, заявленный в докстринге
Докстринг класса утверждает «Natural backpressure when queues fill up», но queue.Queue() создаётся без параметра maxsize, то есть очередь не ограничена. Если потребление сообщений опережает обработку, память будет расти без ограничений. Следует задать maxsize или убрать заявление о backpressure из документации.
| worker.join(timeout=5.0) | ||
|
|
||
|
|
||
| class SimpleQueueProcessingStrategy(ProcessingStrategy[KafkaPayload], Generic[T]): |
There was a problem hiding this comment.
🔵 Докстринг SimpleQueueProcessingStrategy заявляет коммит только после успешной обработки, что не соответствует реализации
Докстринг класса перечисляет гарантию «Only commits offsets after successful processing». Однако, как отмечено в c2, offset фиксируется в finally независимо от результата result_processor. Докстринг вводит в заблуждение будущих разработчиков. Следует либо исправить реализацию, либо обновить документацию.
|
|
||
| committable = self.queue_pool.offset_tracker.get_committable_offsets() | ||
|
|
||
| if committable: |
There was a problem hiding this comment.
🔵 Метрика 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.
🔍 Автоматическое ревью кодаКраткий итогСинтез ревью: 16 уникальных групп замечаний по 2 исходн. выполнениям; 10 для публикации (критических: 0, высоких: 3, средних: 3, низких: 4, мелких: 0). Замечания🟠 Важно
💡 Дополнительно
Автоматическое ревью от Review Engine v2.0 |
| self.shutdown_event.set() | ||
| self.commit_thread.join(timeout=5.0) | ||
| self.queue_pool.shutdown() | ||
|
|
There was a problem hiding this comment.
🔵 Метод 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().
Benchmark fixture for sentry PR #95633.\n\nGolden comments:
benchmark/golden_comments/sentry.json— PR95633.