我正在尝试对 fit_generator 进行多处理。
这些是我面临的问题。
trainable_model.fit_generator(load_random_cached_bottlenecks(BATCH_SIZE, label_map, training_addr_label_map, train_npy_dir, 'h5py', h5py_file_train),epochs = EPOCHS, steps_per_epoch=iterations_per_epoch_t, validation_data = load_random_cached_bottlenecks(BATCH_SIZE, label_map, validation_addr_label_map, val_npy_dir, 'h5py', h5py_file_val), validation_steps=iterations_per_epoch_v, workers = 1, callbacks = callback_list, use_multiprocessing = True, max_queue_size = 32)
Run Code Online (Sandbox Code Playgroud)
造成问题的主要论点:workers和use_multiprocessing。
当 时worker=1,use_multiprocessing=True/False运行没有问题。
如果workers=5,use_multiprocessing=True则抛出错误。奇怪的是它正在运行,但是在一些随机迭代中我遇到了类似的错误
KeyError: 'Unable to open object (bad local heap signature)'
Run Code Online (Sandbox Code Playgroud)
或者
KeyError: 'Unable to open object (wrong B-tree signature)'
Run Code Online (Sandbox Code Playgroud)
我使用 h5py 来读取文件。我为此目的编写了自定义生成器。
def load_random_cached_bottlenecks(batch_size, label_map,
addr_label_map, dirs, comp_type = 'h5py', hdf5_file …Run Code Online (Sandbox Code Playgroud) parallel-processing multithreading deep-learning keras tensorflow
我需要并行执行 Python 脚本,因此我使用以下批处理文件:
start python C:\myfolder\1.py
start python C:\myfolder\2.py
start python C:\myfolder\3.py
Run Code Online (Sandbox Code Playgroud)
它工作正常,但现在我需要在上述前三个脚本完成后并行运行三个脚本。如何在同一个批处理文件中指定它?
假设我有N 个生成器,可以生成项目流gs = [..] # list of generators。
我可以轻松地将它们组合在一起,从:zip中的每个相应生成器中获取元组生成器。gstuple_gen = zip(*gs)
这会按顺序调用next(g)每个并将结果收集到一个元组中。但如果每个项目的生产成本都很高,我们可能希望在多个线程上并行化工作。ggsnext(g)
我怎样才能实现pzip(..)这个功能呢?
python parallel-processing multithreading iterator generator
我正在使用 Pythonmultiprocessing创建并行应用程序。进程需要共享一些数据,为此我使用Manager. 但是,我有一些进程需要调用的常用函数以及需要访问对象存储的数据的函数Manager。我的问题是我是否可以避免需要将Manager实例作为参数传递给这些通用函数,而是像全局函数一样使用它。换句话说,请考虑以下代码:
import multiprocessing as mp
manager = mp.Manager()
global_dict = manager.dict(a=[0])
def add():
global_dict['a'] += [global_dict['a'][-1]+1]
def foo_parallel(var):
add()
print var
num_processes = 5
p = []
for i in range(num_processes):
p.append(mp.Process(target=foo_parallel,args=(global_dict,)))
[pi.start() for pi in p]
[pi.join() for pi in p]
Run Code Online (Sandbox Code Playgroud)
这运行良好并返回p=[0,1,2,3,4,5]到我的机器上。然而,这是“好形式”吗?这是一个好方法吗,就像定义add(var)和调用add(var)一样好?
python parallel-processing multithreading multiprocessing python-multiprocessing
我几天来一直在寻找这个问题的答案,但没有结果。我可能只是不理解那些漂浮在外面的部分,并且该multiprocessing模块的 Python 文档相当大,对我来说不清楚。
假设您有以下 for 循环:
import timeit
numbers = []
start = timeit.default_timer()
for num in range(100000000):
numbers.append(num)
end = timeit.default_timer()
print('TIME: {} seconds'.format(end - start))
print('SUM:', sum(numbers))
Run Code Online (Sandbox Code Playgroud)
输出:
TIME: 23.965870224497916 seconds
SUM: 4999999950000000
Run Code Online (Sandbox Code Playgroud)
对于此示例,假设您有一个 4 核处理器。有没有办法总共创建 4 个进程,其中每个进程都在单独的 CPU 核心上运行,并且完成速度大约快 4 倍,因此 24 秒/4 个进程 = 约 6 秒?
以某种方式将 for 循环分成 4 个相等的块,然后将这 4 个块添加到数字列表中以等于相同的总和?有一个 stackoverflow 线程:Parallel Simple For Loop但我不明白。谢谢大家。
考虑以下张量流代码片段:
import time
import numpy as np
import tensorflow as tf
def fn(i):
# do some junk work
for _ in range(100):
i ** 2
return i
n = 1000
n_jobs = 8
stuff = np.arange(1, n + 1)
eager = False
t0 = time.time()
if eager:
tf.enable_eager_execution()
res = tf.map_fn(fn, stuff, parallel_iterations=n_jobs)
if not eager:
with tf.Session() as sess:
res = sess.run(res)
print(sum(res))
else:
print(sum(res))
dt = time.time() - t0
print("(eager=%s) Took %ims" % (eager, dt * 1000))
Run Code Online (Sandbox Code Playgroud)
如果使用 …
我想并行化一个 for 循环,其中更新 unordered_map 的值:
unordered_map<string,double> umap {{"foo", 0}, {"bar", 0}};
#pragma omp parallel for reduction(my_reduction:umap)
for (int i = 0; i < 100; ++i)
{
// some_string(i) would return either "foo" or "bar"
umap[some_string(i)] += some_double(i);
}
Run Code Online (Sandbox Code Playgroud)
因此,unordered_map 中不会创建新条目,只会更新现有条目的总和。
在这个答案中,用户声明的归约是针对向量的情况定义的。在 unordered_map 的情况下,用户声明的归约是否可以类似地定义?
我正在读取一个视频文件,每 20 帧我将第一帧存储在输入队列中。一旦我在输入队列中获得了所有必需的帧,我就会运行多个进程来对这些帧执行一些操作并将结果存储在输出队列中。但代码总是卡在 join 处,我尝试了针对此类问题提出的不同解决方案,但似乎都不起作用。
import numpy as np
import cv2
import timeit
import face_recognition
from multiprocessing import Process, Queue, Pool
import multiprocessing
import os
s = timeit.default_timer()
def alternative_process_target_func(input_queue, output_queue):
while not output_queue.full():
frame_no, small_frame, face_loc = input_queue.get()
print('Frame_no: ', frame_no, 'Process ID: ', os.getpid(), '----', multiprocessing.current_process())
#canny_frame(frame_no, small_frame, face_loc)
#I am just storing frame no for now but will perform something else later
output_queue.put((frame_no, frame_no))
if output_queue.full():
print('Its Full ---------------------------------------------------------------------------------------')
else:
print('Not Full')
print(timeit.default_timer() - s, ' seconds.') …Run Code Online (Sandbox Code Playgroud) 我的场景是这样的。
我能想到的一个初步解决方案是增加 max-open-requests 设置。但这里的问题是我不知道需要提前发送多少报告。
有人可以建议一个替代解决方案,例如限制通过 Futures.traverse 发生的并行性
据我所知,水晶会通过 io 循环光纤,这意味着如果一根光纤正在等待 io,水晶将切换到另一根光纤。
如果我们生成两个纤程,但其中一个在没有 io 的情况下进行持续计算/循环,该怎么办?
例如,使用下面的代码,服务器不会响应任何 http 请求
spawn do
Kemal.run
end
spawn do
# constant computation/loop with no IO
some_func
end
Fiber.yield
# or sleep
Run Code Online (Sandbox Code Playgroud) python ×6
tensorflow ×2
akka ×1
batch-file ×1
c++ ×1
concurrency ×1
crystal-lang ×1
future ×1
generator ×1
iterator ×1
kemal ×1
keras ×1
openmp ×1
process ×1
python-2.7 ×1
python-3.x ×1
range ×1
scala ×1
windows ×1