You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Amazon Transcribe Python API事件处理器仅在流结束后处理音频

问题根因

流式转写不生效是四个核心用法错误导致的:

  • handle_events() 是持续运行的异步事件监听循环,你把它放在receive回调里每次收到音频块才调用一次,每次仅执行单轮事件检查就退出,只有流关闭触发终态事件时才能拿到结果,无法实时处理中间返回的转写片段。
  • 之前用asyncio.create_task创建监听任务立刻结束,是因为Django Channels的异步消费者会自动回收未绑定到实例生命周期的孤儿协程,任务还没跑起来就被取消了。
  • 为每个音频块创建独立转写流报内部错误属于预期行为:Amazon Transcribe流式接口要求单条连接持续按序发送完整编码序列的音频,碎片化的单块音频缺少编码头、帧序列不连续,服务端无法正常解析。
  • 额外逻辑bug:转写回调里硬编码了channel_name = "test",没有传入当前WebSocket连接的实际channel名称,就算拿到转写结果也无法正确推送到当前前端。
修复方案
  1. 在连接建立、启动Transcribe流之后,立刻在后台启动事件监听协程,把任务绑定到consumer实例属性上避免被回收。
  2. 不要在receive方法里重复调用handle_events(),这个方法本身是无限循环,只要启动一次就会持续监听转写返回。
  3. 断开连接时主动取消监听任务、正确关闭Transcribe流,避免资源泄漏。
  4. 修正参数校验:确认前端发送的音频编码、采样率和调用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 12:51:31