服务端转码不求人:写 FFmpeg 分布式转码 + GOP 切片 + 缩略图

客户端编码好了、推流也通了,但服务端还得把用户上传的视频转成多码率、生成缩略图、GOP 对齐切片,这些「后端活」音视频开发也得会。本文用 Claude Code 写一套生产级转码脚本,从单机到分布式。

1、服务端转码:不只是 「ffmpeg -i input output」

“后端说转码就是跑个 FFmpeg 命令,给你写好了就上线了…”

然后:

  • 📹 1080p 视频转码 3 分钟 → 用户上传后等半天才开始转 → 投诉「视频发不出去」
  • 💣 OOM → 一个 4K 视频把 2GB 内存吃满 → ffmpeg 被 kill → 视频卡在转码中
  • 🔄 GOP 没对齐 → HLS 切片后每片时长不均匀 → 播放器频繁 rebuffer
  • ⏱️ 无并发控制 → 10 个视频同时涌入 → CPU 100% → 所有视频都慢

Claude Code 能做什么? 不只是帮你生成命令——而是帮你设计完整的转码管线:并发控制、内存限制、GOP 对齐、进度追踪、错误重试、多机分布

2、转码管线架构

用户上传视频
     │
     ▼
┌─────────────────┐
│ 消息队列 (Redis)  │  异步解耦,削峰填谷
└────────┬────────┘
     │
     ▼
┌─────────────────┐
│ 调度器 (Go/Node) │  分配任务 + 负载均衡 + 心跳
└────────┬────────┘
     │
     ▼
┌────────────────────────────────────────┐
│ Worker (FFmpeg)                        │
│                                        │
│ 1. 视频分析 (ffprobe) → 元数据          │
│ 2. 多码率转码 → 1080p/720p/480p/360p   │
│ 3. GOP 对齐切片 → HLS (.m3u8 + .ts)     │
│ 4. 缩略图生成 → 雪碧图 (sprite)         │
│ 5. 结果回调 → HTTP / 消息队列           │
└────────────────────────────────────────┘
     │
     ▼
  CDN / OSS / 存储

3、Claude Code 生成转码脚本

3.1、 Prompt

帮我写一套服务端视频转码管线。

要求:
1. 单 Worker 脚本 (Python/Shell):接收一个视频文件,输出多码率 HLS + 缩略图
2. 多码率转码:1080p/720p/480p/360p,H.264 + H.265 各一份
3. GOP 对齐:所有码率的 I 帧在相同位置(保证 HLS 无缝切换)
4. 缩略图雪碧图:每 10 秒一帧,拼成 N×N 网格 + VTT 坐标文件
5. 并发控制:限制同时转码数,超出排队
6. 资源限制:每个 ffmpeg 进程限制 CPU 核心数和内存
7. 错误处理:转码超时 kill、重试机制、错误日志
8. 进度追踪:ffmpeg 输出解析,实时更新转码百分比
9. 中文注释

3.2、 核心转码 Worker

#!/usr/bin/env python3
"""
video_transcoder.py
生产级视频转码 Worker:多码率 HLS + GOP 对齐 + 缩略图
"""

import subprocess
import json
import os
import time
import re
import shutil
from pathlib import Path
from concurrent.futures import ThreadPoolExecutor
from dataclasses import dataclass
from typing import List, Optional
import signal
import resource

# ═══════════════════════════════════════════
# MARK: - 配置
# ═══════════════════════════════════════════

@dataclass
class TranscodeConfig:
    # 输入
    input_path: str
    output_dir: str

    # 多码率档位
    variants: List[dict] = None

    # HLS
    hls_segment_time: int = 4       # 每片 4 秒
    hls_playlist_type: str = "vod"   # VOD 点播

    # 缩略图
    thumbnail_interval: int = 10     # 每 10 秒一帧
    thumbnail_columns: int = 5       # 雪碧图 5 列
    thumbnail_rows: int = 4          # 雪碧图 4 行

    # 资源控制
    max_cpu_cores: int = 4
    max_memory_mb: int = 2048
    timeout_seconds: int = 1800      # 单个视频最多转 30 分钟

    # 重试
    max_retries: int = 3

    def __post_init__(self):
        if self.variants is None:
            self.variants = [
                {"name": "1080p", "width": 1920, "height": 1080, "bitrate": "5000k",
                 "maxrate": "5500k", "bufsize": "8000k"},
                {"name": "720p",  "width": 1280, "height": 720,  "bitrate": "2800k",
                 "maxrate": "3000k", "bufsize": "4000k"},
                {"name": "480p",  "width": 854,  "height": 480,  "bitrate": "1400k",
                 "maxrate": "1500k", "bufsize": "2000k"},
                {"name": "360p",  "width": 640,  "height": 360,  "bitrate": "800k",
                 "maxrate": "850k",  "bufsize": "1200k"},
            ]

