随机性与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 Q 和 var 静止图像中没有任何内容等于 0。这是预期的。
调用qr.create_threads 并设置start=True 后,qr 执行以下操作:
- 开始使用您的
enqueue_ops 填充Q,在本例中为en_q。它将尝试打包尽可能多的en_q 操作,这取决于您的Q 大小。
- 为
Q 中的每个en_q 创建线程。
- 并发运行线程。 (此时,鉴于
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