在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
相关产品推荐
相关产品推荐

