背景
機器學習作業負載與傳統的作業負載相比,一個比較顯著的特點是對 GPU 的需求旺盛,在之前的文章中介紹過(https://mp.weixin.qq.com/s/Nasm-cXLtJObjLwLQHALmw 和 https://mp.weixin.qq.com/s/X4VDynLfKdVp-tyciQccyQ),目前 GPU 的顯存已經不足以跟上模型引數規模的發展,隨著 Transformer 等新的模型結構的出現,這一問題越來越顯著,演算法工程師們訓練模型所需要的資源越來越多,分布式訓練也隨之成為了工業界進行模型訓練的標準方式,
彈性訓練能夠在訓練程序中動態地調整參與訓練的實體數量,極大程度提高集群資源的利用率,同時,配合云上的競價實體等資源型別,能夠以更低的成本進行模型調優,進一步降本增效,在 PyTorch 最新發布的 1.9.0 版本中,其原本分布式訓練的方式torch.distributed.launch 即將被廢棄,轉而推薦用戶使用彈性的分布式訓練介面 torch.distributed.run,
借此機會,我們對這一新特性進行簡單地介紹,并且與 Horovod Elastic 進行簡單地對比和分析,最后總結一下使用彈性訓練時,需要注意的問題,
PyTorch 1.9.0 之前的設計
PyTorch 是目前最流行的深度學習框架之一,它最讓人稱道的是易用性,無論是單機訓練還是分布式訓練,PyTorch 都提供了簡潔的 API,
PyTorch 1.9.0 版本之前,分布式訓練的方式通常是通過如下的方式進行,
python -m torch.distributed.launch
--nnodes=NODE_SIZE
--nproc_per_node=TRAINERS_PER_NODE
--node_rank=NODE_RANK
--master_port=HOST_PORT
--master_addr=HOST_NODE_ADDR
YOUR_TRAINING_SCRIPT.py (--arg1 ... train script args...)
其中 nnodes 是參與訓練的節點個數,nproc_per_node 是每個節點上運行的行程數量,node_rank 是當前節點的識別符號,master_addr 和 master_port 是 master 監聽的地址和埠,torch.distributed.launch 會設定一些環境變數,其中包括 WORLD_SIZE 和 MASTER_PORT、MASTER_ADDR 等,
隨后在當前機器上會創建對應行程進行訓練,當前機器會有 TRAINERS_PER_NODE 個行程,這些行程組成了一個 local worker group,一共有 NODE_SIZE 個機器參與訓練,一共有 NODE_SIZE * TRAINERS_PER_NODE 個行程,如果想要發起一個分布式訓練任務,需要在所有的機器上執行相應的命令,
PyTorch 1.9.0 中的新設計
在 PyTorch 1.9 中,torch.distributed.launch 即將被廢棄,取而代之的是基于 pytorch/elastic 的 torch.distributed.run,這一新的方式與之前相比有一些使用上的改動,如下所示,
python -m torch.distributed.run
--nnodes=MIN_SIZE:MAX_SIZE
--nproc_per_node=TRAINERS_PER_NODE
--rdzv_id=JOB_ID
--rdzv_backend=c10d
--rdzv_endpoint=HOST_NODE_ADDR
YOUR_TRAINING_SCRIPT.py (--arg1 ... train script args...)
它提供了一些新的能力:首先是更好的容錯,當 worker 失敗后會自動重啟繼續訓練;其次是 RANK 和 WORLD_SIZE 這些欄位不再需要手動設定,最后也是最重要的,支持彈性訓練,動態地增加或減少參與訓練的 worker 數量,在上面的例子中,nnodes 的設定不再是一個固定的值,而是一個區間,訓練任務可以容忍在這一區間范圍內的 worker 數量變化,
如果要支持彈性能力,訓練代碼也需要進行一些修改,
def main():
args = parse_args(sys.argv[1:])
state = load_checkpoint(args.checkpoint_path)
initialize(state)
# torch.distributed.run ensure that this will work
# by exporting all the env vars needed to initialize the process group
torch.distributed.init_process_group(backend=args.backend)
for i in range(state.epoch, state.total_num_epochs)
for batch in iter(state.dataset)
train(batch, state.model)
state.epoch += 1
save_checkpoint(state)
其中比較明顯的變化是,用戶需要手動地處理 checkpoint,這是因為當 worker 出現失效時,所有的 worker 都會重啟,所以需要 checkpoint 機制來保證重啟后訓練能夠繼續下去,這一新的分布式訓練方式引入不少新的概念,包括 agent、rendezvous 等,接下來我們自用戶能接觸到的 torch.distributed.run 開始,介紹這些新的設計,
def run(args):
if args.standalone:
args.rdzv_backend = "c10d"
args.rdzv_endpoint = "localhost:29400"
args.rdzv_id = str(uuid.uuid4())
log.info(
f"\n**************************************\n"
f"Rendezvous info:\n"
f"--rdzv_backend={args.rdzv_backend} "
f"--rdzv_endpoint={args.rdzv_endpoint} "
f"--rdzv_id={args.rdzv_id}\n"
f"**************************************\n"
)
config, cmd, cmd_args = config_from_args(args)
elastic_launch(
config=config,
entrypoint=cmd,
)(*cmd_args)
其中主要區分了兩個模式,Standalone 模式和分布式模式,Standalone 模式是分布式模式的一種特例,它主要針對單機多 Worker 的方式提供了一些便利的設定,不再需要設定一些多余的引數如 rdzv_backend 和 rdzv_endpoint 等,
兩者最后都會通過 elastic_launch 發起真正的訓練行程,elastic_launch 會通過 elastic agent 來管理 worker 的生命周期,它的回傳是每個 worker 的輸出,
class elastic_launch:
...
def __call__(self, *args):
return launch_agent(self._config, self._entrypoint, list(args))
def launch_agent(
config: LaunchConfig,
entrypoint: Union[Callable, str, None],
args: List[Any],
) -> Dict[int, Any]:
...
agent = LocalElasticAgent(
spec=spec, start_method=config.start_method, log_dir=config.log_dir
)
...
result = agent.run()
...
return result.return_values
Elastic Agent 的設計:如何管理多個 worker 行程
elastic agent 是一個獨立的行程,負責管理其下的 workers,它起到了類似行程管理系統 supervisor 的作用,會在啟動的時候確保每個 worker 的設定正確,由于有關 WORLD_SIZE 和 RANK 的資訊不再需要用戶提供,elastic agent 會負責處理,
除此之外,worker 的失效也是由 elastic agent 負責捕獲處理,可以說 elastic agent 是彈性訓練中最核心的抽象概念,

上圖展示的是elastic agent 的作業原理,
不同的 elastic agent 之間通過 rendezvous 進行 worker 之間的相互發現和對成員變動的同步,與此同時,通過對 worker 行程的監控,來捕獲訓練程序中的失效,其中核心的邏輯都包裝在 LocalElasticAgent.run() 這一函式呼叫中,
def run(self, role: str = DEFAULT_ROLE) -> RunResult:
...
result = self._invoke_run(role)
return result
def _invoke_run(self, role: str = DEFAULT_ROLE) -> RunResult:
...
self._initialize_workers(self._worker_group)
while True:
...
run_result = self._monitor_workers(self._worker_group)
state = run_result.state
...
if state == WorkerState.SUCCEEDED:
...
return run_result
elif state in {WorkerState.UNHEALTHY, WorkerState.FAILED}:
if self._remaining_restarts > 0:
...
self._restart_workers(self._worker_group)
else:
...
return run_result
elif state == WorkerState.HEALTHY:
...
if num_nodes_waiting > 0:
...
self._restart_workers(self._worker_group)
else:
raise Exception(f"[{role}] Worker group in {state.name} state")
可以看到,核心的邏輯在 _invoke_run 中,其中 _initialize_workers 執行了大部分初始化的作業,其中包括為每個 worker 分配 RANK 等,在默認的實作中 elastic agent 和 workers 行程在同一機器上,因此 self._monitor_workers(self._worker_group) 通過 multiprocessing 對 workers 的運行狀態進行了監控,并且根據不同的狀態,進行不同的處理,
elastic agent 的可擴展性非常好,在 1.9.0 版本中,一共有三個 Agent,分別是 ElasticAgent、SimpleElasticAgent 和 LocalElasticAgent,
其中 ElasticAgent 是一個 Abstract Class,SimpleElasticAgent 對其中的某些函式進行了實作,而 LocalElasticAgent 則實作了管理單機上所有 worker 行程的 elastic agent,
SimpleElasticAgent 這一個抽象主要是為了方便擴展新的 agent 實作,比如如果你想通過一個 agent 管理多機上所有的 worker,而不只是本機上的 worker,則可以通過擴展 SimpleElasticAgent 來實作,
rendezvous 的設計:如何在不同的節點間確定 RANK
接下來,我們再看另外一個核心的抽象 rendezvous,為了實作彈性訓練,worker 之間要能夠動態地進行 membership 的變更,rendezvous 就是實作這一特性的用于同步的組件,
rendezvous 最核心的方法是:
@abstractmethod
def next_rendezvous(
self,
) -> Tuple[Store, int, int]:
"""Main entry-point into the rendezvous barrier.
Blocks until the rendezvous is complete and the current process is
included in the formed worker group, or a timeout occurs, or the
rendezvous was marked closed.
Returns:
A tuple of :py:class:`torch.distributed.Store`, ``rank``, and
``world size``.
Raises:
RendezvousClosedError:
The rendezvous is closed.
RendezvousConnectionError:
The connection to the rendezvous backend has failed.
RendezvousStateError:
The rendezvous state is corrupt.
RendezvousTimeoutError:
The rendezvous did not complete on time.
"""
如注釋所示,這一函式呼叫會被阻塞,直到 worker 的數量達到了要求,在 worker 被初始化,或者重啟的時候,這一函式都會被呼叫,當函式回傳時,不同的 worker 會以回傳中的 rank 作為唯一的標示,rendezvous 一共有四個實作,分別是 etcd、etcd-v2、c10d 和 static,
class EtcdRendezvousHandler(RendezvousHandler):
def next_rendezvous(self):
rdzv_version, rank, world_size = self._rdzv_impl.rendezvous_barrier()
log.info("Creating EtcdStore as the c10d::Store implementation")
store = self._rdzv_impl.setup_kv_store(rdzv_version)
return store, rank, world_size
其中 etcd 相關的是之前推薦使用的實作,在 c10d 出現后就不再推薦了,etcd 的實作中,不同 worker 之間的狀態通過 etcd 的 KV 介面存盤,
確定參與訓練的實體和對應的 RANK 的程序如下圖所示,

首先會在 /rdzv/active_version 下嘗試寫一個值 status: setup,在整個程序中,/rdzv/active_version 會作為存盤 rendezvous 程序中間狀態的 KV store,以及 rendezvous 程序中的排他鎖來使用,
如果寫失敗了,說明目前已經有對應的 rendezvous 程序正在進行中,
在成功后,會更新 /rdzv/version_counter 為原值加一,然后會創建一個目錄 /rdzv/v_${version_counter},這些操作做完后,會將 /rdzv/active_version 的狀態寫為 joinable,這時就進入了 join 階段,
在 join 階段,不同的 agent 在鎖的保護下,會依次更新 /rdzv/active_version 下的 paticipants,分配到遞增的 rank,這里的 rank 并不是每個 worker 行程分配到的 global rank,而是 agent 自己的 rank,worker 行程的 rank 會根據 agent rank 經過一定的計算得到,這也是一個非常容易混淆的設計,竊以為有優化的空間,
def init_phase(self):
try:
active_version = self.try_create_rendezvous()
state = json.loads(active_version.value)
log.info("New rendezvous state created: " + str(state))
except etcd.EtcdAlreadyExist:
# 已經有了一個新的 rendezvous 程序
active_version, state = self.get_rdzv_state()
# Note: it is possible for above query to fail (etcd.EtcdKeyNotFound),
# but this is ok for us - just means we'll restart from beginning.
log.info("Observed existing rendezvous state: " + str(state))
if state["status"] == "closed":
raise RendezvousClosedError()
if state["status"] == "joinable":
return self.join_phase(state["version"])
if state["status"] == "final":
self.handle_existing_rendezvous(state["version"])
raise EtcdRendezvousRetryImmediately()
self.try_wait_for_state_change(etcd_index=active_version.etcd_index + 1)
raise EtcdRendezvousRetryableFailure()
在參與訓練的節點達到 nnodes 的命令列引數中傳入的最小值時,會等待一定時間,在等待時間結束或者參與訓練的節點達到了 nnodes 設定的最大值時,會進入 frozen 階段,
在 fronzen 階段中,每個參與訓練的節點都需要通過在 /rdzv/v_${version_counter}/rank_${agent_rank} 下寫值的方式進行確認,在所有節點都確認完畢后,會進入最后的 final 階段,
在最后的 final 階段中,后續進入的 agent 都會 pending,已經達成 rendezvous 的節點上的 agent 會為其管理的 worker 行程分配 RANK,RANK 0 的實體會作為 master 的角色存在,隨后就會直接創建對應的 worker 行程,在默認的 LocalElasticAgent 中,會利用 python.multiprocessing 在本地創建多個行程,
@prof
def _start_workers(self, worker_group: WorkerGroup) -> Dict[int, Any]:
spec = worker_group.spec
store = worker_group.store
...
for worker in worker_group.workers:
local_rank = worker.local_rank
worker_env = {
"LOCAL_RANK": str(local_rank),
"RANK": str(worker.global_rank),
...
}
...
args[local_rank] = tuple(worker_args)
...
self._pcontext = start_processes(
name=spec.role,
entrypoint=spec.entrypoint,
args=args,
envs=envs,
log_dir=attempt_log_dir,
start_method=self._start_method,
redirects=spec.redirects,
tee=spec.tee,
)
return self._pcontext.pids()
c10d 新的設計
前文介紹了基于 etcd 的 rendezvous 實作,它可以保證多個實體之間對于參與訓練的節點共識的強一致,但是這也為 PyTorch 運行訓練任務引入了額外的依賴,因此 PyTorch 也提供了一個內置的實作 c10d,相比于基于 etcd 的實作,c10d 基于 TCP 來進行同步,
def create_backend(params: RendezvousParameters) -> Tuple[C10dRendezvousBackend, Store]:
...
if store_type == "file":
store = _create_file_store(params)
elif store_type == "tcp":
store = _create_tcp_store(params)
...
backend = C10dRendezvousBackend(store, params.run_id)
def _create_tcp_store(params: RendezvousParameters) -> TCPStore:
host, port = parse_rendezvous_endpoint(params.endpoint, default_port=29400)
...
for is_server in [is_host, False]:
...
store = TCPStore(
host, port, is_master=is_server, timeout=timedelta(seconds=read_timeout)
)
...
break
return store
c10d 是一個 client-server 的架構,其中的一個 agent 上會運行 c10d 的 TCPServer,它監聽給定的埠,提供了 compareAndSet、add 等原語,它也可以被理解為一個簡化的,提供 KV 介面的記憶體資料庫,類似于 Redis,有關 rendezvous 的同步,都是由各個 agent 通過一個中心化的 agent 上的 c10d TCPServer 完成的,可以預見這樣的實作在可用性上相比于 etcd 是有一定差距的,但是勝在易用性,用戶如果使用 c10d,那么不再需要運維一個 etcd 集群,
PyTorch Elastic on Kubernetes
為了能夠享受到彈性訓練帶來的便利,PyTorch 同時提供了在 Kubernetes 上的支持,相比于 1.9.0 之前的版本,新版本的分布式訓練添加了一些新的引數,因此 PyTorch 社區在 Kubeflow PyTorch operator 的基礎上,對 CRD 進行了一些修改,一個典型的彈性訓練示例如下所示:
apiVersion: elastic.pytorch.org/v1alpha1
kind: ElasticJob
metadata:
name: imagenet
namespace: elastic-job
spec:
# Use "etcd-service:2379" if you already apply etcd.yaml
rdzvEndpoint: "<your_etcd_endpoint>:<your_etcd_port>"
minReplicas: 1
maxReplicas: 2
replicaSpecs:
Worker:
replicas: 2
restartPolicy: ExitCode
template:
apiVersion: v1
kind: Pod
spec:
containers:
- name: elasticjob-worker
image: torchelastic/examples:0.2.0
imagePullPolicy: Always
args:
- "--nproc_per_node=1"
- "/workspace/examples/imagenet/main.py"
- "--arch=resnet18"
- "--epochs=20"
- "--batch-size=32"
# number of data loader workers (NOT trainers)
# zero means load the data on the same process as the trainer
# this is set so that the container does not OOM since
# pytorch data loaders use shm
- "--workers=0"
- "/workspace/data/tiny-imagenet-200"
resources:
limits:
nvidia.com/gpu: 1
由于在最開始,基于 c10d 的 rendezvous 還沒有被支持,所以 CRD 中需要定義 rdzvEndpoint,指向一個已經部署好的 etcd 集群,同時,用戶需要指定 minReplicas 和 maxReplicas,其他就與 Kubeflow PyTorchJob 并無二致,
PyTorch Elastic 與 Horovod Elastic
目前,兩者的設計從原理上來說并無二致,相比于 Horovod Elastic,PyTorch Elastic 提供了更靈活的擴展性,它提供了 agent、rendezvous 等介面,用戶可以根據需要進行擴展,但是從另外一個角度講,Horovod 的易用性做的更好,
PyTorch 并沒有提供保存狀態的內置支持,為了能夠在 worker 行程失敗,重建訓練任務的時候,需要用戶自己實作保存會加載 checkpoint 的邏輯;而 Horovod 則提供了內置的實作,
Horovod 和 PyTorch 在同步機制上也具有比較大的差異,Horovod Elastic 需要用戶提供一個腳本 discovery_hosts.sh,幫助其在運行時獲得正在參與訓練的節點,
$ horovodrun -np 8 --host-discovery-script discover_hosts.sh python train.py
...
$ ./discover_hosts.sh
host-1:29500
host-2:29500
host-3:29500
這相當于將節點發現的邏輯交給用戶來實作,反觀 PyTorch,它利用 etcd、自身實作的 c10d 等組件解決節點間的相互發現問題,顯得更為精巧,
總結
在文章最后,我們總結一下目前實作彈性訓練需要注意的問題,
首先,也是最重要的,彈性訓練需要一種機制來解決節點/訓練行程間相互發現的問題,訓練程序中節點會動態地加入或者退出,如何讓其他的節點感知到這一變化,是這一機制主要面對的問題,目前的設計中,Horovod 將這一問題交給用戶來解決,Horovod 定期執行用戶定義的邏輯來發現目前的節點,PyTorch 通過第三方的分布式一致性中間件 etcd 等來實作高可用的節點發現,除此之外,也有一些探索性的作業,利用基于 Gossip 的協議來進行同步,在兼顧高可用的同時也沒有引入過多的組件,
其次,要實作彈性訓練還需要捕獲訓練失效,Horovod 和 PyTorch 都通過一個后臺行程(Horovod 中是 Driver,PyTorch 中是每個節點的 Local Elastic Agent)來實作這一邏輯,當行程 crash,或在梯度通信中遇到問題時,后臺行程會捕獲到失效并且重新進行節點發現,然后重啟訓練,
最后,訓練時的資料切分的邏輯和學習率/ batch size 的設定也要對應進行修改,由于參與訓練的行程會動態的增減,因此可能需要根據新的訓練行程的規模來重新設定學習率和資料分配的邏輯,避免影響模型收斂,
在本文中,我們首先介紹了 PyTorch 1.9.0 版本中彈性訓練的設計與實作,然后分析總結了實作彈性訓練的方式和不同框架之間的設計差異,從我們的角度來看,彈性訓練能夠很好地貼合云原生的趨勢,以極致的彈性來降低成本提高資源利用率,是未來的趨勢,因此目前我們也在積極參與 TensorFlow、PyTorch 和 Kubeflow 等社區的彈性訓練的社區貢獻作業,后續會發布更多的相關文章,感謝關注,
【騰訊云原生】云說新品、云研新術、云游新活、云賞資訊,掃碼關注同名公眾號,及時獲取更多干貨!!
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/296263.html
標籤:其他

