759427794 opened a new issue, #1331: URL: https://github.com/apache/rocketmq-clients/issues/1331
### Before Creating the Bug Report - [x] I found a bug, not just asking a question, which should be created in [GitHub Discussions](https://github.com/apache/rocketmq-clients/discussions). - [x] I have searched the [GitHub Issues](https://github.com/apache/rocketmq-clients/issues) and [GitHub Discussions](https://github.com/apache/rocketmq-clients/discussions) of this repository and believe that this is not a duplicate. - [x] I have confirmed that this bug belongs to the current repository, not other repositories of RocketMQ. ### Programming Language of the Client Python ### Runtime Platform Environment ### Environment - **SDK**: `rocketmq-python-client==5.1.1` - **Python**: 3.8 - **RocketMQ Server**: 5.x - **OS**: Linux (CentOS 7 / Ubuntu 20.04) ### Describe the bug When `PushConsumer`'s process queue cache is full and messages keep arriving faster than they can be consumed, `__execute_receive` and `__execute_receive_later` enter a synchronous mutual recursion. Each call adds a stack frame, and after ~800-1000 seconds (~13-16 minutes), Python's default recursion limit (1000) is exceeded and the process crashes with `RecursionError: maximum recursion depth exceeded`. ### Reproducer Run a PushConsumer with slow message processing against a topic with sustained high throughput. For example, a video-processing consumer that takes 30-90 seconds per message, with the default `consumption_thread_count=20` and `receive_batch_size=32`. ### Stack trace Thread 0x00007f0abe7fc700 (most recent call first): File "logging/init.py", line 424, in usesTime File "logging/init.py", line 633, in usesTime File "logging/init.py", line 665, in format File "logging/init.py", line 925, in format File "logging/init.py", line 1081, in emit File "logging/init.py", line 1183, in emit File "logging/handlers.py", line 71, in emit File "logging/init.py", line 950, in handle File "logging/init.py", line 1649, in callHandlers File "logging/init.py", line 1587, in handle File "logging/init.py", line 1577, in _log File "logging/init.py", line 1446, in warning File "logging/init.py", line 1451, in warn File "rocketmq/v5/model/process_queue.py", line 66, in is_cache_full File "rocketmq/v5/consumer/push/push_consumer.py", line 235, in __execute_receive File "rocketmq/v5/consumer/push/push_consumer.py", line 271, in __execute_receive_later File "rocketmq/v5/consumer/push/push_consumer.py", line 237, in __execute_receive File "rocketmq/v5/consumer/push/push_consumer.py", line 271, in __execute_receive_later File "rocketmq/v5/consumer/push/push_consumer.py", line 237, in __execute_receive File "rocketmq/v5/consumer/push/push_consumer.py", line 271, in __execute_receive_later ... (hundreds more identical frames) ... ### Root Cause The call chain forms unbounded synchronous recursion: __execute_receive (line 237) → is_cache_full() returns True → __execute_receive_later (line 271) → sleep(1s) → __execute_receive (line 237) ← back to start, +1 stack frame 1 stack frame is added per second. After ~800-1000 seconds → `RecursionError` crash. Additionally, `process_queue.py:66` `is_cache_full()` calls `logging.warn()` on each check, whose `emit → format → usesTime` path consumes extra stack frames per iteration, accelerating the overflow. ### Suggested fix Replace the recursive call in `__execute_receive_later` with an iterative `while` loop: ```python # Before def __execute_receive_later(self, message_queue, process_queue, attempt_id): time.sleep(PushConsumer.RECEIVE_RETRY_DELAY) self.__execute_receive(message_queue, process_queue, attempt_id) # After def __execute_receive_later(self, message_queue, process_queue, attempt_id): while self.is_running and not process_queue.dropped: time.sleep(PushConsumer.RECEIVE_RETRY_DELAY) if not self.is_running or process_queue.dropped: return if not process_queue.is_cache_full( self.__queue_threshold(self.__max_cache_message_count), self.__queue_threshold(self.__max_cache_message_size)): self.__execute_receive(message_queue, process_queue, attempt_id) return This preserves identical behavior (sleep + retry) but uses O(1) stack space. I'm happy to submit a PR with this fix if the approach looks good. --- ### RocketMQ Version of the Client/Server rocketmq-python-client==5.1.1 ### Run or Compiler Version _No response_ ### Describe the Bug PushConsumer 在消息处理速度跟不上拉取速度时,`__execute_receive` 和 `__execute_receive_later` 两个方法会形成同步递归。每秒钟增加一层栈帧,约 13-16 分钟后 Python 栈溢出(RecursionError),进程崩溃。 具体链路: 1. `__execute_receive` 检查 `is_cache_full()` → 缓存满了 2. 调用 `__execute_receive_later` → sleep 1秒 → 回调 `__execute_receive` 3. 缓存仍满 → 再次调用 `__execute_receive_later` → 无限循环 4. 每次循环 `is_cache_full` 还调用 `logging.warn()`,额外消耗栈帧 5. 约 800-1000 次递归后超出 Python 默认递归限制 1000,崩溃 线程堆栈中可见数百个重复的: `push_consumer.py:237 __execute_receive → push_consumer.py:271 __execute_receive_later → push_consumer.py:237 __execute_receive → ...` ### Steps to Reproduce 1. 启动 RocketMQ 5.x 服务端 2. 创建一个 PushConsumer,使用默认配置(consumption_thread_count=20, max_cache_message_count=1024) 3. 在 MessageListener.consume() 中加入较长的处理耗时(如 `time.sleep(5)` 模拟慢消费) 4. 向对应 topic 持续高频发送消息(如每秒 100 条) 5. 运行约 15-20 分钟后,进程崩溃并抛出 `RecursionError: maximum recursion depth exceeded` 环境信息: - SDK: rocketmq-python-client==5.1.1 - Python: 3.8 - OS: Linux ### What Did You Expect to See? PushConsumer 在缓存满时应该使用迭代式循环等待(while + sleep),而不是递归调用。进程应该稳定运行,无论缓存满多久都不会因栈溢出而崩溃 ### What Did You See Instead? 进程运行约 15-20 分钟后崩溃,抛出 RecursionError: maximum recursion depth exceeded。 线程堆栈中可见数百个重复帧: push_consumer.py:237 __execute_receive push_consumer.py:271 __execute_receive_later push_consumer.py:237 __execute_receive ...(重复数百次) process_queue.py:66 is_cache_full → logging.warn logging/__init__.py emit → format → usesTime(日志调用耗尽最后栈空间) ### Additional Context _No response_ -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
