视频流(RTSP/RTMP/HLS)(公测中)

更新时间:
复制 MD 格式

视频流连接器通过 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

参数

说明

类型

是否必填

默认值

备注

connector

连接器类型。

STRING

固定填写 video-stream

video-stream.url-list

静态视频流 URL 列表。

STRING

默认使用英文逗号分隔;至少包含一个非空且不重复的 URL。

video-stream.url-list-separator

URL 列表分隔符。

STRING

,

按普通字符串进行分隔,不按正则表达式解析;不能为空。

mode

媒体输出模式。

STRING

可选值为 KEYFRAME_EXTRACTMOVIE_CLIP

sample-interval

抽帧间隔或视频片段目标时长。

DURATION

必须大于或等于 1 毫秒,例如 500ms5s1min

record-queue-capacity

每个 Source Reader 最多缓存的已生成媒体记录数。

INTEGER

16

必须为正整数。值越大越能缓解数据抖动,代价是会占用更多内存。

ffmpeg.properties.*

传递给 FFmpeg Grabber 的通用输入参数。

STRING

去掉前缀后的参数会被转发给 FFmpeg 。

例如,配置 ffmpeg.properties.rtsp_transport = tcp 时,等价于传递 rtsp_transport=tcp命令行参数给 FFmpeg。

说明

完整的可配置项参见 FFmpeg FormatsFFmpeg Protocols

ffmpeg-video.properties.*

传递给 FFmpeg Grabber 的视频输入参数。

STRING

去掉前缀后的参数会被转发给 FFmpeg 视频编码器。

例如,配置 ffmpeg-video.properties.threads = 1 时,等价于传递 threads=1 Codec 参数给 FFmpeg。

说明

完整的可配置项参见 FFmpeg Codecs

元数据字段映射

视频流连接器使用固定 Schema。SQL DDL 必须声明以下全部列,并保持相同顺序。

字段

Flink SQL 类型

是否为空

说明

video_source

STRING

产生当前记录的视频流 URL,与配置中的 URL 一致。

start_ts

BIGINT

图片或视频片段的开始时间,Unix 时间戳,单位为毫秒。首次画面的时间取连接器开始收到该视频流画面时的系统时间,后续时间根据媒体时间线递增。

end_ts

BIGINT

图片或视频片段的结束时间,Unix 时间戳,单位为毫秒。图片记录的值与 start_ts 相同。

blob

BYTES

JPEG 图片或 MP4 视频片段的二进制内容。

mime_type

STRING

媒体类型。图片为 image/jpeg,视频片段为 video/mp4

使用示例

示例一:读取 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

参数

说明

类型

是否必填

默认值

备注

type

连接器类型。

STRING

固定填写 video-stream

name

Source 名称。

STRING

用于标识当前 Source。

video-stream.url-list

静态视频流 URL 列表。

STRING

默认使用英文逗号分隔;至少包含一个非空且不重复的 URL。

video-stream.url-list-separator

URL 列表分隔符。

STRING

,

按普通字符串进行分隔,不按正则表达式解析;不能为空。

mode

媒体输出模式。

STRING

可选值为 KEYFRAME_EXTRACTMOVIE_CLIP

sample-interval

抽帧间隔或视频片段目标时长。

DURATION

必须大于或等于 1 毫秒。

record-queue-capacity

每个 Source Reader 最多缓存的已生成媒体记录数。

INTEGER

16

必须为正整数。

url.to.table-id.matching.pattern

从完整视频 URL 中提取 Table ID 的正则表达式。

STRING

必须包含 1~3 个捕获组并匹配完整 URL;未匹配的 URL 写入 default 表。

ffmpeg.properties.*

传递给 FFmpeg Grabber 的通用输入参数。

STRING

去掉前缀后的名称作为 FFmpeg 参数名。

ffmpeg-video.properties.*

传递给 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_01shop002.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

常见问题

作业启动后没有输出数据

依次检查以下项目:

  1. TaskManager 是否能够访问视频 URL。

  2. FFmpeg 是否支持该协议、封装格式和编码格式。

  3. 视频流是否包含可解码的视频轨道。

  4. sample-interval 是否设置得过大。

  5. 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 占用率。