用 Cloudflare Stream、Workers 与 Durable Objects 编排持续视频处理流水线

2026-10-03 35 预计阅读时间: 1 分钟
来源: blog.cloudflare.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:11 分钟

持续视频处理和普通的 HTTP 请求不是同一类工作:转码、叠加字幕、内容检测或格式转换可能运行数小时,而 Worker 更适合接收请求、鉴权和调度。Streamline 展示的核心思路,是让 Cloudflare Stream 负责视频输入与分发,让 Workers 和 Durable Objects 管理流水线状态,再把真正消耗 CPU 的媒体处理交给容器化引擎。

这种拆分的价值不只是“能跑 FFmpeg”,而是让一个长时间运行的视频任务具备明确的身份、状态和恢复路径。

把控制面与媒体数据面分开

一条持续流水线可以拆成三层:

  • Cloudflare Stream:承担视频接入和播放分发。
  • Worker:提供创建、查询、停止流水线的 HTTP API,并执行鉴权、参数校验和限流。
  • Durable Object:为每条流水线保存唯一状态,串行处理启动、停止和重试命令。
  • 容器化媒体引擎:运行 FFmpeg、GStreamer 或自研处理程序,执行真正的长时间媒体计算。

关键边界在于:不要把完整转码进程塞进一次 Worker 请求。Worker 应当快速接受任务并返回;媒体进程则在容器中持续运行。Durable Object 充当协调器,例如维护以下状态机:

idle -> starting -> running -> stopping -> stopped
                    |
                    +-> failed -> restarting

使用流水线 ID 创建 Durable Object,可以让同一条流水线的并发命令到达同一个协调点。这样能避免两个启动请求同时创建两个 FFmpeg 进程,也方便实现幂等重试。

一个可改造的最小协调器

下面是一个示意项目。这里做了两个明确假设:

  1. 容器化媒体引擎暴露了受令牌保护的 POST /jobs/start 接口。
  2. 输入是可由 FFmpeg 读取的 HTTP/HLS 地址,输出是已经创建好的 RTMP 或 RTMPS 推流地址。

实际部署时,应按照当前 Cloudflare 容器或服务绑定方式替换 MEDIA_ENGINE_URL,并通过 Stream 的实际 API 或控制台管理输入、输出地址。

wrangler.toml:

name = 'video-pipeline-control'
main = 'src/index.js'
compatibility_date = '2025-01-01'

[durable_objects]
bindings = [
  { name = 'PIPELINES', class_name = 'PipelineCoordinator' }
]

[[migrations]]
tag = 'v1'
new_classes = ['PipelineCoordinator']

[vars]
MEDIA_ENGINE_URL = 'https://media-engine.example.internal'

src/index.js:

export default {
  async fetch(request, env) {
    const url = new URL(request.url);
    const match = url.pathname.match(/^\/pipelines\/([^/]+)\/(start|status)$/);

    if (!match) {
      return Response.json({ error: 'not_found' }, { status: 404 });
    }

    const [, pipelineId, action] = match;
    const objectId = env.PIPELINES.idFromName(pipelineId);
    const coordinator = env.PIPELINES.get(objectId);
    const body = request.method === 'POST' ? await request.text() : undefined;

    return coordinator.fetch(`https://coordinator/${action}`, {
      method: request.method,
      headers: {
        'content-type': 'application/json',
        'x-pipeline-id': pipelineId
      },
      body
    });
  }
};

export class PipelineCoordinator {
  constructor(state, env) {
    this.state = state;
    this.env = env;
  }

  async fetch(request) {
    const url = new URL(request.url);
    const pipelineId = request.headers.get('x-pipeline-id');

    if (url.pathname === '/status' && request.method === 'GET') {
      const job = await this.state.storage.get('job');
      return Response.json(job ?? { pipelineId, status: 'idle' });
    }

    if (url.pathname !== '/start' || request.method !== 'POST') {
      return Response.json({ error: 'method_not_allowed' }, { status: 405 });
    }

    const existing = await this.state.storage.get('job');
    if (existing && ['starting', 'running'].includes(existing.status)) {
      return Response.json(existing, { status: 200 });
    }

    const spec = await request.json();
    if (!spec.inputUrl || !spec.outputUrl) {
      return Response.json(
        { error: 'inputUrl and outputUrl are required' },
        { status: 400 }
      );
    }

    const job = {
      pipelineId,
      status: 'starting',
      inputUrl: spec.inputUrl,
      createdAt: new Date().toISOString()
    };
    await this.state.storage.put('job', job);

    try {
      const response = await fetch(`${this.env.MEDIA_ENGINE_URL}/jobs/start`, {
        method: 'POST',
        headers: {
          'content-type': 'application/json',
          'authorization': `Bearer ${this.env.MEDIA_ENGINE_TOKEN}`
        },
        body: JSON.stringify({
          pipeline_id: pipelineId,
          input_url: spec.inputUrl,
          output_url: spec.outputUrl
        })
      });

      if (!response.ok) {
        throw new Error(`media engine returned ${response.status}`);
      }

      const engineJob = await response.json();
      const running = {
        ...job,
        status: 'running',
        engineJobId: engineJob.job_id,
        startedAt: new Date().toISOString()
      };
      await this.state.storage.put('job', running);
      return Response.json(running, { status: 202 });
    } catch (error) {
      const failed = {
        ...job,
        status: 'failed',
        error: String(error),
        failedAt: new Date().toISOString()
      };
      await this.state.storage.put('job', failed);
      return Response.json(failed, { status: 502 });
    }
  }
}

不要把引擎令牌写入 wrangler.toml。可以这样写入 Secret:

npx wrangler secret put MEDIA_ENGINE_TOKEN
npx wrangler deploy

