
Daft 图像分类基准测试80 万张图片的分布式 GPU 推理流水线以及与 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 仓库中的 图像分类基准测试文档完整解析这一多模态 AI 基准测试的工作负载设计、三引擎Daft / Ray Data / Spark实测结果以及 Daft 实现中的关键源码机制daft.clsGPU 模型封装、download/decode_image原生算子、批量推理方法与集群资源配置。读完本篇你可以理解如何在 Ray 集群上跑通一条完整的读 Parquet → 下载图片 → 解码 → 预处理 → ResNet18 推理 → 写 Parquet的多模态推理流水线并掌握其中的关键参数与调优点。基准测试工作负载设计图像分类基准测试的目标是模拟真实的大规模 AI 推理场景。其官方定义如下引自 README数据规模共803,580 行数据由 80,358 张唯一图片各重复 10 次构成用于放大推理负载、验证分布式扩展性模型TorchVision 的 ResNet18ImageNet 预训练权重处理流程下载图片 → 应用预处理变换 → 运行推理预测 ImageNet 类别标签任务分布在多个 GPU 节点上输入数据集ImageNet 基准数据集S3 Parquet 格式位于s3://daft-oss-public-datasets/imagenet/benchmark每行含image_url字段输出格式Parquet包含图片 URL 与预测标签两列集群规格8 个 GPU worker 节点AWSg6.xlarge实例每节点 1 块 L4 GPU基准测试日期2024 年 9 月 22 日框架版本Daft 0.6.2、Ray Data 2.49.2、AWS EMR Spark 7.10.0。该基准是仓库 AI 基准测试套件 的四个工作负载之一同套件还包括音频转写、文档嵌入与视频目标检测共同覆盖多模态数据处理的不同侧面。性能结果同一数据集、同一集群规格下三引擎的实测运行时间如下来自 README引擎运行时间Daft4m 23sRay Data23m 30sSpark45m 7sAI 基准汇总页 中给出了同一数据规模803,580 张图下的一致记录Daft 4m 23s、Ray Data 23m 30s、Spark 45m 7s。这一结果展示了在多模态下载 解码 模型推理混合负载下各引擎的端到端差异——其中图片下载与解码这类 CPU/IO 密集步骤以及 GPU 推理步骤之间的流水化程度是影响总时长的关键。Daft 实现解析daft_main.py核心实现在 daft_main.py整条流水线约 80 行代码可拆分为五个环节。1. 集群就绪等待与运行器配置NUM_GPU_NODES 8 daft.set_runner_ray() # Wait for Ray cluster to be ready ray.remote def warmup(): pass ray.get([warmup.remote() for _ in range(64)])脚本先通过daft.set_runner_ray()指定 Ray 作为执行后端随后提交 64 个空的ray.remote任务8 节点 × 8等待 Ray 自动扩容把 8 个 worker 全部拉起后再开始正式计时避免把集群冷启动时间计入基准结果。2. GPU 模型类daft.cls 封装 ResNet18weights ResNet18_Weights.DEFAULT transform transforms.Compose([transforms.ToTensor(), weights.transforms()]) daft.cls( max_concurrencyNUM_GPU_NODES, gpus1, ) class ResNetModel: def __init__(self): self.weights weights self.device torch.device(cuda if torch.cuda.is_available() else cpu) self.model resnet18(weightsweights).to(self.device) self.model.eval() daft.method.batch( return_dtypedaft.DataType.string(), batch_sizeBATCH_SIZE, ) def __call__(self, images): if len(images) 0: return [] torch_batch torch.from_numpy(np.array(images.to_pylist())).to(self.device) with torch.inference_mode(): prediction self.model(torch_batch) predicted_classes prediction.argmax(dim1).detach().cpu() predicted_labels [self.weights.meta[categories][i] for i in predicted_classes] return predicted_labels结合 Daft 类 UDF 的源码定义这里涉及两个关键机制daft.cls模型复用该装饰器把普通 Python 类转为 Daft 用户自定义类每个实例的__init__只在查询执行时被惰性调用一次之后同一实例复用于多批数据。对于加载模型权重并绑定 CUDA 设备这类昂贵初始化按 8 个 actor 各加载一次而非按行加载是 GPU 推理 UDF 的标准写法。参数方面gpus1每个实例申请 1 块 GPU源码支持 0~1 之间的小数用于多个小模型共享一张 GPU但不支持大于 1 的小数max_concurrencyNUM_GPU_NODES同步方法下控制 actor 池规模即最多 8 个并发推理实例正好对应 8 个 GPU 节点另有cpus、use_process、max_retries、on_errorraise/log/ignore、ray_options等可选参数可进一步调优。daft.method.batch批量推理将方法声明为按批调用batch_size100与文件顶部BATCH_SIZE 100一致即每次向 GPU 送 100 张图片return_dtypedaft.DataType.string()显式声明输出列类型为字符串标签名。torch.inference_mode()与argmax取预测类别再通过weights.meta[categories]把类别索引映射为 ImageNet 标签名。3. 端到端流水线表达式daft.set_planning_config( default_io_configdaft.io.IOConfig(s3daft.io.S3Config.from_env().replace(requester_paysTrue)) ) df daft.read_parquet(INPUT_PATH) df df.with_column( decoded_image, df[image_url].download().decode_image(modedaft.ImageMode.RGB), ) df df.with_column( norm_image, df[decoded_image].apply( funclambda image: transform(image), return_dtypedaft.DataType.tensor(dtypedaft.DataType.float32(), shapeIMAGE_DIM), ), ) df df.with_column(label, ResNetModel()(col(norm_image))) df df.select(image_url, label) df.write_parquet(OUTPUT_PATH)这条链路的每个环节对应一个原生或 UDF 算子daft.read_parquet(INPUT_PATH)读取 S3 上的 Parquet。注意IOConfig中requester_paysTrue——输入桶属于 S3 请求者付费桶必须显式开启计费模式才能访问image_url.download()把 URL 字符串列逐行下载为字节列。从 表达式定义 看默认max_connections32每分区并发连接数、on_errorraise下载逻辑本身是引擎级并发 IO而非逐行 Python 循环.decode_image(modedaft.ImageMode.RGB)将字节解码为统一 RGB 模式的图像列同样是 Rust 引擎内的原生算子见 decode_image 表达式支持on_errorraise/null错误策略apply预处理用 TorchVision 的transforms.Compose([ToTensor(), weights.transforms()])把图片转为(3, 224, 224)的 float32 张量列return_dtypedaft.DataType.tensor(dtypefloat32, shapeIMAGE_DIM)显式声明张量 schema让引擎可以在后续 stage 间正确传递该列ResNetModel()(col(norm_image))调用上一步定义的 Daft 类 UDF 做 GPU 推理产出label列selectwrite_parquet只保留image_url与label两列写出 S3对应 README 中Parquet with image URLs and predicted labels的输出格式。整条流水线是惰性声明的从读表到写表没有显式的.collect()write_parquet触发一次性执行下载、解码、预处理、推理各阶段在引擎内部流水化重叠执行。4. 计时方式脚本用time.time()在read_parquet之前开始计时、在write_parquet之后结束打印总秒数。而在 CI 场景中仓库提供了统一入口 run_ai_benchmark.py它以 Ray Job Submission 方式提交DAFT_RUNNERray DAFT_PROGRESS_BAR0 python daft_main.py作为 entrypoint先做一次 warmup 运行再正式运行 2 次取平均最后把结果连同 Daft 版本等元数据上传记录。Ray Data 对照实现ray_data_main.py 用相同的 ResNet18 模型、相同输入输出路径与BATCH_SIZE 100构建了对照流水线paths ray.data.read_parquet(INPUT_PATH).take_all() paths [row[image_url] for row in paths] ds ( ray.data.read_images(paths, include_pathsTrue, ignore_missing_pathsTrue) .map(fntransform_image) .map_batches(fnResNetActor, batch_sizeBATCH_SIZE, num_gpus1.0, concurrencyNUM_GPU_NODES) .select_columns([path, label]) ) ds.write_parquet(OUTPUT_PATH)两个实现的可对照差异点图片获取Ray Data 侧先用take_all()把全部 80 万条 URL 一次性拉到驱动进程再交给read_images下载Daft 侧则保持 URL 列在分布式数据内用引擎级download()算子在各分区并发下载逐行 vs 批式算子Ray Data 用.map逐行函数transform_image内部Image.fromarray(row[image]).convert(RGB)做预处理.map_batches(ResNetActor, num_gpus1.0, concurrency8)做推理Daft 侧预处理走apply 张量列推理走daft.clsdaft.method.batchGPU 资源表达num_gpus1.0/concurrency8与 Daft 的gpus1/max_concurrency8语义对应两边的 actor 池规模一致。Spark 对照实现spark.ipynb 基于 AWS EMR Spark 7.10.0 完成同一任务其要点通过%%configure设置spark.sql.execution.arrow.maxRecordsPerBatch 100将 Arrow 批大小对齐到 100与前两个实现的BATCH_SIZE 100保持一致用模块级_model_cache字典在 executor 进程内缓存resnet18模型、权重与设备并把TORCH_HOME/XDG_CACHE_HOME指到/tmp避免重复加载权重图片解码与预处理封装为pandas_udfdecode_and_preprocess_image_udf逐条Image.open字节流、convert(RGB)后过transform以ArrayType(FloatType())返回推理同样以 pandas UDF 形式完成。集群与依赖配置Ray 集群配置cluster.yaml 定义了 CI 使用的 Ray 集群AWSus-west-2节点类型head 节点不占 CPU/GPU 资源resources: {CPU: 0, GPU: 0}worker 固定min_workers: 8/max_workers: 8保证扩容后精确为 8 个 GPU worker避免自动缩放在测试期间扰动结果实例规格head 与 worker 均为g6.xlarge每节点 1 块 L4 GPUPyTorch AMIImageId: ami-0976479b866d22613100GB gp3 加密云盘安全与 IAM统一安全组ray-autoscaler-c1、IAM 角色ray-autoscaler-v1SSH 密钥ci-github-actions-ray-cluster-key环境准备setup_commands中把 Ray 的 tmp 与 object spilling 目录指到/opt/ray防止对象溢出打爆根盘并安装固定版本依赖ray[default]2.49.2、numpy1.26.4、torchvision0.22.0cu128、pillow11.3.0Daft 则以pip install daft --pre --extra-index-url ${DAFT_INDEX_URL}从预发布索引安装。依赖锁定pyproject.toml 锁定本基准的运行时环境Python3.12daft0.6.2ray[default]2.49.2与 README 中的Framework Versions一一对应保证结果可复现。结论这个基准说明了什么该基准测量的是多模态推理的端到端能力而非单纯的模型吞吐图片下载网络 IO、解码CPU、预处理CPU/张量转换、GPU 推理CUDA与结果落盘S3 写串在一条查询里。Daft 在此场景下的 4m 23s相对 Ray Data 23m 30s 与 Spark 45m 7s的优势从源码结构看主要来自三方面下载/解码作为引擎原生算子与 GPU 推理 stage 流水化重叠执行而非 Ray Data 中先take_all再read_images的两段式路径daft.clsdaft.method.batch让模型按 actor 池粒度加载一次、按 100 条/批送 GPU以及显式张量列 schemaDataType.tensor使预处理输出能以紧凑的 Arrow 张量列在 stage 间传递。若你在自己的仓库中复现类似负载可直接以 daft_main.py 为模板替换INPUT_PATH/OUTPUT_PATH与模型类并按 cluster.yaml 的思路固定 GPU worker 数量以保证结果稳定。需要注意的适用前提输入桶为 S3 请求者付费桶requester_paysTrue模型为 TorchVision 官方 ImageNet 权重且结果基于 2024 年 9 月 22 日的 Daft 0.6.2 / Ray Data 2.49.2 / EMR Spark 7.10.0 版本组合跨版本对比时应先对齐框架版本与集群规格。【免费下载链接】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),仅供参考