【问题标题】:In-graph replication for Distributed Tensorflow分布式 TensorFlow 的图内复制
【发布时间】:2017-08-19 13:23:53
【问题描述】:

我正在学习分布式Tensorflow,我实现了In-graph replication的简单版本代码如下(task_parallel.py):

import argparse
import logging

import tensorflow as tf


log = logging.getLogger(__name__)

# Job Names
PARAMETER_SERVER = "ps"
WORKER_SERVER = "worker"

# Cluster Details
CLUSTER_SPEC = {
    PARAMETER_SERVER: ["localhost:2222"],
    WORKER_SERVER: ["localhost:1111", "localhost:1112", "localhost:1113"]}


def parse_command_arguments():
    """ Set up and parse the command line arguments passed for experiment. """
    parser = argparse.ArgumentParser(
        description="Parameters and Arguments for the Test.")

    parser.add_argument(
        "--ps_hosts",
        type=str,
        default="",
        help="Comma-separated list of hostname:port pairs"
    )
    parser.add_argument(
        "--worker_hosts",
        type=str,
        default="",
        help="Comma-separated list of hostname:port pairs"
    )
    parser.add_argument(
        "--job_name",
        type=str,
        default="",
        help="One of 'ps', 'worker'"
    )
    # Flags for defining the tf.train.Server
    parser.add_argument(
        "--task_index",
        type=int,
        default=0,
        help="Index of task within the job"
    )

    return parser.parse_args()


def start_server(
        job_name, ps_hosts, task_index, worker_hosts):
    """ Create a server based on a cluster spec. """
    cluster_spec = {
        PARAMETER_SERVER: ps_hosts,
        WORKER_SERVER: worker_hosts}
    cluster = tf.train.ClusterSpec(cluster_spec)

    server = tf.train.Server(
        cluster, job_name=job_name, task_index=task_index)

    return server


def model():
    """ Build up a simple estimator model. """
    with tf.device("/job:%s/task:0" % PARAMETER_SERVER):
        log.info("111")
        # Build a linear model and predict values
        W = tf.Variable([.3], tf.float32)
        b = tf.Variable([-.3], tf.float32)
        x = tf.placeholder(tf.float32)
        linear_model = W * x + b
        y = tf.placeholder(tf.float32)
        global_step = tf.Variable(0)

    with tf.device("/job:%s/task:0" % WORKER_SERVER):
        # Loss sub-graph
        loss = tf.reduce_sum(tf.square(linear_model - y))
        log.info("222")
        # optimizer
        optimizer = tf.train.GradientDescentOptimizer(0.01)

    with tf.device("/job:%s/task:1" % WORKER_SERVER):
        log.info("333")
        train = optimizer.minimize(loss, global_step=global_step)

    return W, b, loss, x, y, train, global_step


def main():
    # Parse arguments from command line.
    arguments = parse_command_arguments()

    # Initializing logging with level "INFO".
    logging.basicConfig(level=logging.INFO)

    ps_hosts = arguments.ps_hosts.split(",")
    worker_hosts = arguments.worker_hosts.split(",")
    job_name = arguments.job_name
    task_index = arguments.task_index

    # Start a server.
    server = start_server(
        job_name, ps_hosts, task_index, worker_hosts)

    W, b, loss, x, y, train, global_step = model()
    # with sv.prepare_or_wait_for_session(server.target) as sess:
    with tf.train.MonitoredTrainingSession(
            master=server.target,
            is_chief=(arguments.task_index == 0 and (
                        arguments.job_name == 'ps')),
            config=tf.ConfigProto(log_device_placement=True)) as sess:
        step = 0
        # training data
        x_train = [1, 2, 3, 4]
        y_train = [0, -1, -2, -3]
        while not sess.should_stop() and step < 1000:
            _, step = sess.run(
                [train, global_step], {x: x_train, y: y_train})

        # evaluate training accuracy
        curr_W, curr_b, curr_loss = sess.run(
            [W, b, loss], {x: x_train, y: y_train})
        print("W: %s b: %s loss: %s" % (curr_W, curr_b, curr_loss))

if __name__ == "__main__":
    main()

我在一台机器上用 3 个不同的进程运行代码(只有 CPU 的 MacPro):

PS:$python task_parallel.py --task_index 0 --ps_hosts localhost:2222 --worker_hosts localhost:1111,localhost:1112 --job_name ps

工人 1:$python task_parallel.py --task_index 0 --ps_hosts localhost:2222 --worker_hosts localhost:1111,localhost:1112 --job_name worker

工人 2:$python task_parallel.py --task_index 1 --ps_hosts localhost:2222 --worker_hosts localhost:1111,localhost:1112 --job_name worker

我注意到结果不是我所期望的。具体来说,我希望进程“PS”只打印111,“Worker 1”只打印222,而“Worker 3”只打印333,因为我为每个进程指定了任务。然而,我得到的是所有 3 个进程都打印了完全相同的东西:

INFO:__main__:111
INFO:__main__:222
INFO:__main__:333

进程PS 不是只执行块with tf.device("/job:%s/task:0" % PARAMETER_SERVER 内的代码吗?工人也一样?我想知道我是否遗漏了代码中的某些内容。

我还发现我必须先运行所有工作进程,然后再运行 ps 进程。否则,训练完成后,工作进程无法正常退出。所以我想知道我的代码中这个问题的任何原因。非常感谢您的帮助:) 谢谢!

【问题讨论】:

    标签: python-3.x tensorflow distributed


    【解决方案1】:

    请注意,在您的 sn-p 中,MonitoredTrainingSession 之前的代码用于描述和构建运行图,参数服务器和工作人员都会执行这些代码来生成图。在创建 MonitoredTrainingSession 时,图表将被冻结。

    如果您只想在PS 中看到111,您的代码可能会这样工作:

    FLAGS = tf.app.flags.FLAGS
    if FLAGS.job_name == 'ps':
        print('111')
        server.join()
    else:
        print('222')
    

    如果你想在workers中设置replicas模型,在model()函数中:

    with tf.device('/job:ps/task:0'):
        # define variable in parameter
    with tf.device('/job:worker/task:%d' % FLAGS.task_index):
        # define model in worker % task_index
    

    此外,replica_device_setter 将在构造 Operation 对象时自动将设备分配给它们。

    有一些tensorflow提供的例子,比如:

    1. hello distributed,tensorflow教程中的基础指南。

    2. mnist_replica.py,分布式 MNIST 训练和验证,带有模型副本。

    3. cifar10_multi_gpu_train.py,一个二进制文件,用于使用具有同步更新的多个 GPU 训练 CIFAR-10。

    希望这对您有所帮助。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多