# ═══════════════════════════════════════════
# MARK: - 主转码类
# ═══════════════════════════════════════════

class VideoTranscoder:
    def __init__(self, config: TranscodeConfig):
        self.config = config
        self.progress = 0.0
        self.status = "pending"
        self.error_message = None

    def run(self):
        """执行完整转码管线"""
        try:
            # Step 0: 准备
            os.makedirs(self.config.output_dir, exist_ok=True)

            # Step 1: 视频分析(ffprobe)
            metadata = self._analyze_video()
            print(f"[Transcoder] 源视频: {metadata['width']}x{metadata['height']}, "
                  f"{metadata['duration']:.1f}s, {metadata['bitrate']}")

            # Step 2: GOP 对齐多码率转码 (并行)
            self.status = "transcoding"
            self._transcode_variants(metadata)

            # Step 3: 生成 Master Playlist
            self._generate_master_playlist()

            # Step 4: 生成缩略图雪碧图
            self._generate_thumbnails(metadata)

            self.status = "completed"
            self.progress = 100.0
            print(f"[Transcoder] ✅ 转码完成 → {self.config.output_dir}")

        except Exception as e:
            self.status = "failed"
            self.error_message = str(e)
            print(f"[Transcoder] ❌ 转码失败: {e}")
            raise

    # ═══════════════════════════════════════
    # Step 1: 视频分析
    # ═══════════════════════════════════════

    def _analyze_video(self) -> dict:
        """用 ffprobe 提取视频元数据"""
        cmd = [
            "ffprobe",
            "-v", "quiet",
            "-print_format", "json",
            "-show_format",
            "-show_streams",
            self.config.input_path
        ]

        try:
            result = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
            data = json.loads(result.stdout)

            # 找视频流
            video_stream = None
            audio_stream = None
            for stream in data.get("streams", []):
                if stream["codec_type"] == "video":
                    video_stream = stream
                elif stream["codec_type"] == "audio":
                    audio_stream = stream

            if not video_stream:
                raise ValueError("源文件无视频流")

            return {
                "width": video_stream.get("width", 0),
                "height": video_stream.get("height", 0),
                "codec": video_stream.get("codec_name", "unknown"),
                "fps": eval(video_stream.get("r_frame_rate", "30/1")),
                "duration": float(data.get("format", {}).get("duration", 0)),
                "bitrate": data.get("format", {}).get("bit_rate", "0"),
                "has_audio": audio_stream is not None,
            }

        except subprocess.TimeoutExpired:
            raise RuntimeError("ffprobe 超时")

    # ═══════════════════════════════════════
    # Step 2: GOP 对齐多码率转码
    # ═══════════════════════════════════════

    def _transcode_variants(self, metadata: dict):
        """并行转码各码率档位"""
        # GOP 对齐策略:
        # 所有码率使用相同的 I 帧间隔 (keyint) 和场景切换检测 (sc_threshold)
        # HLS 切片必须与 GOP 边界对齐: hls_time == GOP 时长

        gop_duration = self.config.hls_segment_time  # 4 秒
        gop_frames = int(gop_duration * metadata["fps"])

        # 限制并发数(避免内存爆炸)
        max_concurrent = min(len(self.config.variants), self.config.max_cpu_cores)

        with ThreadPoolExecutor(max_workers=max_concurrent) as executor:
            futures = []
            for variant in self.config.variants:
                future = executor.submit(
                    self._transcode_single_variant,
                    metadata, variant, gop_frames
                )
                futures.append(future)

            # 等待全部完成
            for i, future in enumerate(futures):
                try:
                    future.result(timeout=self.config.timeout_seconds)
                    self.progress = (i + 1) / len(self.config.variants) * 100
                except Exception as e:
                    print(f"[Transcoder] ⚠️ 档位 {self.config.variants[i]['name']} 失败: {e}")
                    # 降级:缺少一个档位不影响整体
                    # (HLS Master Playlist 会自动排除失败的档位)

    def _transcode_single_variant(self, metadata: dict, variant: dict, gop_frames: int):
        """转码单个码率档位到 HLS"""
        output_dir = os.path.join(self.config.output_dir, variant["name"])
        os.makedirs(output_dir, exist_ok=True)
        output_playlist = os.path.join(output_dir, "index.m3u8")

        # 构建 ffmpeg 命令
        cmd = [
            "ffmpeg",
            "-y",
            "-i", self.config.input_path,

            # GOP 对齐参数(所有码率相同!)
            "-force_key_frames", f"expr:gte(n,n_forced*{gop_frames})",
            "-sc_threshold", "0",           # 禁用场景切换自动 I 帧
            "-g", str(gop_frames),           # GOP 大小
            "-keyint_min", str(gop_frames), # 最小 GOP 大小

            # 视频编码
            "-c:v", "libx264",
            "-preset", "medium",            # 平衡速度与压缩率
            "-profile:v", "main",
            "-b:v", variant["bitrate"],
            "-maxrate", variant["maxrate"],
            "-bufsize", variant["bufsize"],
            "-vf", f"scale={variant['width']}:{variant['height']}:force_original_aspect_ratio=decrease,pad={variant['width']}:{variant['height']}:(ow-iw)/2:(oh-ih)/2",

            # 音频编码
            "-c:a", "aac",
            "-b:a", "128k",
            "-ar", "48000",
            "-ac", "2",

            # HLS 输出
            "-f", "hls",
            "-hls_time", str(self.config.hls_segment_time),
            "-hls_playlist_type", self.config.hls_playlist_type,
            "-hls_list_size", "0",               # VOD = 列出所有分片
            "-hls_segment_filename",
            os.path.join(output_dir, "segment_%05d.ts"),
            output_playlist
        ]

        # 资源限制(ulimit)
        self._set_resource_limits()

        # 执行
        print(f"[{variant['name']}] 开始转码...")
        process = subprocess.Popen(
            cmd,
            stdout=subprocess.PIPE,
            stderr=subprocess.PIPE,
            universal_newlines=True
        )

        # 监控进度
        try:
            self._monitor_progress(process, variant["name"])
            process.wait(timeout=self.config.timeout_seconds)

            if process.returncode != 0:
                stderr = process.stderr.read()
                raise RuntimeError(f"ffmpeg 返回 {process.returncode}: {stderr[-500:]}")

            print(f"[{variant['name']}] ✅ 转码完成")

        except subprocess.TimeoutExpired:
            process.kill()
            raise RuntimeError(f"[{variant['name']}] 转码超时 ({self.config.timeout_seconds}s)")

    # ═══════════════════════════════════════
    # Step 3: Master Playlist
    # ═══════════════════════════════════════

    def _generate_master_playlist(self):
        """生成 HLS Master Playlist (多码率自适应)"""
        master_path = os.path.join(self.config.output_dir, "master.m3u8")

        lines = ["#EXTM3U", "#EXT-X-VERSION:3"]

        for variant in self.config.variants:
            variant_dir = os.path.join(self.config.output_dir, variant["name"])
            playlist_path = os.path.join(variant_dir, "index.m3u8")

            if not os.path.exists(playlist_path):
                continue  # 跳过失败的档位

            # 计算总带宽(视频 + 音频)
            video_br = int(re.sub(r'[kK]', '000', variant["bitrate"]))
            total_br = video_br + 128000

            lines.append(
                f'#EXT-X-STREAM-INF:'
                f'BANDWIDTH={total_br},'
                f'RESOLUTION={variant["width"]}x{variant["height"]},'
                f'CODECS="avc1.64001f,mp4a.40.2"'
            )
            lines.append(f'{variant["name"]}/index.m3u8')

        with open(master_path, 'w') as f:
            f.write('\n'.join(lines) + '\n')

        print(f"[Transcoder] ✅ Master Playlist → {master_path}")

    # ═══════════════════════════════════════
    # Step 4: 缩略图雪碧图
    # ═══════════════════════════════════════

    def _generate_thumbnails(self, metadata: dict):
        """生成缩略图雪碧图 (sprite.jpg) + VTT 坐标文件"""
        cols = self.config.thumbnail_columns
        rows = self.config.thumbnail_rows
        frames_per_sprite = cols * rows
        interval = self.config.thumbnail_interval

        sprite_dir = os.path.join(self.config.output_dir, "thumbnails")
        os.makedirs(sprite_dir, exist_ok=True)

        thumb_width = 160   # 每个缩略图宽
        thumb_height = 90   # 每个缩略图高

        # 1. 提取缩略图帧
        temp_dir = os.path.join(sprite_dir, "temp")
        os.makedirs(temp_dir, exist_ok=True)

        cmd_extract = [
            "ffmpeg", "-y",
            "-i", self.config.input_path,
            "-vf", f"fps=1/{interval},scale={thumb_width}:{thumb_height}",
            os.path.join(temp_dir, "thumb_%04d.jpg")
        ]
        subprocess.run(cmd_extract, check=True)

        # 2. 拼成雪碧图
        # 每 frames_per_sprite 帧拼一张大图
        # (简化版本:用 ImageMagick montage,生产环境可用 GPU 加速)
        thumb_files = sorted(Path(temp_dir).glob("thumb_*.jpg"))
        sprite_index = 0
        vtt_entries = []

        for i in range(0, len(thumb_files), frames_per_sprite):
            batch = thumb_files[i:i + frames_per_sprite]
            sprite_path = os.path.join(sprite_dir, f"sprite_{sprite_index:03d}.jpg")

            # montage: 拼成 grid
            subprocess.run([
                "montage",
                *[str(f) for f in batch],
                "-tile", f"{cols}x{rows}",
                "-geometry", f"{thumb_width}x{thumb_height}+0+0",
                sprite_path
            ], check=True)

            # 生成 VTT
            for j, thumb_file in enumerate(batch):
                time_sec = (i + j) * interval
                col = j % cols
                row = j // cols
                x = col * thumb_width
                y = row * thumb_height

                vtt_entries.append(
                    f"{self._format_vtt_time(time_sec)} → "
                    f"sprite_{sprite_index:03d}.jpg#"
                    f"xywh={x},{y},{thumb_width},{thumb_height}"
                )

            sprite_index += 1

        # 写 VTT 文件
        vtt_path = os.path.join(sprite_dir, "thumbnails.vtt")
        with open(vtt_path, 'w') as f:
            f.write("WEBVTT\n\n")
            for entry in vtt_entries:
                f.write(f"{entry}\n\n")

        # 清理临时帧
        shutil.rmtree(temp_dir)
        print(f"[Transcoder] ✅ 雪碧图: {sprite_index} 张, VTT → {vtt_path}")

    # ═══════════════════════════════════════
    # 工具方法
    # ═══════════════════════════════════════

    def _monitor_progress(self, process, variant_name: str):
        """解析 ffmpeg stderr 中的 time= 字段来追踪进度"""
        for line in process.stderr:
            match = re.search(r"time=(\d+):(\d+):(\d+\.\d+)", line)
            if match:
                hours = int(match.group(1))
                minutes = int(match.group(2))
                seconds = float(match.group(3))
                elapsed_sec = hours * 3600 + minutes * 60 + seconds
                print(f"  [{variant_name}] 进度: {elapsed_sec:.1f}s")

    def _set_resource_limits(self):
        """限制 ffmpeg 子进程的资源"""
        # CPU: 通过 ffmpeg -threads 参数控制,这里设置辅助限制
        resource.setrlimit(
            resource.RLIMIT_CPU,
            (self.config.timeout_seconds, self.config.timeout_seconds)
        )
        # 内存
        resource.setrlimit(
            resource.RLIMIT_AS,
            (self.config.max_memory_mb * 1024 * 1024,
             self.config.max_memory_mb * 1024 * 1024)
        )

    @staticmethod
    def _format_vtt_time(seconds: int) -> str:
        h = seconds // 3600
        m = (seconds % 3600) // 60
        s = seconds % 60
        return f"{h:02d}:{m:02d}:{s:02d}.000"


