googleapis / googleapis/python-aiplatform
[Agent Engines] `_wrap_async_stream_query_operation` blocks asyncio event loop due to synchronous gRPC iteration
- 主要语言
- Python
- 星标
- 905
- 派生
- 465
- 平均合并
- 1 天 13 小时
- 30 天内合并 PR
- 44
描述
`agent_engine.async_stream_query(...)` blocks the calling `asyncio` event loop thread during stream generation.
In `vertexai/agent_engines/_agent_engines.py`, `_wrap_async_stream_query_operation` wraps a synchronous client call in an `async def` and uses a blocking `for` loop:
```python
# CURRENT IMPLEMENTATION (vertexai/agent_engines/_agent_engines.py)
def _wrap_async_stream_query_operation(*, method_name: str):
async def _method(self, **kwargs):
# 1. Uses sync client instead of self.execution_async_client
response = self.execution_api_client.stream_query_reasoning_engine(...)
# 2. Synchronous iteration blocks the asyncio event loop on socket reads
for chunk in response:
for parsed_json in _utils.yield_parsed_json(chunk):
if parsed_json is not None:
yield parsed_json
return _method
```
### Reproduction:
```python
import asyncio
from vertexai import agent_engines
agent = agent_engines.get("projects/
/locations//reasoningEngines/")
async def main():
async for chunk in agent.async_stream_query(user_id="u", message="Long prompt"):
pass
asyncio.run(main(), debug=True)
# Result: Emits "Executing took X.XX seconds" because the thread is blocked on socket reads without yielding.
```
贡献指南
调研方向
Start in vertexai/agent_engines/_agent_engines.py at _wrap_async_stream_query_operation and compare the synchronous execution_api_client call with self.execution_async_client. Reproduce the stream with asyncio.run(..., debug=True); done means async_stream_query no longer blocks the event-loop thread during socket reads while still yielding parsed chunks.
由索引模型根据 Issue 内容生成。
评估
- 技术栈
- grpc, python
- 领域
- api, backend
- Issue 类型
- 缺陷
- 难度
- 3/5
- 预计耗时
- 1-2 天
- 活跃度
- 活跃
- 描述清晰度
- 描述清楚
- 新手友好度
- 72/100