本文介绍 Python DataFrame API 音频算子清单、使用限制、使用方法、音频数据形态,以及音频处理 Pipeline 示例。
使用限制
仅实时计算引擎 VVR 11.8 及以上版本支持。
音频算子统一从
pyflink.multimodal.operators导入。音频处理算子依赖
soundfile库,已内置到引擎内部。
算子清单
单击算子名称,可以跳转到对应详情页的算子说明。
分类 | 算子 | 说明 |
解码 | 将 | |
将 float32 PCM | ||
格式转换 | 先解码输入,再按目标格式重新编码为新字节。适合在下游要求特定容器或编码格式时使用,例如将来源不统一的音频归一成 WAV、FLAC 等格式。 | |
先转换声道,再重采样,将音频波形统一为指定采样率和声道数。适合在 ASR、VAD 或模型推理前规范不同来源的输入。 | ||
使用 soxr HQ 改变音频波形的采样率,同时保留原声道数。适合在下游只要求特定采样率、但需要保持原声道布局时使用,例如把 48 kHz 立体声录音转换为 16 kHz 立体声。 | ||
按数组顺序拼接多个采样率、声道数和采样格式一致的音频波形。适合重组切分片段或合并分段处理结果,例如把并行降噪后的连续片段恢复为一段音频。 | ||
有效性检查 | 通过媒体头探测或完整解码,按指定校验级别判断编码音频能否继续处理并返回布尔结果。适合在外部音频接入、批量回灌或 ASR 前置质量检查,例如先用 metadata 低成本粗筛,再用 decode 排除截断、空音频和解码大小超限的数据。 | |
属性提取 | 读取音频的采样率、声道数、帧数、时长、格式、编码器、码率和对象大小等元数据。适合用于音频盘点、质量统计和按属性路由,例如按采样率选择预处理链路,或按格式识别异常数据来源。 | |
计算并返回音频的秒级时长。适合按音频长短分桶、估算处理成本或制定切分策略,例如把超长录音送入固定时长切分,将短录音直接交给 ASR。 | ||
数据过滤 | 判断音频时长是否位于指定闭区间,可只设下限、只设上限或同时设置。适合在模型推理或训练前过滤过短、过长样本,例如排除低于 1 秒的无效录音,或拦截超过模型输入上限的音频。 | |
按字节判断 | ||
音频检测 | 检测音频波形中的低振幅静音区间。 | |
使用 WebRTC VAD 检测音频波形中的语音活动区间。 | ||
音频切分 | 按指定固定时长依次切分音频,并根据 | |
按调用方提供的开始和结束时间范围切分音频,并按原顺序输出音频波形或片段引用。适合已经有字幕、日志、会议纪要或人工标注时间轴的场景,例如按字幕时间提取对应音频,或按会议纪要定位发言片段。 | ||
根据已有语音活动范围,对区间做扩边、合并、短段过滤和长段拆分,再输出音频波形或片段引用。适合处理含大量静音的通话、访谈或 Podcast,例如配合 |
使用方法
基本调用
不含专属配置参数的算子可以直接应用到 DataFrame 列。
from pyflink.dataframe import col
from pyflink.multimodal.operators import audio_metadata
result = df.with_column(
"metadata",
audio_metadata(col("audio_bytes"))
)参数调用
算子可以在传入列的同时配置参数。
from pyflink.multimodal.operators import audio_decode
result = df.with_column(
"waveform",
audio_decode(
col("audio_bytes"),
on_error="null"
)
)也可以先创建配置好的算子,再应用到一个或多个列:
decode = audio_decode(on_error="null")
result = df.with_column(
"waveform",
decode(col("audio_bytes"))
)数据类型说明
音频数据形态
数据形态 | DataFrame 类型 | 说明 |
编码音频 |
| WAV、FLAC、OGG、AIFF 或 MP3 等编码音频字节。 |
音频 URI |
| Flink FileSystem 可访问的路径或对象存储 URI。 |
音频片段引用 |
| 保存 URI 和时间范围,不复制音频内容。 |
解码波形 |
| 保存 float32 PCM Tensor、采样率、声道数和来源信息。 |
音频元数据 |
|
|
音频片段引用 (AUDIO_CLIP_REF_TYPE)
DataType.struct({
"uri": DataType.string(),
"start_time_ms": DataType.int64(),
"end_time_ms": DataType.int64()
})音频片段使用半开时间区间 [start_time_ms, end_time_ms)。None 起点或终点分别表示文件开头或结尾。
解码后的音频波形(AUDIO_WAVEFORM_TYPE)
DataType.struct({
"data": DataType.tensor("float32"),
"sample_rate": DataType.int32(),
"channels": DataType.int32(),
"frames": DataType.int64(),
"sample_format": DataType.string(),
"layout": DataType.string(),
"duration_ms": DataType.int64(),
"source_uri": DataType.string(),
"start_time_ms": DataType.int64(),
"end_time_ms": DataType.int64(),
})data 保存可变形状的 float32 PCM Tensor。
Python 算子将其暴露为形状为 (frames, channels) 的 NumPy 数组。
音频元数据 (AUDIO_METADATA_TYPE)
DataType.struct({
"sample_rate": DataType.int32(),
"channels": DataType.int32(),
"frames": DataType.int64(),
"duration_ms": DataType.int64(),
"format": DataType.string(),
"codec": DataType.string(),
"bit_rate": DataType.int64(),
"size": DataType.int64()
})字段 | 说明 |
| 采样率,单位 Hz。 |
| 声道数。 |
| 每个声道的采样帧数。 |
| 时长,单位毫秒。 |
| 编码容器格式。 |
| 编解码器短名称。 |
| 编码码率,无法获取时为 |
| 编码对象大小,单位字节;波形输入为 |
音频时间范围
静音和语音检测算子返回:
DataType.list(
DataType.struct({
"start_ms": DataType.int64(),
"end_ms": DataType.int64(),
"duration_ms": DataType.int64()
})
)时间戳切分和语音切分也接受不含 duration_ms 的形式:
DataType.list(
DataType.struct({
"start_ms": DataType.int64(),
"end_ms": DataType.int64()
})
)时间范围采用半开区间 [start_ms, end_ms)。检测和切分参数中的时间通常使用毫秒;audio_duration 和 audio_duration_filter 使用秒。物化波形片段时,毫秒边界需要换算为采样帧并进行取整,因此非整帧边界可能产生轻微的时间偏差。
空值与错误处理
场景 | 配置方式 |
解码失败时抛出异常 |
|
解码失败时返回空值 |
|
只探测容器或文件头 |
|
完整解码并校验媒体内容 |
|
说明:
on_error="null"只处理可识别的媒体数据错误;URI 无法访问、音频片段引用非法、输入类型错误、依赖缺失和参数错误仍会抛出异常;
audio_duration_filter对空值、损坏输入或未知时长返回False;audio_size_filter只检查字节长度或对象大小,因此内容损坏但大小符合条件时仍可能返回True。
完整 Pipeline 示例
以客服录音预处理为例:DataFrame 存放原始录音字节,下游需要按语音区间切出片段用于 ASR。
以下 Pipeline 依次完成解码、标准化到 16 kHz 单声道、检测语音活动区间、切分为独立片段并计算每段时长。
from pyflink.dataframe import col
from pyflink.multimodal.operators import (
audio_decode,
audio_detect_speech,
audio_duration,
audio_split_by_speech,
audio_standardize,
)
# 解码 -> 标准化 16 kHz 单声道 -> 过滤空值 -> 检测语音活动区间
prepared = (
df.with_column(
"waveform",
audio_standardize(
audio_decode(
col("audio_bytes"),
on_error="null",
),
sample_rate=16000,
channels=1,
),
)
.filter(col("waveform").is_not_null)
.with_column(
"speech_ranges",
audio_detect_speech(
col("waveform"),
aggressiveness=2,
frame_ms=30,
min_speech_ms=100,
merge_gap_ms=50,
max_segments=1024,
),
)
)
# 根据语音活动区间切分独立片段
segments = prepared.join_lateral(
audio_split_by_speech(
col("waveform"),
col("speech_ranges"),
segment_type="audio"
).alias("segment")
)
# 计算每个片段的时长
result = (
segments
.with_column(
"segment_duration_seconds",
audio_duration(col("segment"))
)
)