视频流连接器通过 FFmpeg 持续读取 RTSP、RTMP、HTTP/HLS 等视频流,将直播画面按固定间隔转换为 JPEG 图片或 MP4 视频片段,并以流式记录输出。本文介绍如何在 Flink SQL 和数据摄入 YAML 中使用视频流连接器。
背景信息
视频流连接器支持的信息如下。
类别 | 详情 |
支持类型 | 源表 |
运行模式 | 仅支持流模式 |
Changelog 模式 | 仅 INSERT |
数据格式 | 固定 Schema |
API 种类 | Flink SQL 和数据摄入 YAML |
输出内容 | JPEG 图片或 MP4 视频片段 |
您需要升级到实时计算引擎 VVR 11.9(含 Preview 预览版本)及更高版本,方可使用视频流连接器。
特色功能
两种输出模式:
KEYFRAME_EXTRACT模式按间隔抽取 JPEG 图片,MOVIE_CLIP模式按时长切分并输出 MP4 视频片段。多路视频流并行读取:支持配置静态 URL 列表。每个 URL 对应一个 Source Split,可分配给不同 Source Reader 并行处理。
多协议接入:协议解析交由 FFmpeg 完成,可接入 RTSP、RTMP、HTTP/HLS 等 FFmpeg 支持的视频源。
FFmpeg 参数透传:支持通过
ffmpeg.properties.和ffmpeg-video.properties.配置传输协议、超时、解码线程等 FFmpeg 参数。按 URL 路由:数据摄入 YAML 支持通过正则表达式从视频 URL 提取 Table ID,将不同视频流路由到不同的源表。
背压缓冲:每个 Source Reader 使用有界队列缓存已经生成的媒体记录,下游拥塞时可限制内存中的待消费记录数。
前提条件
使用视频流连接器前,需满足以下条件。
配置项 | 要求 |
视频源 | 已准备持续可用的视频流地址,且视频流中包含视频轨道。 |
网络连通性 | Flink TaskManager 所在网络能够访问视频流地址及其关联的媒体地址。 |
协议与编码 | 视频流使用 FFmpeg 支持的协议、封装格式和编码格式。 |
访问凭证(可选) | 若视频源需要鉴权,已准备可供 FFmpeg 使用的鉴权 URL 或协议参数。 |
使用限制
仅支持源表,不支持维表和结果表。
仅支持无界流式读取,不支持快照、批量读取和历史回溯。
视频 URL 列表在作业启动时确定,不支持运行时自动发现、增加或删除视频源。
SQL
语法结构
CREATE TEMPORARY TABLE video_stream_source (
video_source STRING NOT NULL,
start_ts BIGINT NOT NULL,
end_ts BIGINT NOT NULL,
`blob` BYTES NOT NULL,
mime_type STRING NOT NULL
) WITH (
'connector' = 'video-stream',
'video-stream.url-list' = '<videoStreamUrl>',
'mode' = 'KEYFRAME_EXTRACT',
'sample-interval' = '5s'
);DDL 必须按照以上顺序声明固定 Schema 中的全部五列。列的 NULL 约束不影响校验,但建议与连接器输出一致,将所有列声明为 NOT NULL。
WITH 参数
以下为视频流连接器在 Flink SQL 中的 WITH 参数。connector 固定填写 video-stream。
参数 | 说明 | 类型 | 是否必填 | 默认值 | 备注 |
| 连接器类型。 | STRING | 是 | 无 | 固定填写 |
| 静态视频流 URL 列表。 | STRING | 是 | 无 | 默认使用英文逗号分隔;至少包含一个非空且不重复的 URL。 |
| URL 列表分隔符。 | STRING | 否 |
| 按普通字符串进行分隔,不按正则表达式解析;不能为空。 |
| 媒体输出模式。 | STRING | 是 | 无 | 可选值为 |
| 抽帧间隔或视频片段目标时长。 | DURATION | 是 | 无 | 必须大于或等于 1 毫秒,例如 |
| 每个 Source Reader 最多缓存的已生成媒体记录数。 | INTEGER | 否 |
| 必须为正整数。值越大越能缓解数据抖动,代价是会占用更多内存。 |
| 传递给 FFmpeg Grabber 的通用输入参数。 | STRING | 否 | 无 | 去掉前缀后的参数会被转发给 FFmpeg 。 例如,配置 说明 完整的可配置项参见 FFmpeg Formats、FFmpeg Protocols。 |
| 传递给 FFmpeg Grabber 的视频输入参数。 | STRING | 否 | 无 | 去掉前缀后的参数会被转发给 FFmpeg 视频编码器。 例如,配置 说明 完整的可配置项参见 FFmpeg Codecs。 |
元数据字段映射
视频流连接器使用固定 Schema。SQL DDL 必须声明以下全部列,并保持相同顺序。
字段 | Flink SQL 类型 | 是否为空 | 说明 |
| STRING | 否 | 产生当前记录的视频流 URL,与配置中的 URL 一致。 |
| BIGINT | 否 | 图片或视频片段的开始时间,Unix 时间戳,单位为毫秒。首次画面的时间取连接器开始收到该视频流画面时的系统时间,后续时间根据媒体时间线递增。 |
| BIGINT | 否 | 图片或视频片段的结束时间,Unix 时间戳,单位为毫秒。图片记录的值与 |
| BYTES | 否 | JPEG 图片或 MP4 视频片段的二进制内容。 |
| STRING | 否 | 媒体类型。图片为 |
使用示例
示例一:读取 RTSP 视频并抽取 JPEG 图片
以下示例每 5 秒抽取一张画面,并使用 TCP 方式连接 RTSP 服务。
CREATE TEMPORARY TABLE video_stream_source (
video_source STRING NOT NULL,
start_ts BIGINT NOT NULL,
end_ts BIGINT NOT NULL,
`blob` BYTES NOT NULL,
mime_type STRING NOT NULL
) WITH (
'connector' = 'video-stream',
'video-stream.url-list' = 'rtsp://<host>:<port>/<path>',
'mode' = 'KEYFRAME_EXTRACT',
'sample-interval' = '5s',
-- 以下参数直接透传给 FFmpeg 库
'ffmpeg.properties.rtsp_transport' = 'tcp',
'ffmpeg-video.properties.threads' = '1'
);
INSERT INTO paimon_catalog.default.video_stream_table
SELECT video_source, start_ts, end_ts, `blob`, mime_type
FROM video_stream_source;示例二:读取 RTMP 视频并生成 MP4 片段
以下示例将直播流切分为目标时长为 10 秒的 MP4 片段。
CREATE TEMPORARY TABLE video_clip_source (
video_source STRING NOT NULL,
start_ts BIGINT NOT NULL,
end_ts BIGINT NOT NULL,
`blob` BYTES NOT NULL,
mime_type STRING NOT NULL
) WITH (
'connector' = 'video-stream',
'video-stream.url-list' = 'rtmp://<host>:<port>/<app>/<stream>',
'mode' = 'MOVIE_CLIP',
'sample-interval' = '10s',
'record-queue-capacity' = '4'
);示例三:读取多路视频流
以下示例使用 | 分隔两个 URL。当 URL 查询参数中含有逗号时,建议采用这种写法。
CREATE TEMPORARY TABLE multi_video_source (
video_source STRING NOT NULL,
start_ts BIGINT NOT NULL,
end_ts BIGINT NOT NULL,
`blob` BYTES NOT NULL,
mime_type STRING NOT NULL
) WITH (
'connector' = 'video-stream',
'video-stream.url-list' = 'https://<host>/live.m3u8?tracks=video,audio|rtsp://<host>/<camera>',
'video-stream.url-list-separator' = '|',
'mode' = 'KEYFRAME_EXTRACT',
'sample-interval' = '3s'
);示例四:将视频内容写入 Paimon Blob
以下示例将视频流连接器输出的媒体二进制写入 Paimon Blob 列。Paimon Blob 表为 append-only 表,不要定义主键。
CREATE TABLE my_catalog.my_db.video_archive (
video_source STRING,
start_ts BIGINT,
end_ts BIGINT,
`blob` BYTES,
mime_type STRING
) WITH (
'blob-field' = 'blob',
'row-tracking.enabled' = 'true',
'data-evolution.enabled' = 'true'
);
INSERT INTO paimon_catalog.default.video_stream_table
SELECT video_source, start_ts, end_ts, `blob`, mime_type
FROM video_stream_source;数据摄入
语法结构
source:
type: video-stream
name: Video Stream Source
video-stream.url-list: rtsp://<host>:<port>/<path>
mode: KEYFRAME_EXTRACT
sample-interval: 5s
transform:
- source-table: default
projection: video_source, start_ts, end_ts, blob, mime_type
table-options: blob-field=blob;row-tracking.enabled=true;data-evolution.enabled=true
route:
- source-table: default
sink-table: <sinkDatabase>.<sinkTable>
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: <catalogEndpoint>
catalog.properties.warehouse: <warehouseName>未配置 url.to.table-id.matching.pattern 时,所有视频流均写入名为 default 的源表。
配置项
以下为视频流连接器在数据摄入 YAML 中的配置项。type 固定填写 video-stream。
参数 | 说明 | 类型 | 是否必填 | 默认值 | 备注 |
| 连接器类型。 | STRING | 是 | 无 | 固定填写 |
| Source 名称。 | STRING | 否 | 无 | 用于标识当前 Source。 |
| 静态视频流 URL 列表。 | STRING | 是 | 无 | 默认使用英文逗号分隔;至少包含一个非空且不重复的 URL。 |
| URL 列表分隔符。 | STRING | 否 |
| 按普通字符串进行分隔,不按正则表达式解析;不能为空。 |
| 媒体输出模式。 | STRING | 是 | 无 | 可选值为 |
| 抽帧间隔或视频片段目标时长。 | DURATION | 是 | 无 | 必须大于或等于 1 毫秒。 |
| 每个 Source Reader 最多缓存的已生成媒体记录数。 | INTEGER | 否 |
| 必须为正整数。 |
| 从完整视频 URL 中提取 Table ID 的正则表达式。 | STRING | 否 | 无 | 必须包含 1~3 个捕获组并匹配完整 URL;未匹配的 URL 写入 |
| 传递给 FFmpeg Grabber 的通用输入参数。 | STRING | 否 | 无 | 去掉前缀后的名称作为 FFmpeg 参数名。 |
| 传递给 FFmpeg Grabber 的视频输入参数。 | STRING | 否 | 无 | 去掉前缀后的名称作为 FFmpeg 视频参数名。 |
URL 与 Table ID 映射
url.to.table-id.matching.pattern 使用 Java 正则表达式对完整 URL 进行匹配,并按照捕获组数量生成 Table ID。例如,视频地址为 rtmp://server/video/shop001/camera_01,可使用以下配置将其映射为 shop001.camera_01:
url.to.table-id.matching.pattern: ^rtmp://[^/]+/video/([^/]+)/([^/?]+)$正则表达式必须有且仅有 1~3 个捕获组,且捕获结果不能为空。未配置正则表达式或 URL 未匹配时,连接器会使用 default 作为 Table ID。
使用示例
示例一:将单路视频写入 Paimon
以下示例每 5 秒抽取一张 JPEG 图片,并将图片内容存入 Paimon Blob。
source:
type: video-stream
name: Camera Source
video-stream.url-list: rtsp://<host>:<port>/<camera>
mode: KEYFRAME_EXTRACT
sample-interval: 5s
record-queue-capacity: 8
ffmpeg.properties.rtsp_transport: tcp
ffmpeg.properties.timeout: 10000000
transform:
- source-table: default
projection: video_source, start_ts, end_ts, blob, mime_type
table-options: blob-field=blob;row-tracking.enabled=true;data-evolution.enabled=true
route:
- source-table: default
sink-table: video_db.camera_frames
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: <catalogEndpoint>
catalog.properties.warehouse: <warehouseName>
pipeline:
name: Camera Frames Pipeline
parallelism: 1示例二:按 URL 将多路视频映射为不同表
以下示例将两路 RTMP 视频分别映射为 shop001.camera_01 和 shop002.camera_02。之后可在 route 中使用 Pipeline 支持的表匹配规则,将这些源表路由至对应结果表。
source:
type: video-stream
name: Multi-camera Source
video-stream.url-list: rtmp://server/video/shop001/camera_01,rtmp://server/video/shop002/camera_02
mode: MOVIE_CLIP
sample-interval: 10s
url.to.table-id.matching.pattern: ^rtmp://[^/]+/video/([^/]+)/([^/?]+)$
route:
- source-table: shop001.camera_01
sink-table: video_db.shop001_camera_01
- source-table: shop002.camera_02
sink-table: video_db.shop002_camera_02
sink:
type: paimon
name: Paimon Sink
catalog.properties.metastore: rest
catalog.properties.uri: <catalogEndpoint>
catalog.properties.warehouse: <warehouseName>
pipeline:
name: Multi-camera Clips Pipeline
parallelism: 2常见问题
作业启动后没有输出数据
依次检查以下项目:
TaskManager 是否能够访问视频 URL。
FFmpeg 是否支持该协议、封装格式和编码格式。
视频流是否包含可解码的视频轨道。
sample-interval是否设置得过大。RTSP 服务是否要求配置
ffmpeg.properties.rtsp_transport: tcp。
为什么作业恢复后没有从中断位置继续读取
由于直播流没有可稳定恢复的媒体位点,Checkpoint 只能恢复 URL Split 分配,故障期间可能存在未采集的画面,重新连接后也需要说明是否可能产生重复记录。因此,视频流连接器不提供从中断位置续读或端到端 exactly-once 保证。
如何降低 CPU 和内存消耗
增大
sample-interval,减少图片数量或视频切片频率。只需要单帧分析时优先使用
KEYFRAME_EXTRACT。调小
record-queue-capacity,限制积压媒体记录数。将多路 URL 分配到适当的 Source 并行度,避免单个 Source Reader 同时承担过多视频流。
建议将
ffmpeg-video.properties.threads设置为 1 以降低 CPU 占用率。