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

基于Golang与AWS实现视频实时转码HLS并流式上传至S3

实时流式转码视频为HLS并上传至AWS S3的实现方案

现有代码的核心问题

  • 用io.ReadAll把整个视频下载到内存后才处理,完全违背了实时流式处理的要求
  • 转码结果写入本地文件,未对接S3的流式上传能力
  • FFmpeg命令参数不符合HLS格式要求

解决方案核心思路

实现下载→转码→上传的全流式流水线:

  1. 视频流从HTTP响应直接流入FFmpeg的标准输入,无需缓存完整文件
  2. FFmpeg实时转码为HLS格式,输出分片和索引文件
  3. 实时监听转码输出的文件变化,将新生成的分片和更新的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 15:03:11