more robust parallel stripes

* use finally to ensure shared memory is closed/unliked
* use events to coordinate jobs
* "minor" refactoring
This commit is contained in:
Bruno Madeira
2024-08-06 11:36:33 +01:00
parent df0b1ee9d6
commit f39b88fc2b
+168 -125
View File
@@ -1,6 +1,7 @@
from .patch_search import get_find_patch_to_the_right_method, get_find_patch_below_method, get_find_patch_both_method
from .jena2020.generate import getMinCutPatchHorizontal, getMinCutPatchVertical, getMinCutPatchBoth
from multiprocessing.shared_memory import SharedMemory
from dataclasses import dataclass
from .types import UiCoordData
from math import ceil
import numpy as np
@@ -53,11 +54,11 @@ def fill_column(image, initial_block, overlap, rows: int, tolerance, version, rn
texture_map = np.zeros(
((block_size + rows * (block_size - overlap)), block_size, image.shape[2])).astype(image.dtype)
texture_map[:block_size, :block_size, :] = initial_block
for i, blkIdx in enumerate(range((block_size - overlap), texture_map.shape[0] - overlap, (block_size - overlap))):
ref_block = texture_map[(blkIdx - block_size + overlap):(blkIdx + overlap), :block_size]
for i, blk_idx in enumerate(range((block_size - overlap), texture_map.shape[0] - overlap, (block_size - overlap))):
ref_block = texture_map[(blk_idx - block_size + overlap):(blk_idx + overlap), :block_size]
patch_block = find_patch_below(ref_block, image, block_size, overlap, tolerance, rng)
min_cut_patch = getMinCutPatchVertical(ref_block, patch_block, block_size, overlap)
texture_map[blkIdx:(blkIdx + block_size), :block_size] = min_cut_patch
texture_map[blk_idx:(blk_idx + block_size), :block_size] = min_cut_patch
return texture_map
@@ -67,11 +68,11 @@ def fill_row(image, initial_block, overlap, columns: int, tolerance, version: in
texture_map = np.zeros(
(block_size, (block_size + columns * (block_size - overlap)), image.shape[2])).astype(image.dtype)
texture_map[:block_size, :block_size, :] = initial_block
for i, blkIdx in enumerate(range((block_size - overlap), texture_map.shape[1] - overlap, (block_size - overlap))):
ref_block = texture_map[:block_size, (blkIdx - block_size + overlap):(blkIdx + overlap)]
for i, blk_idx in enumerate(range((block_size - overlap), texture_map.shape[1] - overlap, (block_size - overlap))):
ref_block = texture_map[:block_size, (blk_idx - block_size + overlap):(blk_idx + overlap)]
patch_block = find_patch_to_the_right(ref_block, image, block_size, overlap, tolerance, rng)
min_cut_patch = getMinCutPatchHorizontal(ref_block, patch_block, block_size, overlap)
texture_map[:block_size, blkIdx:(blkIdx + block_size)] = min_cut_patch
texture_map[:block_size, blk_idx:(blk_idx + block_size)] = min_cut_patch
return texture_map
@@ -105,7 +106,7 @@ def fill_quad(rows: int, columns: int, block_size, overlap, texture_map, image,
# region parallel solution
def generate_texture_parallel(image, block_size, overlap, outH, outW, tolerance, version: int, nps,
def generate_texture_parallel(image, block_size, overlap, out_h, out_w, tolerance, version: int, nps,
rng: np.random.Generator, uicd: UiCoordData | None):
"""
@param uicd: contains:
@@ -139,8 +140,8 @@ def generate_texture_parallel(image, block_size, overlap, outH, outW, tolerance,
# note: inverted both ways is computed later when filling the top-left quadrant/section
# auxiliary variables
quad_row_width = ceil(outW / 2 + block_size / 2)
quad_column_height = ceil(outH / 2 + block_size / 2)
quad_row_width = ceil(out_w / 2 + block_size / 2)
quad_column_height = ceil(out_h / 2 + block_size / 2)
cols_per_quad = int(ceil((quad_row_width - block_size) / (block_size - overlap))) # minding the overlap
rows_per_quad = int(ceil((quad_column_height - block_size) / (block_size - overlap)))
# might get 1 more row or column than needed here
@@ -180,7 +181,7 @@ def generate_texture_parallel(image, block_size, overlap, outH, outW, tolerance,
texture[q1.shape[0] - bmo:, :q1.shape[1] - bmo] = q4[overlap:, :q1.shape[1] - bmo]
texture[q1.shape[0] - bmo:, q1.shape[1] - bmo:] = q3[overlap:, overlap:]
return texture[:outH, :outW]
return texture[:out_h, :out_w]
def quad1(vis, his, hi_image, rows: int, columns: int, overlap, tolerance, version, p_strips, rng,
@@ -190,186 +191,228 @@ def quad1(vis, his, hi_image, rows: int, columns: int, overlap, tolerance, versi
@param vis: vertical inverted stripe
@param p_strips: number of sub processes to build the quadrant (only when para_lvl higher than 1)
"""
shm_text = None
vi_hi_s = np.ascontiguousarray(np.flipud(his)) # vertical inversion of the horizontal inverted stripe
hi_vi_s = np.ascontiguousarray(np.fliplr(vis))
vhi_image = np.ascontiguousarray(np.flipud(hi_image))
def place_stripes_on_texture():
texture[:vi_hi_s.shape[0], :vi_hi_s.shape[1]] = vi_hi_s[:, :]
texture[vi_hi_s.shape[0]:hi_vi_s.shape[0], :hi_vi_s.shape[1]] = hi_vi_s[vi_hi_s.shape[0]:, :]
if p_strips > 1:
size = vis.shape[0] * his.shape[1] * hi_image.shape[2] * hi_image.dtype.itemsize
shm_text = SharedMemory(create=True, size=size)
texture = np.ndarray((vis.shape[0], his.shape[1], hi_image.shape[2]), dtype=hi_image.dtype, buffer=shm_text.buf)
try:
texture = np.ndarray((vis.shape[0], his.shape[1], hi_image.shape[2]), dtype=hi_image.dtype,
buffer=shm_text.buf)
place_stripes_on_texture()
fill_quad_ps(rows, columns, vi_hi_s.shape[0], overlap, version, shm_text.name, vhi_image, tolerance,
p_strips,
rng, None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 0))
texture = np.ascontiguousarray(np.flip(texture, axis=(0, 1)))
finally:
shm_text.close()
shm_text.unlink()
else:
texture = np.zeros((vis.shape[0], his.shape[1], hi_image.shape[2])).astype(hi_image.dtype)
texture[:vi_hi_s.shape[0], :vi_hi_s.shape[1]] = vi_hi_s[:, :]
texture[vi_hi_s.shape[0]:hi_vi_s.shape[0], :hi_vi_s.shape[1]] = hi_vi_s[vi_hi_s.shape[0]:, :]
if p_strips > 1:
fill_quad_ps(rows, columns, vi_hi_s.shape[0], overlap, version, shm_text.name, vhi_image, tolerance, p_strips,
rng, None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 0))
else:
place_stripes_on_texture()
texture = fill_quad(rows, columns, vi_hi_s.shape[0], overlap, texture, vhi_image, tolerance, version, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + 0))
texture = np.ascontiguousarray(np.flip(texture, axis=(0, 1)))
if p_strips > 1:
shm_text.close()
shm_text.unlink()
texture = np.ascontiguousarray(np.flip(texture, axis=(0, 1)))
return texture
def quad2(vis, hs, vi_image, rows: int, columns: int, overlap, tolerance, version, p_strips, rng,
uicd: UiCoordData | None):
shm_text = None
vi_hs = np.ascontiguousarray(np.flipud(hs))
def place_stripes_on_texture():
texture[:hs.shape[0], :hs.shape[1]] = vi_hs[:, :]
texture[hs.shape[0]:vis.shape[0], :vis.shape[1]] = vis[hs.shape[0]:, :]
if p_strips > 1:
size = vis.shape[0] * hs.shape[1] * vi_image.shape[2] * vi_image.dtype.itemsize
shm_text = SharedMemory(create=True, size=size)
texture = np.ndarray((vis.shape[0], hs.shape[1], vi_image.shape[2]), dtype=vi_image.dtype, buffer=shm_text.buf)
try:
texture = np.ndarray((vis.shape[0], hs.shape[1], vi_image.shape[2]), dtype=vi_image.dtype, buffer=shm_text.buf)
place_stripes_on_texture()
fill_quad_ps(rows, columns, hs.shape[0], overlap, version, shm_text.name, vi_image, tolerance, p_strips, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 1))
texture = np.ascontiguousarray(np.flipud(texture))
finally:
shm_text.close()
shm_text.unlink()
else:
texture = np.zeros((vis.shape[0], hs.shape[1], vi_image.shape[2])).astype(vi_image.dtype)
texture[:hs.shape[0], :hs.shape[1]] = vi_hs[:, :]
texture[hs.shape[0]:vis.shape[0], :vis.shape[1]] = vis[hs.shape[0]:, :]
if p_strips > 1:
fill_quad_ps(rows, columns, hs.shape[0], overlap, version, shm_text.name, vi_image, tolerance, p_strips, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 1))
else:
place_stripes_on_texture()
texture = fill_quad(rows, columns, hs.shape[0], overlap, texture, vi_image, tolerance, version, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + 1))
texture = np.ascontiguousarray(np.flipud(texture))
if p_strips > 1:
shm_text.close()
shm_text.unlink()
texture = np.ascontiguousarray(np.flipud(texture))
return texture
def quad4(vs, his, hi_image, rows: int, columns: int, overlap, tolerance, version, p_strips, rng,
uicd: UiCoordData | None):
shm_text = None
hi_vs = np.ascontiguousarray(np.fliplr(vs))
def place_stripes_on_texture():
texture[:his.shape[0], :his.shape[1]] = his[:, :]
texture[his.shape[0]:vs.shape[0], :vs.shape[1]] = hi_vs[his.shape[0]:, :]
if p_strips > 1:
size = vs.shape[0] * his.shape[1] * hi_image.shape[2] * hi_image.dtype.itemsize
shm_text = SharedMemory(create=True, size=size)
texture = np.ndarray((vs.shape[0], his.shape[1], hi_image.shape[2]), dtype=hi_image.dtype, buffer=shm_text.buf)
try:
texture = np.ndarray((vs.shape[0], his.shape[1], hi_image.shape[2]), dtype=hi_image.dtype, buffer=shm_text.buf)
place_stripes_on_texture()
fill_quad_ps(rows, columns, his.shape[0], overlap, version, shm_text.name, hi_image, tolerance, p_strips, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 2))
texture = np.ascontiguousarray(np.fliplr(texture))
finally:
shm_text.close()
shm_text.unlink()
else:
texture = np.zeros((vs.shape[0], his.shape[1], hi_image.shape[2])).astype(hi_image.dtype)
texture[:his.shape[0], :his.shape[1]] = his[:, :]
texture[his.shape[0]:vs.shape[0], :vs.shape[1]] = hi_vs[his.shape[0]:, :]
if p_strips > 1:
fill_quad_ps(rows, columns, his.shape[0], overlap, version, shm_text.name, hi_image, tolerance, p_strips, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 2))
else:
place_stripes_on_texture()
texture = fill_quad(rows, columns, his.shape[0], overlap, texture, hi_image, tolerance, version, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + 2))
texture = np.ascontiguousarray(np.fliplr(texture))
if p_strips > 1:
shm_text.close()
shm_text.unlink()
texture = np.ascontiguousarray(np.fliplr(texture))
return texture
def quad3(vs, hs, image, rows: int, columns: int, overlap, tolerance, version, p_strips, rng,
uicd: UiCoordData | None):
shm_text = None
def place_stripes_on_texture():
texture[:hs.shape[0], :hs.shape[1]] = hs[:, :]
texture[hs.shape[0]:vs.shape[0], :vs.shape[1]] = vs[hs.shape[0]:, :]
if p_strips > 1:
size = vs.shape[0] * hs.shape[1] * image.shape[2] * image.dtype.itemsize
shm_text = SharedMemory(create=True, size=size)
texture = np.ndarray((vs.shape[0], hs.shape[1], image.shape[2]), dtype=image.dtype, buffer=shm_text.buf)
try:
texture = np.ndarray((vs.shape[0], hs.shape[1], image.shape[2]), dtype=image.dtype, buffer=shm_text.buf)
place_stripes_on_texture()
fill_quad_ps(rows, columns, vs.shape[1], overlap, version, shm_text.name, image, tolerance, p_strips, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 3))
texture = texture.copy()
finally:
shm_text.close()
shm_text.unlink()
else:
texture = np.zeros((vs.shape[0], hs.shape[1], image.shape[2])).astype(image.dtype)
texture[:hs.shape[0], :hs.shape[1]] = hs[:, :]
texture[hs.shape[0]:vs.shape[0], :vs.shape[1]] = vs[hs.shape[0]:, :]
if p_strips > 1:
fill_quad_ps(rows, columns, vs.shape[1], overlap, version, shm_text.name, image, tolerance, p_strips, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + p_strips * 3))
else:
return fill_quad(rows, columns, vs.shape[1], overlap, texture, image, tolerance, version, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + 3))
texture = texture.copy()
if p_strips > 1:
shm_text.close()
shm_text.unlink()
place_stripes_on_texture()
texture = fill_quad(rows, columns, vs.shape[1], overlap, texture, image, tolerance, version, rng,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + 3))
return texture
def fill_quad_ps(rows, columns, block_size, overlap, version:int,
# region parallel cascading stripes
def fill_quad_ps(rows, columns, block_size, overlap, version: int,
texture_shared_mem_name: str, image, tolerance, total_procs, rng,
uicd: UiCoordData | None):
from joblib import Parallel, delayed
bmo = block_size - overlap # taking into account the overlap, what is the filled area at each iteration ?
b_o = ceil(block_size / bmo) # how many of the above are needed to fill a block ?
from multiprocessing import Manager
# note : ShareableList seems to have problems even though there are no concurrent writes to the same position...
# setup shared array to store & check each job row & column
size = 2 * total_procs * np.int32(1).itemsize
shm_coord = SharedMemory(create=True, size=size)
np_coord = np.ndarray((2 * total_procs,), dtype=np.int32, buffer=shm_coord.buf)
for ip in range(total_procs):
np_coord[2 * ip] = 1 + ip
np_coord[2 * ip + 1] = 1
try:
np_coord = np.ndarray((2 * total_procs,), dtype=np.int32, buffer=shm_coord.buf)
for ip in range(total_procs):
np_coord[2 * ip] = 1 + ip
np_coord[2 * ip + 1] = 1
def fill_rows(pid: int, coord_shared_list_name: str, texture_shm_name: str, uicd: UiCoordData | None):
find_patch_both = get_find_patch_both_method(version)
job_data = ParaRowsJobInfo(
total_procs=total_procs,
coord_shared_list_name=shm_coord.name,
texture_shm_name=texture_shared_mem_name,
src=image, block_size=block_size, overlap=overlap, tolerance=tolerance, rng=rng, version=version,
columns=columns, rows=rows
)
# get data in shared memory
shm_coord_ref = SharedMemory(name=coord_shared_list_name)
coord_list = np.ndarray((2 * total_procs,), dtype=np.int32, buffer=shm_coord_ref.buf)
prior_proc_base_index = (pid + total_procs - 1) % total_procs
shm_texture = SharedMemory(name=texture_shm_name)
texture = np.ndarray((block_size + rows * bmo, block_size + columns * bmo, image.shape[2]), dtype=image.dtype,
buffer=shm_texture.buf)
for i in range(1 + pid, rows + 1, total_procs):
coord_list[pid * 2 + 0] = i
for j in range(1, columns + 1):
coord_list[pid * 2 + 1] = j
# if previous row hasn't processed the adjacent section yet wait for it to advance.
# -1 is used to shortcircuit this check when the job on the prior row has completed all rows.
while -1 < coord_list[prior_proc_base_index * 2 + 0] < i and \
coord_list[prior_proc_base_index * 2 + 1] - b_o <= j:
pass
# The same as source implementation ( similar to fill_quad )
blk_index_i = i * bmo
blk_index_j = j * bmo
ref_block_left = texture[
blk_index_i:(blk_index_i + block_size),
(blk_index_j - block_size + overlap):(blk_index_j + overlap)]
ref_block_top = texture[
(blk_index_i - block_size + overlap):(blk_index_i + overlap),
blk_index_j:(blk_index_j + block_size)]
patch_block = find_patch_both(ref_block_left, ref_block_top, image, block_size, overlap, tolerance, rng)
min_cut_patch = getMinCutPatchBoth(ref_block_left, ref_block_top, patch_block, block_size, overlap)
texture[blk_index_i:(blk_index_i + block_size), blk_index_j:(blk_index_j + block_size)] = min_cut_patch
if uicd is not None and uicd.add_to_job_data_slot_and_check_interrupt(columns):
break # is required to set job as complete and avoid deadlock in jobs waiting for prior row completion
coord_list[pid * 2 + 0] = -1 # set job as completed
Parallel(n_jobs=total_procs, backend="loky", timeout=None)(
delayed(fill_rows)(i, shm_coord.name, texture_shared_mem_name,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + i))
for i in range(total_procs))
shm_coord.close()
shm_coord.unlink()
with Manager() as manager:
events = {pid: manager.Event() for pid in range(total_procs)} # the stripe sub-job index (not the os pid)
Parallel(n_jobs=total_procs, backend="loky", timeout=None)(
delayed(fill_rows_ps)(i, job_data, events,
None if uicd is None else UiCoordData(uicd.jobs_shm_name, uicd.job_id + i))
for i in range(total_procs))
finally:
shm_coord.close()
shm_coord.unlink()
# endregion
@dataclass
class ParaRowsJobInfo:
total_procs: int
coord_shared_list_name: str
texture_shm_name: str
src: np.ndarray
block_size: int
overlap: int
tolerance: float
rng: np.random.Generator
version: int
columns: int
rows: int
def fill_rows_ps(pid: int, job: ParaRowsJobInfo, jobs_events: list, uicd: UiCoordData | None):
find_patch_both = get_find_patch_both_method(job.version)
# unwrap data
block_size, overlap, tolerance, rng = job.block_size, job.overlap, job.tolerance, job.rng
total_procs, image, rows, columns = job.total_procs, job.src, job.rows, job.columns
bmo = block_size - overlap
b_o = ceil(block_size / bmo)
# get data in shared memory
shm_coord_ref = SharedMemory(name=job.coord_shared_list_name)
coord_list = np.ndarray((2 * job.total_procs,), dtype=np.int32, buffer=shm_coord_ref.buf)
prior_pid = (pid + job.total_procs - 1) % job.total_procs
shm_texture = SharedMemory(name=job.texture_shm_name)
texture = np.ndarray((block_size + rows * bmo, block_size + columns * bmo, image.shape[2]), dtype=image.dtype,
buffer=shm_texture.buf)
for i in range(1 + pid, rows + 1, total_procs):
coord_list[pid * 2 + 1] = -b_o # indicates no column yet processed on this row
coord_list[pid * 2 + 0] = i
for j in range(1, columns + 1):
# if previous row hasn't processed the adjacent section yet wait for it to advance.
# -1 is used to shortcircuit this check when the job on the prior row has completed all rows.
while -1 < coord_list[prior_pid * 2 + 0] < i and \
coord_list[prior_pid * 2 + 1] - b_o <= j:
jobs_events[prior_pid].wait()
jobs_events[prior_pid].clear()
# The same as source implementation ( similar to fill_quad )
blk_index_i = i * bmo
blk_index_j = j * bmo
ref_block_left = texture[
blk_index_i:(blk_index_i + block_size),
(blk_index_j - block_size + overlap):(blk_index_j + overlap)]
ref_block_top = texture[
(blk_index_i - block_size + overlap):(blk_index_i + overlap),
blk_index_j:(blk_index_j + block_size)]
patch_block = find_patch_both(ref_block_left, ref_block_top, image, block_size, overlap, tolerance, rng)
min_cut_patch = getMinCutPatchBoth(ref_block_left, ref_block_top, patch_block, block_size, overlap)
texture[blk_index_i:(blk_index_i + block_size), blk_index_j:(blk_index_j + block_size)] = min_cut_patch
coord_list[pid * 2 + 1] = j
jobs_events[pid].set()
if uicd is not None and uicd.add_to_job_data_slot_and_check_interrupt(columns):
break # is required to set job as complete and avoid deadlock in jobs waiting for prior row completion
coord_list[pid * 2 + 0] = -1 # set job as completed
jobs_events[pid].set()
# endregion parallel cascading stripes
# endregion parallel solution
# region non-parallel solution adapted from jena2020 for node compliance & error function optionality