Amazon Transcribe Python API事件处理器仅在流结束后处理音频
问题根因
流式转写不生效是四个核心用法错误导致的:
handle_events()是持续运行的异步事件监听循环,你把它放在receive回调里每次收到音频块才调用一次,每次仅执行单轮事件检查就退出,只有流关闭触发终态事件时才能拿到结果,无法实时处理中间返回的转写片段。- 之前用
asyncio.create_task创建监听任务立刻结束,是因为Django Channels的异步消费者会自动回收未绑定到实例生命周期的孤儿协程,任务还没跑起来就被取消了。 - 为每个音频块创建独立转写流报内部错误属于预期行为:Amazon Transcribe流式接口要求单条连接持续按序发送完整编码序列的音频,碎片化的单块音频缺少编码头、帧序列不连续,服务端无法正常解析。
- 额外逻辑bug:转写回调里硬编码了
channel_name = "test",没有传入当前WebSocket连接的实际channel名称,就算拿到转写结果也无法正确推送到当前前端。
修复方案
- 在连接建立、启动Transcribe流之后,立刻在后台启动事件监听协程,把任务绑定到consumer实例属性上避免被回收。
- 不要在
receive方法里重复调用handle_events(),这个方法本身是无限循环,只要启动一次就会持续监听转写返回。 - 断开连接时主动取消监听任务、正确关闭Transcribe流,避免资源泄漏。
- 修正参数校验:确认前端发送的音频编码、采样率和调用Transcribe接口时传入的
media_encoding、media_sample_rate_hz完全匹配,浏览器端MediaRecorder输出的ogg/opus默认是48kHz,注意不要传裸Opus帧却声明是ogg-opus封装。
修正后可运行代码
import json import asyncio from channels.generic.websocket import AsyncWebsocketConsumer from amazon_transcribe.client import TranscribeStreamingClient from amazon_transcribe.handlers import TranscriptResultStreamHandler from amazon_transcribe.model import TranscriptEvent from channels.layers import get_channel_layer stream_client = TranscribeStreamingClient(region="us-west-2") class AWSTranscriptHandler(TranscriptResultStreamHandler): def __init__(self, transcript_result_stream, channel_layer, channel_name): self.channel_layer = channel_layer self.channel_name = channel_name super().__init__(transcript_result_stream) async def handle_transcript_event(self, transcript_event: TranscriptEvent): results = transcript_event.transcript.results for result in results: # 去掉is_partial判断可以返回实时中间转写结果,保留则只返回最终确定的片段 for alt in result.alternatives: await self.channel_layer.send( self.channel_name, {"type": "send_transcript", "message": alt.transcript}, ) class ChatConsumer(AsyncWebsocketConsumer): async def connect(self): await self.accept() await self.send( text_data=json.dumps( {"type": "connection_established", "channel_name": self.channel_name} ) ) self.stream = await stream_client.start_stream_transcription( language_code="en-US", media_sample_rate_hz=48000, media_encoding="ogg-opus", # 可选:开启部分结果稳定返回,降低转写延迟 enable_partial_results_stabilization=True ) # 初始化handler时传入当前连接的真实channel_name self.handler = AWSTranscriptHandler( self.stream.output_stream, get_channel_layer(), self.channel_name ) # 启动后台监听任务,绑定到实例属性避免被回收 self.listen_task = asyncio.create_task(self.handler.handle_events()) async def disconnect(self, close_code): # 关闭流、取消监听任务 await self.stream.input_stream.end_stream() if hasattr(self, "listen_task") and not self.listen_task.done(): self.listen_task.cancel() try: await self.listen_task except asyncio.CancelledError: pass async def receive(self, text_data=None, bytes_data=None): if bytes_data: # 只负责发送音频,不要在这里调用handle_events await self.stream.input_stream.send_audio_event(audio_chunk=bytes_data) async def send_transcript(self, event): await self.send( text_data=json.dumps({"type": "transcript", "message": event["message"]}) )
额外注意事项
- 如果需要更低的延迟,可以把
enable_partial_results_stabilization参数设为True,同时去掉if not (result.is_partial)的判断,就能拿到用户说话过程中实时更新的中间转写文本。 - 前端发送音频块时不要攒太大的blob,建议每100-200ms发送一次音频帧,延迟过高会导致Transcribe返回结果的实时性下降。
- 如果持续出现音频解析错误,可以先把收到的音频块存成本地文件,用ffprobe校验编码格式、采样率是否和接口声明的参数一致。
内容的提问来源于stack exchange,提问作者vettipayyan
相关产品推荐
相关产品推荐