# ═══════════════════════════════════════════
# MARK: - 分布式调度器(简化版)
# ═══════════════════════════════════════════

class TranscodeScheduler:
    """简单的并发控制 + 任务队列"""

    def __init__(self, max_concurrent: int = 2):
        self.max_concurrent = max_concurrent
        self.executor = ThreadPoolExecutor(max_workers=max_concurrent)
        self.active_tasks = {}

    def submit(self, task_id: str, config: TranscodeConfig) -> bool:
        """提交转码任务,如果队列满了返回 False"""
        if len(self.active_tasks) >= self.max_concurrent:
            return False  # 队列满,需要等待

        transcoder = VideoTranscoder(config)
        future = self.executor.submit(self._run_task, task_id, transcoder)
        self.active_tasks[task_id] = {
            "transcoder": transcoder,
            "future": future,
            "status": "running"
        }
        return True

    def _run_task(self, task_id: str, transcoder: VideoTranscoder):
        """在 worker 线程中执行转码"""
        try:
            transcoder.run()
            self.active_tasks[task_id]["status"] = "completed"
        except Exception as e:
            self.active_tasks[task_id]["status"] = "failed"
            print(f"[Scheduler] ❌ 任务 {task_id} 失败: {e}")

    def get_progress(self, task_id: str) -> float:
        if task_id in self.active_tasks:
            return self.active_tasks[task_id]["transcoder"].progress
        return 0.0

    def shutdown(self):
        self.executor.shutdown(wait=True)

