持续视频处理和普通的 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 进程,也方便实现幂等重试。
一个可改造的最小协调器
下面是一个示意项目。这里做了两个明确假设:
- 容器化媒体引擎暴露了受令牌保护的
POST /jobs/start接口。 - 输入是可由 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 和容器化媒体引擎的组合,正好把这三个问题拆到了各自合适的执行环境中。