使用Python asyncio向Dialogflow流式传输音频时触发CancelledError
问题:Dialogflow流式音频交互触发asyncio.CancelledError
我正在编写自定义Dialogflow集成程序,基于现有示例修改实现音频流式传输至Dialogflow的streaming_detect_intent()接口,音频采集采用sounddevice库。程序设置两个Python异步任务:一个负责采集音频,另一个处理Dialogflow交互,二者通过共享队列通信。但流式传输音频时触发asyncio.exceptions.CancelledError(),报错栈如下:
File "/home/andrew/experiments/messaging/a_recording.py", line 95, in sample_streaming_detect_intent async for response in stream: File "/home/andrew/venv/lib/python3.11/site-packages/google/api_core/grpc_helpers_async.py", line 102, in _wrapped_aiter async for response in self._call: # pragma: no branch File "/home/andrew/venv/lib/python3.11/site-packages/grpc/aio/_call.py", line 327, in _fetch_stream_responses await self._raise_for_status() File "/home/andrew/venv/lib/python3.11/site-packages/grpc/aio/_call.py", line 233, in _raise_for_status raise asyncio.CancelledError() asyncio.exceptions.CancelledError
Dialogflow处理任务代码
async def sample_streaming_detect_intent( loop, audio_queue, project_id, session_id, sample_rate ): client = dialogflow.SessionsAsyncClient() audio_config = dialogflow.InputAudioConfig( audio_encoding=dialogflow.AudioEncoding.AUDIO_ENCODING_LINEAR_16, language_code="en", sample_rate_hertz=sample_rate, ) async def request_generator(loop, project_id, session_id, audio_config, audio_queue): query_input = dialogflow.QueryInput(audio_config=audio_config) # Initialize request argument(s) yield dialogflow.StreamingDetectIntentRequest( session=client.session_path(project_id, session_id), query_input=query_input ) while True: chunk = await audio_queue.get() if not chunk: break # The later requests contains audio data. yield dialogflow.StreamingDetectIntentRequest(input_audio=chunk) # Make the request client_task = asyncio.create_task( client.streaming_detect_intent( requests=request_generator( loop, project_id, session_id, audio_config, audio_queue ) ) ) try: stream = await client_task except Exception as e: print(f"failed with {e.__cause__}") try: async for response in stream: print(response) except Exception as e: print(f"failed with {e.__cause__}") query_result = response.query_result print("=" * 20) print("Query text: {}".format(query_result.query_text)) print( "Detected intent: {} (confidence: {})\n".format( query_result.intent.display_name, query_result.intent_detection_confidence ) ) print("Fulfillment text: {}\n".format(query_result.fulfillment_text))
调用代码片段
audio_queue = asyncio.Queue() # to assert that we are using the same event loop loop = asyncio.get_event_loop() await asyncio.gather( record_audio(fp, loop, audio_queue, sample_rate, device, channels), sample_streaming_detect_intent( loop, audio_queue, project_id, session_id, sample_rate ), )
已尝试的方案
- 添加任务将流式音频写入文件
- 编写任务读取文件音频并发送至Dialogflow任务
- 验证确保使用同一事件循环
- 在
read_audio()任务中使用loop.call_soon_threadsafe
已确认音频设置正确,调试发现问题源自GRPC,查阅CancelledError相关资料后仍未解决。
内容的提问来源于stack exchange,提问作者user2607755
相关产品推荐
相关产品推荐