创建流水线时,为 inputUrl 和 outputUrl 换成自己的临时签名地址或直播输入地址:

curl -X POST 'https://worker.example.com/pipelines/camera-42/start' \
  -H 'content-type: application/json' \
  -d '{
    "inputUrl": "https://input.example.com/live/index.m3u8",
    "outputUrl": "rtmps://output.example.com/live/REPLACE_WITH_SECRET_KEY"
  }'

curl 'https://worker.example.com/pipelines/camera-42/status'

这个示例的幂等范围是“正在启动或运行时不重复创建”。生产系统还需要请求幂等键、版本号或期望状态,避免停止与重启命令互相覆盖。

容器中的媒体引擎

下面的 FastAPI 服务可以作为本地可运行的最小引擎。它只演示启动 FFmpeg,不代表完整的生产实现。

engine/app.py:

import os
import re
import subprocess
import uuid

from fastapi import FastAPI, Header, HTTPException
from pydantic import BaseModel

app = FastAPI()
TOKEN = os.environ['ENGINE_TOKEN']
JOBS = {}


class StartJob(BaseModel):
    pipeline_id: str
    input_url: str
    output_url: str


def authorize(value: str | None) -> None:
    if value != f'Bearer {TOKEN}':
        raise HTTPException(status_code=401, detail='unauthorized')


@app.post('/jobs/start')
def start_job(spec: StartJob, authorization: str | None = Header(default=None)):
    authorize(authorization)

    if not re.fullmatch(r'[A-Za-z0-9_-]{1,80}', spec.pipeline_id):
        raise HTTPException(status_code=400, detail='invalid pipeline_id')
    if not spec.input_url.startswith(('https://', 'http://')):
        raise HTTPException(status_code=400, detail='unsupported input protocol')
    if not spec.output_url.startswith(('rtmps://', 'rtmp://')):
        raise HTTPException(status_code=400, detail='unsupported output protocol')

    current = JOBS.get(spec.pipeline_id)
    if current and current.poll() is None:
        return {'job_id': spec.pipeline_id, 'status': 'already_running'}

    command = [
        'ffmpeg', '-nostdin', '-hide_banner', '-loglevel', 'warning',
        '-reconnect', '1', '-reconnect_streamed', '1',
        '-reconnect_delay_max', '5',
        '-i', spec.input_url,
        '-c:v', 'libx264', '-preset', 'veryfast',
        '-c:a', 'aac', '-f', 'flv', spec.output_url
    ]

    process = subprocess.Popen(
        command,
        stdin=subprocess.DEVNULL,
        stdout=subprocess.DEVNULL,
        stderr=subprocess.STDOUT
    )
    job_id = f'{spec.pipeline_id}-{uuid.uuid4().hex[:8]}'
    JOBS[spec.pipeline_id] = process
    return {'job_id': job_id, 'pid': process.pid, 'status': 'started'}

engine/Dockerfile:

FROM python:3.12-slim

RUN apt-get update \
    && apt-get install -y --no-install-recommends ffmpeg \
    && rm -rf /var/lib/apt/lists/*

WORKDIR /app
RUN pip install --no-cache-dir fastapi==0.115.0 uvicorn==0.30.6
COPY app.py /app/app.py

EXPOSE 8080
CMD ["uvicorn", "app:app", "--host", "0.0.0.0", "--port", "8080"]

本地构建和启动:

docker build -t media-engine ./engine
docker run --rm -p 8080:8080 \
  -e ENGINE_TOKEN='dev-token' \
  media-engine

如果处理逻辑只是复制编码、抽帧或叠加滤镜,可以修改 FFmpeg 参数。不要允许 API 调用者直接提交任意命令行参数,否则很容易形成命令注入或资源滥用入口。

长任务真正困难的是恢复

“容器成功接受启动请求”不等于“流水线一直健康”。更完整的实现通常还需要:

  • 引擎定期向协调器发送带签名的心跳,包含进程状态、最后一个视频时间戳和处理延迟。
  • Durable Object 使用 alarm 检查租约;超过一定时间没有心跳,就把任务标记为失联并决定是否重启。
  • 容器重启后重新读取期望状态,而不是依赖内存中的 JOBS 字典。
  • 启动和停止操作携带单调递增的 generation,防止旧回调覆盖新任务。
  • 每个进程设置 CPU、内存、码率和最长重启次数,避免坏输入造成无限重启。
  • 日志统一记录 pipelineId、engineJobId 和输入事件 ID,便于跨 Worker、Durable Object 与容器排查。

安全方面尤其要关注 URL:即使已经限制为 HTTP,也仍可能产生 SSRF。生产环境应校验域名白名单、阻止私有网段访问,并避免在日志中输出完整推流密钥。签名播放地址和输入地址还可能过期,因此续签策略也必须成为流水线状态的一部分。

采用时的检查清单

这套架构适合持续直播处理、动态水印、多路输出和模型推理等长任务,但会增加协调与恢复成本。落地前至少确认:

  • Worker 只承担短时控制逻辑,没有等待媒体进程结束。
  • 每条流水线映射到稳定的 Durable Object ID。
  • 启动、停止、重试都具备幂等语义。
  • 容器退出、网络中断和输入断流都有明确状态转换。
  • 推流密钥、引擎令牌和签名 URL 不进入代码仓库或普通日志。
  • 已定义处理延迟、断帧时间、重启次数和资源消耗等监控指标。

真正可靠的视频流水线不是一个长期运行的 FFmpeg 命令,而是一套能回答“任务现在在哪里、应该处于什么状态、失败后由谁恢复”的控制系统。Workers、Durable Objects 和容器化媒体引擎的组合,正好把这三个问题拆到了各自合适的执行环境中。


相关推荐