【发布时间】:2022-08-18 11:39:43
【问题描述】:
我们正在尝试使用 Google Cloud Dataflow 构建一个简单的基于 GPU 的分类管道,如下所示:Pub/Sub 请求带有指向 GCS 上的文件的链接 → 从 GCS 读取数据 → 切碎和批处理数据 → 运行推理火炬。
背景
我们使用改编自pytorch-minimalsample 的自定义 Docker 映像在 Dataflow 上部署我们的管道。
我们使用 pathy 从 GCS 中提取 Pub/Sub 消息并下载数据音频文件,然后将音频切成块进行分类。
我们改编了 Beam 相对较新的RunInference 功能。目前,Dataflow 上的 RunInference 不支持 GPU
(见公开问题https://issues.apache.org/jira/browse/BEAM-13986)。在部署到 Dataflow 之前在本地构建 Beam 管道时,模型初始化步骤无法识别 CUDA 环境并默认使用 CPU 设备进行推理。此配置会传播到正确启用 GPU 的 Dataflow 执行环境。因此,如果在没有 CUDA 设备检查的情况下请求,我们会强制使用 GPU 设备。除此之外,代码与一般的RunInference 代码相同:BatchElements 操作后跟调用模型的ParDo。
问题
一切正常,但 GPU 推理速度非常慢 - 比我们在 Google Cloud Compute Engine 上处理批次的相同 GPU 实例的时间要慢得多。
我们正在寻找有关如何调试和加速管道的建议。我们怀疑这个问题可能与线程以及 Beam/Dataflow 如何管理流水线阶段的负载有关。在ParDo 函数中尝试访问 GPU 的多个线程不断遇到 CUDA OOM 问题。我们使用--num_workers=1 --experiment=\"use_runner_v2\" --experiment=\"no_use_multiple_sdk_containers\" 启动我们的工作,以完全避免多处理。我们看到这个2021 beam summit talk on using Dataflow for local ML batch inference 建议更进一步,只使用单个工作线程--number_of_worker_harness_threads=1。但是,理想情况下,我们不想这样做:在像这样的 ML 管道中,让多个线程执行从存储桶下载数据和准备批处理的 I/O 工作是很常见的做法,这样 GPU 就永远不会坐下来。闲置的。不幸的是,似乎没有办法告诉 beam 使用某个最大线程数每阶段(?),所以我们能想出的最佳解决方案是使用信号量保护 GPU,如下所示:
class _RunInferenceDoFn(beam.DoFn, Generic[ExampleT, PredictionT]):
...
def _get_semaphore(self):
def get_semaphore():
logging.info(\'intializing semaphore...\')
return Semaphore(1)
return self._shared_semaphore.acquire(get_semaphore)
def setup(self):
...
self._model = self._load_model()
self._semaphore = self._get_semaphore()
def process(self, batch, inference_args):
...
logging.info(\'trying to acquire semaphore...\')
self._semaphore.acquire()
logging.info(\'semaphore acquired\')
start_time = _to_microseconds(self._clock.time_ns())
result_generator = self._model_handler.run_inference(
batch, self._model, inference_args)
end_time = _to_microseconds(self._clock.time_ns())
self._semaphore.release()
...
我们在该设置中做了三个奇怪的观察:
- Beam 始终使用我们允许的最小批量大小;如果我们指定最小 8 最大 32 的批处理大小,它总是会选择最多 8 的批处理大小,有时会更低。
- 这里的推理时间在允许多线程 (
--number_of_worker_harness_threads=10) 时仍然比单线程 (--number_of_worker_harness_threads=1) 慢得多。每批 2.7 秒与每批 0.4 秒,两者都比直接在计算引擎上运行要慢一些。 - 在多线程设置中,尽管使用了保守的批量大小,但我们仍会偶尔看到 CUDA OOM 错误。
将不胜感激任何和所有调试指导如何使这项工作!现在,整个管道太慢了,以至于我们不得不再次在 Compute Engine 上批量运行:/ – 但必须有一种方法可以在 Dataflow 上完成这项工作,对吧?
以供参考:
- 单线程作业:
catalin-debug-classifier-test-1660143139 (Job ID: 2022-08-10_07_53_06-5898402459767488826) - 多线程作业:
catalin-debug-classifier-10threads-32batch-1660156741 (Job ID: 2022-08-10_11_39_50-2452382118954657386)
- 单线程作业:
标签: python google-cloud-platform google-cloud-dataflow apache-beam