基于RxJS实现Amazon Transcribe无阻塞轮询方案咨询
基于RxJS实现Amazon Transcribe无阻塞轮询方案
原用setInterval实现的轮询有两大硬伤:一是不管前一次请求是否返回,到点就发新请求,容易造成请求堆积;二是异步回调嵌套在定时器里,会阻塞事件循环,并发场景下根本扛不住。用RxJS可以完美解决这些问题,以下是具体实现:
完整替换代码
import { interval, from } from 'rxjs'; import { concatMap, filter, take, catchError } from 'rxjs/operators'; // 启动转录任务的代码保持不变 const jobId = nanoid(); await amazonTrascribeClient .startTranscriptionJob({ IdentifyMultipleLanguages: true, TranscriptionJobName: jobId, Media: { MediaFileUri: "s3://file-location", }, Subtitles: { OutputStartIndex: 1, Formats: ["vtt", "srt"], }, OutputBucketName: `file-location`, OutputKey: `transcriptions/${jobId}/`, }) .promise(); // RxJS轮询实现 const pollTranscriptionJob = (jobId) => { // 每2秒发起一次查询请求 return interval(2000).pipe( // concatMap确保前一次请求完成后再发下一次,避免请求堆积 concatMap(() => from(amazonTrascribeClient.getTranscriptionJob({ TranscriptionJobName: jobId }).promise())), // 过滤出任务完成或失败的状态 filter(response => ["COMPLETED", "FAILED"].includes(response.TranscriptionJob.TranscriptionJobStatus)), // 拿到结果后立即停止轮询 take(1), // 处理查询过程中的错误 catchError(error => { console.error('轮询任务状态失败:', error); throw error; }) ); }; // 订阅轮询流 pollTranscriptionJob(jobId).subscribe({ next: (response) => { const { TranscriptionJob } = response; if (TranscriptionJob.TranscriptionJobStatus === 'COMPLETED') { // 任务完成:写入数据库、处理转录结果等操作 console.log('转录任务完成:', TranscriptionJob.Transcript.TranscriptFileUri); } else { // 任务失败:处理失败逻辑 console.error('转录任务失败:', TranscriptionJob.FailureReason); } }, error: (error) => { // 处理全局错误 console.error('轮询流程出错:', error); } });
关键优势说明
- 无阻塞:RxJS的操作符基于异步流实现,不会阻塞事件循环,并发场景下可同时处理多个转录任务的轮询。
- 避免请求堆积:用
concatMap替代定时器的强制触发,确保只有前一次查询请求完成后,才会发起下一次请求,不会额外增加服务器压力。 - 自动终止:通过
take(1)在拿到最终状态后自动停止轮询,无需手动清理定时器。 - 统一错误处理:内置的
catchError可以集中处理轮询过程中的网络错误或API异常,避免遗漏错误场景。
内容的提问来源于stack exchange,提问作者Mirusky
相关产品推荐
相关产品推荐