# ═══════════════════════════════════════════
# MARK: - 入口
# ═══════════════════════════════════════════

if __name__ == "__main__":
    import argparse

    parser = argparse.ArgumentParser(description="视频转码 Worker")
    parser.add_argument("input", help="输入视频路径")
    parser.add_argument("-o", "--output-dir", default="./output", help="输出目录")
    parser.add_argument("--h264-only", action="store_true", help="仅 H.264")
    parser.add_argument("--cpu", type=int, default=4, help="最大 CPU 核数")
    parser.add_argument("--memory", type=int, default=2048, help="最大内存 (MB)")

    args = parser.parse_args()

    config = TranscodeConfig(
        input_path=args.input,
        output_dir=args.output_dir,
        max_cpu_cores=args.cpu,
        max_memory_mb=args.memory
    )

    transcoder = VideoTranscoder(config)
    transcoder.run()

4、单机 → 分布式改造

4.1、Claude Code Prompt

上面的脚本是单机版的,帮我改造为分布式版本。

要求:
1. 消息队列:Redis List / RabbitMQ 接收转码任务
2. 多 Worker:每个 Worker 独立拉取任务
3. 心跳上报:Worker 每 10 秒上报状态(处理中/空闲/挂了)
4. 结果回调:转码完成后 HTTP POST 结果 URL 给业务方
5. Worker 挂掉自动重分配:心跳超时 30s → 任务重新入队

