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

在Python asyncio中保存Deepgram流式识别回调的转录内容

解决Deepgram流式转录文本的收集与格式化问题

要将Deepgram流式识别的转录块收集起来并进行格式化,核心是让回调函数能安全地与主协程共享数据,以下是两种可行的实现方案:

方案一:用类维护转录状态

通过自定义类来保存转录片段和完整文本,类的实例可以被回调函数直接访问,无需复杂的同步操作(适用于单线程asyncio场景)。

修改后的完整代码:

from deepgram import Deepgram
import asyncio
import aiohttp

DEEPGRAM_API_KEY = '****'
URL = 'http://stream.live.vc.bbcmedia.co.uk/bbc_world_service'

class TranscriptCollector:
    def __init__(self):
        self.transcript_chunks = []  # 存储单个转录块
        self.full_transcript = ""    # 拼接后的完整文本

    def add_chunk(self, json_data):
        chunk = json_data['channel']['alternatives'][0]['transcript']
        if chunk:  # 过滤空转录内容
            self.transcript_chunks.append(chunk)
            self.full_transcript += chunk + " "

    def format_transcript(self):
        # 自定义格式化逻辑,比如去除多余空格、分段等
        return self.full_transcript.strip()

async def main():
    collector = TranscriptCollector()
    deepgram = Deepgram(DEEPGRAM_API_KEY)

    # 建立Deepgram流式连接
    deepgramLive = await deepgram.transcription.live({ 'language': 'en-US' })

    # 注册连接关闭回调
    deepgramLive.registerHandler(deepgramLive.event.CLOSE, lambda c: print(f'连接关闭,状态码:{c}.'))

    # 注册转录接收回调,将块添加到收集器
    deepgramLive.registerHandler(deepgramLive.event.TRANSCRIPT_RECEIVED, collector.add_chunk)

    # 读取音频流并发送给Deepgram
    async with aiohttp.ClientSession() as session:
        async with session.get(URL) as audio:
            while True:
                data = await audio.content.readany()
                deepgramLive.send(data)

                # 可在此处定期检查并格式化输出,比如每收集5个块
                if len(collector.transcript_chunks) >= 5:
                    formatted_text = collector.format_transcript()
                    print("当前格式化转录内容:")
                    print(formatted_text)
                    # 可选:清空收集器,继续收集新内容
                    # collector.transcript_chunks.clear()
                    # collector.full_transcript = ""

                if not data:
                    break

    await deepgramLive.finish()

    # 处理最终完整转录内容
    final_formatted = collector.format_transcript()
    print("\n最终完整转录:")
    print(final_formatted)

asyncio.run(main())

方案二:用asyncio.Queue异步处理转录块

如果需要在异步任务中处理转录内容(如写入文件、调用异步API),可以使用asyncio.Queue。注意回调是同步函数,需通过事件循环的线程安全方法操作队列。

完整代码示例:

from deepgram import Deepgram
import asyncio
import aiohttp

DEEPGRAM_API_KEY = '****'
URL = 'http://stream.live.vc.bbcmedia.co.uk/bbc_world_service'

async def transcript_processor(queue):
    # 异步任务:处理队列中的转录块
    full_transcript = ""
    while True:
        chunk = await queue.get()
        if chunk is None:  # 接收结束信号,终止任务
            break
        if chunk:
            full_transcript += chunk + " "
            # 每累积一定长度就格式化输出
            if len(full_transcript) > 200:
                print("当前格式化转录内容:")
                print(full_transcript.strip())
                full_transcript = ""
    # 处理剩余未输出的内容
    if full_transcript:
        print("\n最终剩余转录内容:")
        print(full_transcript.strip())

def queue_callback(queue, json_data):
    # 同步回调:安全地将转录块加入队列
    chunk = json_data['channel']['alternatives'][0]['transcript']
    asyncio.get_event_loop().call_soon_threadsafe(queue.put_nowait, chunk)

async def main():
    queue = asyncio.Queue()
    # 启动转录处理异步任务
    processor_task = asyncio.create_task(transcript_processor(queue))

    deepgram = Deepgram(DEEPGRAM_API_KEY)
    deepgramLive = await deepgram.transcription.live({ 'language': 'en-US' })

    deepgramLive.registerHandler(deepgramLive.event.CLOSE, lambda c: print(f'连接关闭,状态码:{c}.'))

    # 注册转录接收回调,绑定队列
    deepgramLive.registerHandler(deepgramLive.event.TRANSCRIPT_RECEIVED, lambda data: queue_callback(queue, data))

    # 发送音频流给Deepgram
    async with aiohttp.ClientSession() as session:
        async with session.get(URL) as audio:
            while True:
                data = await audio.content.readany()
                deepgramLive.send(data)
                if not data:
                    break

    await deepgramLive.finish()
    # 发送结束信号给处理器任务
    await queue.put(None)
    # 等待处理器任务完成
    await processor_task

asyncio.run(main())

关于之前遇到的AttributeError

你之前遇到的错误大概率是因为使用异步函数作为Deepgram的回调(Deepgram要求回调为同步函数),或者回调中访问了未正确初始化的变量。上述两种方案均使用同步回调,避免了该问题。

内容的提问来源于stack exchange,提问作者nineplymaple

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 08:45:08