Daft 多模态 AI 基准实测Whisper 音频转写管线在 Daft、Ray Data 与 Spark 上的性能对比【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft本文以 Daft 仓库中的音频转写基准Audio Transcription Benchmark为核心完整还原该基准的测试对象Common Voice 17 数据集 Whisper-tiny 模型、集群环境8 节点 AWS g6.xlarge与三引擎实现细节。读完本文你将掌握如何用 Daft 的daft.func.batch、daft.cls/daft.method.batch构建带 GPU 并发控制的音频转写管线并理解它与 Ray Data、Spark 对应实现的差异来源。基准总览转写 113,800 个音频文件根据 基准说明文档该任务的核心设定如下项目内容任务使用 Whisper-tiny 模型对 113,800 个音频文件做语音转写处理内容音频重采样resampling、特征提取feature extraction、语音转文本推理输入数据集Common Voice 17S3 上的 Parquet 格式输出格式带转写文本与元数据的 Parquet集群规模8 个 worker 节点实例类型为 g6.xlarge基准日期2024-09-22框架版本Daft 0.6.2、Ray Data 2.49.2、AWS EMR Spark 7.10.0管线在逻辑上分为四个阶段三引擎的实现都围绕同一套阶段展开重采样把 FLAC 音频字节解码为波形并重采样到 16 kHzWhisper 要求的输入采样率特征提取用 Hugging FaceAutoProcessor计算 mel 频谱特征input_features模型推理Whisper-tiny 在 GPU 上批量执行model.generate得到 token id解码输出processor.batch_decode将 token id 还原为转写文本并附带transcription_length等元数据列。该基准是仓库 AI 基准系列音频转写、文档嵌入、图像分类、视频目标检测之一系列总览见 AI 基准说明。性能结果Daft 6m 22s vs Ray Data 29m 20s vs Spark 25m 46sREADME 原文记录的实测结果如下由项目方在上述集群上实测此处原样引用引擎运行时间Daft6m 22sRay Data29m 20sSpark25m 46s在 AI 基准系列的整体结果中见 AI 基准汇总音频转写是 Daft 与 Ray Data/Spark 差距最大的工作负载之一工作负载数据规模DaftRay DataSpark音频转写113,800 个音频文件6m 22s29m 20s25m 46s文档嵌入10,000 个 PDF1m 54s14m 32s8m 4s图像分类803,580 张图片4m 23s23m 30s45m 7s视频目标检测1,000 个视频11m 46s25m 54s3h 36m需要说明的适用前提结果强依赖硬件与版本——8 个 g6.xlarge每节点 1 块 GPU合计 8 GPU集群、Daft 0.6.2 / Ray Data 2.49.2 / EMR Spark 7.10.0 的具体组合。更换实例规格、模型规模或框架版本后各引擎间的相对差距可能发生变化。复现环境Ray 集群配置与依赖锁定基准的硬件环境由 cluster.yaml 定义区域AWSus-west-2通过 Ray autoscaler 管理节点节点类型head 节点与 worker 节点均为g6.xlargeNVIDIA L4 GPU 实例head 节点 CPU/GPU 资源声明为 0仅承担调度角色worker 规模min_workers: 8, max_workers: 8与 README 中8 worker nodes一致也对应代码中的NUM_GPUS 8存储每节点 100 GB gp3 加密 EBS 卷setup 命令创建/opt/ray/tmp与/opt/ray/spill目录并设置RAY_TMPDIR、RAY_object_spilling_directory随后pip install ray[default]2.49.2 numpy1.26.4 accelerate1.10.1 transformers4.56.2 torchaudio2.7.0cu128 soundfile0.13.1最后通过pip install daft --pre --extra-index-url ${DAFT_INDEX_URL}安装 Daft 预发布版本。依赖侧由 pyproject.toml 锁定daft0.6.2与ray[default]2.49.2Python 版本固定为 3.12保证结果可复现。基准的运行入口是 run_ai_benchmark.py其流程为通过JobSubmissionClient向 Ray 集群默认http://localhost:8265提交作业entrypoint 为DAFT_RUNNERray DAFT_PROGRESS_BAR0 python daft_main.pyworking_dir指向./benchmarking/ai/audio_transcription并使用uvruntime env 管理依赖先执行一次warmup run随后再正式运行 2 次取平均消除冷启动影响结果连同 Daft 版本等元数据一并整理供后续上传记录该脚本同时服务于audio_transcription、document_embedding、image_classification、video_object_detection四个基准。Daft 实现逐段解析daft_main.py完整实现见 daft_main.py。关键常量如下TRANSCRIPTION_MODEL openai/whisper-tiny NUM_GPUS 8 NEW_SAMPLING_RATE 16000 INPUT_PATH s3://daft-oss-public-datasets/common_voice_17 OUTPUT_PATH s3://eventual-dev-benchmarking-results/ai-benchmark-results/audio-transcription daft.set_runner_ray()脚本首先调用daft.set_runner_ray()把 Daft 切到 Ray 分布式执行器并用一个技巧等待集群就绪提交 64 个空ray.remotewarmup 任务并ray.get从源码结构看这是为了确认 worker 节点都已上线、资源分配完成后再开始计时。阶段一重采样普通 Python 函数 applydef resample(audio_bytes): waveform, sampling_rate torchaudio.load(io.BytesIO(audio_bytes), formatflac) waveform T.Resample(sampling_rate, NEW_SAMPLING_RATE)(waveform).squeeze() return np.array(waveform)audio.bytes列存的是 FLAC 字节。这里用torchaudio.load直接从内存字节流解码再用torchaudio.transforms.Resample线性插值到 16 kHz返回 numpy 数组以便后续 Arrow 化。阶段二特征提取daft.func.batchprocessor AutoProcessor.from_pretrained(TRANSCRIPTION_MODEL) daft.func.batch(return_dtypedaft.DataType.tensor(daft.DataType.float32())) def whisper_preprocess(resampled): extracted_features processor( resampled.to_arrow().to_numpy(zero_copy_onlyFalse).tolist(), sampling_rateNEW_SAMPLING_RATE, devicecpu, ).input_features return extracted_featuresdaft.func.batch声明这是一个逐批执行的 UDF输入以 Arrow 批量传入resampled.to_arrow().to_numpy(zero_copy_onlyFalse)把列还原为 numpy 再转 list处理器在 CPU 上计算 mel 特征返回类型显式声明为tensor(float32)让 Daft 在编译期就确定该列的数据类型与内存布局。阶段三GPU 推理daft.cls daft.method.batch这是整个管线中并发控制最关键的部分daft.cls(max_concurrencyNUM_GPUS, gpus1) class Transcriber: def __init__(self) - None: self.device cuda if torch.cuda.is_available() else cpu self.dtype torch.float16 self.model AutoModelForSpeechSeq2Seq.from_pretrained( TRANSCRIPTION_MODEL, torch_dtypeself.dtype, low_cpu_mem_usageTrue, use_safetensorsTrue, ) self.model.to(self.device) daft.method.batch( return_dtypedaft.DataType.list(daft.DataType.int32()), batch_size64, ) def __call__(self, extracted_features): spectrograms np.array(extracted_features) spectrograms torch.tensor(spectrograms).to(self.device, dtypeself.dtype) with torch.no_grad(): token_ids self.model.generate(spectrograms) return token_ids.cpu().numpy()从源码结构可以看出其并发语义daft.cls(gpus1)告诉 Daft 每个Transcriber实例占用 1 块 GPUmax_concurrencyNUM_GPUS限制全局最多同时存在 8 个实例——正好与集群的 8 块 GPU 一一对应避免模型副本超配模型在__init__中加载一次float16 safetensors 降低显存与加载开销被该实例的所有批次复用避免重复加载daft.method.batch(batch_size64)规定每次model.generate处理 64 条样本批量推理把 GPU 喂满同时torch.no_grad()关闭梯度以节省显存返回类型声明为list(int32)即每行一个 token id 列表供下一阶段的解码 UDF 消费。阶段四解码与 IO 配置daft.func.batch(return_dtypedaft.DataType.string()) def decoder(token_ids): transcription processor.batch_decode(token_ids, skip_special_tokensTrue) return transcription解码是轻量 CPU 操作用普通的批量函数即可。值得注意的是 S3 读取配置daft.set_planning_config(default_io_configdaft.io.IOConfig(s3daft.io.S3Config.from_env().replace(requester_paysTrue)))Common Voice 17 数据集所在的公共桶开启了请求者付费requester paysS3Config.from_env().replace(requester_paysTrue)从环境变量读取凭证并追加开启请求者付费标志否则 S3 读取会直接报 403。这是读取该公开数据集的必要步骤。管线组装df daft.read_parquet(INPUT_PATH) df df.with_column( resampled, df[audio][bytes].apply(resample, return_dtypedaft.DataType.list(daft.DataType.float32())), ) df df.with_column(extracted_features, whisper_preprocess(df[resampled])) df df.with_column(token_ids, Transcriber()(df[extracted_features])) df df.with_column(transcription, decoder(df[token_ids])) df df.with_column(transcription_length, df[transcription].length()) df df.exclude(token_ids, extracted_features, resampled) df.write_parquet(OUTPUT_PATH)几个值得注意的细节df[audio][bytes]使用结构化列访问说明 Common Voice 17 在 Parquet 中把音频组织为 struct含bytes、path等字段无需在 Python 侧手动拆解apply(resample, return_dtype...)对普通 Python 函数也要求显式声明返回类型与daft.func.batch的声明式风格一致保证 Daft 在执行前就完成 schema 推导末尾exclude掉中间列token_ids、extracted_features、resampled最终 Parquet 只保留转写文本、长度等轻元数据避免把巨大的波形/特征列写回存储。计时采用朴素的time.time()差值围绕read_parquet → write_parquet的完整执行包含 S3 读写在内的端到端时间。对照实现Ray Data 与 Spark 的差异点Ray Data 版本ray_data_main.pyray_data_main.py 的阶段划分与 Daft 版一一对应核心差异在算子表达与资源声明方式ds ray.data.read_parquet(INPUT_PATH) ds ds.map(unnest) ds ds.map(resample) ds ds.map_batches(whisper_preprocess) ds ds.map_batches( Transcriber, batch_sizeBATCH_SIZE, # 64 concurrencyNUM_GPUS, # 8 num_gpus1, ) ds ds.map_batches(decoder) ds ds.drop_columns([input_features, token_ids, arr]) ds.write_parquet(OUTPUT_PATH)GPU 并发同样通过concurrency8, num_gpus1, batch_size64表达参数语义与 Daft 的max_concurrency/gpus/batch_size对应脚本开头有一段带注释的unnest函数是绕开 Ray Data 扩展类型转换限制arrow_tensor_v2与arrow_variable_shaped_tensor之间不允许直接转换的 workaround——从源码注释看需要先落到存储类型再转回扩展类型这类数据布局转换开销是 Ray Data 路径上特有的摩擦每阶段都是独立的map/map_batches中间列arr、input_features、token_ids在多个算子间以对象/扩展类型传递最后才drop_columns。Spark 版本spark.ipynbspark.ipynb 用pandas_udf组织四个阶段配置上有两处关键调优// %%configure -f { executorCores: 1, conf: { spark.sql.execution.arrow.maxRecordsPerBatch: 64, spark.executorEnv.HF_HOME: /tmp/hf_home 的临时缓存目录 /tmp/huggingface } }executorCores: 1配合maxRecordsPerBatch: 64把 Arrow batch 大小对齐到 GPU 推理的 64 条批次处理器与模型用模块级缓存get_processor/get_model延迟加载避免在每个 executor 上重复 import 与下载推理侧与另两版一致float16、low_cpu_mem_usageTrue、use_safetensorsTrue、model.generate批量推理数据路径使用s3://daft-public-datasets/common_voice_17注意桶名与 Daft/Ray 脚本中的daft-oss-public-datasets不同输出用write.mode(append)写回同一结果桶。三版实现都完成了相同的五阶段流水线但 Daft 版本通过daft.cls/daft.method.batch把模型实例生命周期 GPU 资源 批大小封装进类装饰器schema 在函数定义处即被完整声明Ray Data 版本需要对扩展类型转换做额外 workaroundSpark 版本则依赖 executor 配置与手动缓存来近似等价的批量行为。如何运行该基准基于仓库现有文件复现步骤可归纳为仓库为只读参考以下操作在本地克隆中进行按 cluster.yaml 用ray up cluster.yaml拉起 AWS 集群需自备 IAM profileray-autoscaler-v1、安全组与 SSH key这些是原作者 CI 环境的资源需替换为自己的在集群可达的环境执行python benchmarking/scripts/run_ai_benchmark.py audio_transcription脚本会先 warmup 再正式跑 2 轮取平均作业实际执行的 entrypoint 是DAFT_RUNNERray DAFT_PROGRESS_BAR0 python daft_main.py工作目录自动切到benchmarking/ai/audio_transcription结果 Parquet 写入s3://eventual-dev-benchmarking-results/ai-benchmark-results/audio-transcription结果桶实际复现时可改daft_main.py中的OUTPUT_PATH为自己有写权限的桶。本地小规模验证时可把daft_main.py中的INPUT_PATH指向任意 S3/本地 Parquet 音频目录并相应调低NUM_GPUS、max_concurrency与batch_size无 GPU 时Transcriber会自动落到 CPU 与 float16→CUDA 不可用时的设备选择分支torch.cuda.is_available()判定。小结该基准以 113,800 个 Common Voice 17 音频文件为输入在 8 节点 g6.xlarge8 GPU集群上对比了 Daft 0.6.2、Ray Data 2.49.2 与 EMR Spark 7.10.0 完成 Whisper-tiny 转写的端到端时间README 记录的结果为 6m 22s / 29m 20s / 25m 46sDaft 侧的管线由普通函数 apply重采样daft.func.batch特征提取/解码daft.cls(gpus1)类 UDFGPU 推理三段式构成所有中间列类型在声明期确定GPU 并发数与集群 GPU 数严格对齐复现所需的全部工件集中在 benchmarking/ai/audio_transcription/ 目录README结果与设定、daft_main.py/ray_data_main.py/spark.ipynb三引擎实现、cluster.yaml集群定义、pyproject.toml依赖锁定配合 benchmarking/scripts/run_ai_benchmark.py 即可完成 warmup 2 轮取平均的标准化跑分流程。【免费下载链接】DaftHigh-performance data engine for AI and multimodal workloads. Process images, audio, video, and structured data at any scale项目地址: https://gitcode.com/GitHub_Trending/da/Daft创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考