【问题标题】:Can't understand the result of TensorOverflow train.QueueRunner无法理解 TensorOverflow train.QueueRunner 的结果
【发布时间】:2019-09-03 04:04:50
【问题描述】:

抱歉,我是 TensorFlow 的初学者。下面只是一个 TensorFlow 的代码,它使用两个线程独立地入队和出队。

import tensorflow as tf

Q = tf.compat.v1.FIFOQueue(1000, tf.float32)
var = tf.Variable(0.0)
data = tf.compat.v1.assign_add(var, tf.constant(1.0))
en_q = Q.enqueue(data)
qr = tf.train.QueueRunner(Q, enqueue_ops=[en_q])
init_op = tf.compat.v1.global_variables_initializer()
with tf.compat.v1.Session() as sess:
    sess.run(init_op)
    coord = tf.train.Coordinator()
    threads = qr.create_threads(sess, coord=coord, start=True)
    for i in range(300):
        print(sess.run(Q.dequeue()))
    coord.request_stop()
    coord.join(threads)

结果如下:

3.0
7.0
10.0
14.0
18.0
21.0
26.0
....

我对这个结果感到很困惑。由于Q 是一个先进先出队列,即使出队和入队在两个不同的线程中,入队的数量仍然应该是1,2,3,4,5,6,...。为什么出队的号码有可能是3,7,10,14,....1, 2, 4, 5, ...的号码在哪里?

【问题讨论】:

    标签: python multithreading tensorflow queue


    【解决方案1】:

    随机性与FIFOqueue 无关,而是线程执行并发的预期行为。肯定不止两个线程。

    为了简单起见,让我们尝试在没有for 循环的情况下运行您的代码,添加一些print 语句和sleep 间隔,看看会发生什么:

    import tensorflow as tf
    import time
    
    Q_size = 1000 # Q size
    t = 1 # t in sec
    
    Q = tf.compat.v1.FIFOQueue(Q_size, tf.float32)
    var = tf.Variable(0.0)
    data = tf.compat.v1.assign_add(var, tf.constant(1.0))
    en_q = Q.enqueue(data)
    
    qr = tf.train.QueueRunner(Q, enqueue_ops=[en_q])
    
    init_op = tf.compat.v1.global_variables_initializer()
    
    with tf.compat.v1.Session() as sess:
        sess.run(init_op)
        coord = tf.train.Coordinator()
    
        print('*** before qr.create_threads ***')
        print('Q.size()', sess.run(Q.size()))
        print('var', var.eval())
    
        threads = qr.create_threads(sess, coord=coord, start=True)
    
        print('*** after qr.create_threads ***')
        time.sleep(t) # sleep t sec in main to wait for all threads to finish running.
        print('Q.size()', sess.run(Q.size()))
        print('var', var.eval())
    
        coord.request_stop()
        coord.join(threads)    
    

    输出:

    *** before qr.create_threads ***
    Q.size() 0
    var 0.0
    
    *** after qr.create_threads ***
    Q.size() 1000
    var 1001.0
    

    在调用qr 之前,没有任何反应。 FIFOQueue Qvar 静止图像中没有任何内容等于 0。这是预期的。

    调用qr.create_threads 并设置start=True 后,qr 执行以下操作:

    1. 开始使用您的enqueue_ops 填充Q,在本例中为en_q。它将尝试打包尽可能多的en_q 操作,这取决于您的Q 大小。
    2. Q 中的每个en_q 创建线程。
    3. 并发运行线程。 (此时,鉴于Q 的大小,很明显不能只涉及 2 个线程,如 OP 所述。)

    sleep 的引入是为了让所有这些线程在我们显示输出之前完成执行。我们想知道qr创建的所有线程运行完毕后Q的大小和var的值是多少。

    正如预期的那样,Q 填充了 1000 个 en_q 操作,var 被添加了 1001 次。

    现在,如果我们删除 sleep 间隔,我们会得到如下所示的随机输出。

    没有sleep的随机输出:

    *** before qr.create_threads ***
    Q.size() 0
    var 0.0
    
    *** after qr.create_threads ***
    Q.size() 32
    var 35.0
    

    上面随机输出的意思是,当我们在主线程中显示输出时,qr 仍在填充Q,在后台同时创建和运行线程。

    现在,让我们把你的 for 循环放回去,同时在循环中引入另一个 sleep 间隔:

    import tensorflow as tf
    import time
    
    Q_size = 1000 # Q size
    N =10 # range of for loop
    t = 1 # t in secs
    
    Q = tf.compat.v1.FIFOQueue(Q_size, tf.float32)
    var = tf.Variable(0.0)
    data = tf.compat.v1.assign_add(var, tf.constant(1.0))
    en_q = Q.enqueue(data)
    
    qr = tf.train.QueueRunner(Q, enqueue_ops=[en_q])
    
    init_op = tf.compat.v1.global_variables_initializer()
    
    with tf.compat.v1.Session() as sess:
        sess.run(init_op)
        coord = tf.train.Coordinator()
    
        print('*** before qr.create_threads ***')
        print('Q.size()', sess.run(Q.size()))
        print('var', var.eval())
    
        threads = qr.create_threads(sess, coord=coord, start=True)
    
        print('*** after qr.create_threads ***')
        time.sleep(t) # sleep t sec in main to wait for all threads to finish running.
        print('Q.size()', sess.run(Q.size()))
        print('var', var.eval())
    
        print('*** for loop ***')
        for i in range(N):
            sess.run(Q.dequeue())
            time.sleep(t) # sleep t sec in main to wait for all threads to finish running.
            print('var', var.eval())
            print('.size()', sess.run(Q.size()))
    
        coord.request_stop()
        coord.join(threads)    
    

    引入sleep 允许后台线程在我们在主线程中显示变量之前完成。

    您可以在下面的输出中看到,var 现在按顺序增加,而Q 的大小保持在 1000。每个 sess.run(Q.dequeue()) 调用现在从 Q 执行一个 en_q 操作,将 var 加 1。

    我希望这可以澄清您观察到的随机性与 Q 是否为 FIFOqueue 无关。随机性是由线程并行的预期行为引起的。

    输出:

    *** before qr.create_threads ***
    Q.size() 0
    var 0.0
    
    *** after qr.create_threads ***
    Q.size() 1000
    var 1001.0
    
    *** for loop ***
    var 1002.0
    .size() 1000
    var 1003.0
    .size() 1000
    var 1004.0
    .size() 1000
    var 1005.0
    .size() 1000
    var 1006.0
    .size() 1000
    var 1007.0
    .size() 1000
    var 1008.0
    .size() 1000
    var 1009.0
    .size() 1000
    var 1010.0
    .size() 1000
    var 1011.0
    .size() 1000
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-08-20
      • 1970-01-01
      • 2022-09-21
      • 1970-01-01
      • 1970-01-01
      • 2016-06-25
      • 2014-01-21
      相关资源
      最近更新 更多