fix the diff between async mode and async_half mode (#19535)
* test=develop, communicator merge add => merge averagesigmoid_bug
parent
e9233d1c1e
commit
2f037c3189
@ -0,0 +1,55 @@
|
||||
# Copyright (c) 2019 PaddlePaddle Authors. All Rights Reserved.
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
from paddle.fluid.incubate.fleet.parameter_server.distribute_transpiler import fleet
|
||||
from paddle.fluid.contrib.utils import HDFSClient
|
||||
import os
|
||||
|
||||
|
||||
def check_all_trainers_ready(ready_path, epoch):
|
||||
trainer_num = fleet.worker_num()
|
||||
trainer_id = fleet.worker_index()
|
||||
|
||||
hadoop_home = os.getenv("HADOOP_HOME")
|
||||
configs = {
|
||||
"fs.default.name": os.getenv("FS_NAME"),
|
||||
"hadoop.job.ugi": os.getenv("FS_UGI")
|
||||
}
|
||||
|
||||
node_ready = "ready.{}.{}.done".format(epoch, trainer_id)
|
||||
|
||||
with open(node_ready, "w") as node:
|
||||
node.write("")
|
||||
|
||||
client = HDFSClient(hadoop_home, configs)
|
||||
if not client.is_dir(ready_path):
|
||||
client.makedirs(ready_path)
|
||||
client.upload(
|
||||
hdfs_path=ready_path,
|
||||
local_path=node_ready,
|
||||
overwrite=True,
|
||||
retry_times=0)
|
||||
|
||||
print("PUT {} ON HDFS {} OK".format(node_ready, ready_path))
|
||||
|
||||
while True:
|
||||
ready_num = len(client.ls(ready_path))
|
||||
print("have {} trainers need to be ready".format(trainer_num - ready_num
|
||||
% trainer_num))
|
||||
if ready_num % trainer_num == 0:
|
||||
break
|
||||
time.sleep(10)
|
||||
ready_num = len(client.ls(ready_path))
|
||||
|
||||
print("All trainers are ready, continue training")
|
@ -0,0 +1,35 @@
|
||||
# Copyright (c) 2018 PaddlePaddle Authors. All Rights Reserved.
|
||||
#
|
||||
# Licensed under the Apache License, Version 2.0 (the "License");
|
||||
# you may not use this file except in compliance with the License.
|
||||
# You may obtain a copy of the License at
|
||||
#
|
||||
# http://www.apache.org/licenses/LICENSE-2.0
|
||||
#
|
||||
# Unless required by applicable law or agreed to in writing, software
|
||||
# distributed under the License is distributed on an "AS IS" BASIS,
|
||||
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
# See the License for the specific language governing permissions and
|
||||
# limitations under the License.
|
||||
|
||||
from __future__ import print_function
|
||||
import paddle.fluid as fluid
|
||||
import unittest
|
||||
import paddle.fluid.incubate.fleet.base.role_maker as role_maker
|
||||
from paddle.fluid.incubate.fleet.parameter_server.distribute_transpiler import fleet
|
||||
from paddle.fluid.incubate.fleet.utils.fleet_barrier_util import check_all_trainers_ready
|
||||
|
||||
|
||||
class TestFleetUtils(unittest.TestCase):
|
||||
def test_fleet_barrier(self):
|
||||
role = role_maker.UserDefinedRoleMaker(
|
||||
current_id=0,
|
||||
role=role_maker.Role.WORKER,
|
||||
worker_num=1,
|
||||
server_endpoints=['127.0.0.1'])
|
||||
fleet.init(role)
|
||||
check_all_trainers_ready("/ready_path/", 0)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
unittest.main()
|
Loading…
Reference in new issue