4.2、分布式 Worker

# distributed_worker.py (关键部分)

import redis
import requests
import threading
import time

class DistributedTranscodeWorker:
    def __init__(self, redis_url: str, worker_id: str, callback_url: str):
        self.redis = redis.from_url(redis_url)
        self.worker_id = worker_id
        self.callback_url = callback_url
        self.is_running = True

    def start(self):
        """主循环:拉任务 → 转码 → 回调"""
        # 启动心跳线程
        heartbeat_thread = threading.Thread(target=self._send_heartbeat)
        heartbeat_thread.daemon = True
        heartbeat_thread.start()

        while self.is_running:
            # 从队列拉取任务 (BLPOP 阻塞等待)
            result = self.redis.blpop("transcode:queue", timeout=5)
            if result is None:
                continue

            _, task_json = result
            task = json.loads(task_json)

            try:
                # 执行转码
                config = TranscodeConfig(
                    input_path=task["input_url"],
                    output_dir=task["output_dir"]
                )
                transcoder = VideoTranscoder(config)
                transcoder.run()

                # 回调业务方
                self._callback_success(task["task_id"], {
                    "master_url": f"{task['output_dir']}/master.m3u8",
                    "duration": transcoder.progress,
                })

            except Exception as e:
                self._callback_failure(task["task_id"], str(e))

                # 重试逻辑
                retries = task.get("retries", 0)
                if retries < 3:
                    task["retries"] = retries + 1
                    self.redis.rpush("transcode:queue", json.dumps(task))

    def _send_heartbeat(self):
        """每 10 秒发送心跳"""
        while self.is_running:
            self.redis.setex(
                f"worker:heartbeat:{self.worker_id}",
                30,  # TTL = 30s
                json.dumps({"timestamp": time.time(), "status": "alive"})
            )
            time.sleep(10)

    def _callback_success(self, task_id: str, result: dict):
        requests.post(f"{self.callback_url}/transcode/{task_id}/success", json=result)

    def _callback_failure(self, task_id: str, error: str):
        requests.post(f"{self.callback_url}/transcode/{task_id}/failed",
                      json={"error": error})

5、转码优化清单

优化项效果方法
-preset faster 替代 medium速度 +40%,画质略降短视频/实时场景用 faster,点播用 medium
跳过无用的音频流省 5-10% CPU-map 0:v -map 0:a:0 只取第一个音频流
HLS 分片复用同 GOP 的 ts 可直接复用不同码率的同 GOP 分片 md5 相同 → 硬链替代拷⻉
GPU 加速 (NVENC / VideoToolbox)3-5x 速度画质略低于 x264 软编,需评估
预分析 + 分割4K 视频切成 4 段并行转切分点必须在 GOP 边界

学习和提升音视频开发技术,欢迎你加入我们的知识星球

服务端转码不求人:写 FFmpeg 分布式转码 + GOP 切片 + 缩略图

版权声明:本文内容转自互联网,本文观点仅代表作者本人。本站仅提供信息存储空间服务,所有权归原作者所有。如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至1393616908@qq.com 举报,一经查实,本站将立刻删除。

(0)

相关推荐