基于Golang与AWS实现视频实时转码HLS并流式上传至S3
实时流式转码视频为HLS并上传至AWS S3的实现方案
现有代码的核心问题
- 用
io.ReadAll把整个视频下载到内存后才处理,完全违背了实时流式处理的要求 - 转码结果写入本地文件,未对接S3的流式上传能力
- FFmpeg命令参数不符合HLS格式要求
解决方案核心思路
实现下载→转码→上传的全流式流水线:
- 视频流从HTTP响应直接流入FFmpeg的标准输入,无需缓存完整文件
- FFmpeg实时转码为HLS格式,输出分片和索引文件
- 实时监听转码输出的文件变化,将新生成的分片和更新的m3u8索引立即上传至S3
修正后的实现代码
import ( "context" "io" "net/http" "os" "os/exec" "path/filepath" "github.com/VinukaThejana/go-utils/logger" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/config" "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/fsnotify/fsnotify" ) func main() { // 初始化AWS S3客户端 cfg, err := config.LoadDefaultConfig(context.TODO(), config.WithRegion("us-east-1")) if err != nil { logger.Fatalf("加载AWS配置失败: %v", err) } s3Client := s3.NewFromConfig(cfg) bucketName := "your-hls-bucket" s3Prefix := "stream/" // S3上存放HLS文件的前缀 // 发起视频流请求,保持响应Body为流式状态 res, err := http.Get("http://localhost/video.mp4") if err != nil { logger.Fatalf("获取视频流失败: %v", err) } defer res.Body.Close() // 创建临时目录存储转码后的HLS文件 tempDir, err := os.MkdirTemp("", "hls-transcode-*") if err != nil { logger.Fatalf("创建临时目录失败: %v", err) } defer os.RemoveAll(tempDir) // 程序结束后清理临时文件 // 配置FFmpeg实时转码为HLS的命令 // 关键参数:ultrafast预设、zerolatency调优保证低延迟,10秒分片,保留所有分片索引 cmd := exec.Command("ffmpeg", "-i", "pipe:0", // 从标准输入读取视频流 "-c:v", "libx264", "-preset", "ultrafast", "-tune", "zerolatency", "-c:a", "aac", "-hls_time", "10", // 每10秒生成一个TS分片 "-hls_list_size", "0", // 在m3u8中保留所有分片记录 "-hls_segment_filename", tempDir+"/segment_%03d.ts", // 分片命名模板 "-f", "hls", tempDir+"/index.m3u8", // 输出的m3u8索引文件路径 ) cmd.Stderr = os.Stderr // 重定向FFmpeg日志到标准错误输出 // 建立FFmpeg的标准输入管道 stdin, err := cmd.StdinPipe() if err != nil { logger.Fatalf("创建FFmpeg输入管道失败: %v", err) } defer stdin.Close() // 启动FFmpeg转码进程 if err := cmd.Start(); err != nil { logger.Fatalf("启动FFmpeg失败: %v", err) } // 后台协程:将视频流持续写入FFmpeg输入管道 go func() { _, err := io.Copy(stdin, res.Body) if err != nil && err != io.EOF { logger.Errorf("写入FFmpeg输入失败: %v", err) } stdin.Close() // 关闭输入,触发FFmpeg结束转码 }() // 初始化文件监听器,监听临时目录的文件变化 watcher, err := fsnotify.NewWatcher() if err != nil { logger.Fatalf("创建文件监听器失败: %v", err) } defer watcher.Close() if err := watcher.Add(tempDir); err != nil { logger.Fatalf("添加目录监听失败: %v", err) } // 后台协程:监听文件变化,实时上传到S3 go func() { uploadedSegments := make(map[string]bool) // 记录已上传的分片,避免重复上传 for { select { case event, ok := <-watcher.Events: if !ok { return } // 只处理文件创建和写入完成的事件 if (event.Has(fsnotify.Create) || event.Has(fsnotify.Write)) && !event.Has(fsnotify.IsDir) { fileName := filepath.Base(event.Name) // 分片文件只上传一次,m3u8每次更新都重新上传 if fileName != "index.m3u8" && uploadedSegments[fileName] { continue } // 打开文件准备上传 file, err := os.Open(event.Name) if err != nil { logger.Errorf("打开转码文件失败: %v", err) continue } // 上传到S3 objectKey := s3Prefix + fileName _, err = s3Client.PutObject(context.TODO(), &s3.PutObjectInput{ Bucket: aws.String(bucketName), Key: aws.String(objectKey), Body: file, }) file.Close() if err != nil { logger.Errorf("上传文件%s到S3失败: %v", objectKey, err) } else { logger.Infof("成功上传文件到S3: %s", objectKey) if fileName != "index.m3u8" { uploadedSegments[fileName] = true } } } case err, ok := <-watcher.Errors: if !ok { return } logger.Errorf("文件监听错误: %v", err) } } }() // 等待FFmpeg转码进程结束 if err := cmd.Wait(); err != nil { logger.Errorf("FFmpeg转码进程异常退出: %v", err) } logger.Infof("视频流转码并上传S3完成") }
关键细节说明
- 流式传输:通过
io.Copy将HTTP响应Body直接写入FFmpeg输入,全程无完整文件缓存,实现实时处理 - FFmpeg参数:
-preset ultrafast和-tune zerolatency是实时转码的核心,大幅降低转码延迟;-hls_time控制分片大小,平衡延迟和请求次数 - S3上传逻辑:使用
fsnotify监听临时目录,新分片生成后立即上传;m3u8索引文件会随新分片不断更新,所以每次变化都重新上传,保证客户端能获取最新的播放列表 - 资源清理:临时目录通过
defer os.RemoveAll自动清理,避免磁盘资源泄漏
内容的提问来源于stack exchange,提问作者VinukaKodituwakku
相关产品推荐
相关产品推荐

