Skip to content

Commit 93f54b1

Browse files
Merge pull request #3023 from devitocodes/revamp-buf-async-degree
compiler: Revamp bug-async-degree opt-option
2 parents 760a476 + d64d7bc commit 93f54b1

6 files changed

Lines changed: 162 additions & 28 deletions

File tree

‎.github/workflows/lint.yaml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ jobs:
4040
4141
- name: Lint codebase with flake8
4242
run: |
43-
flake8 --builtins=ArgumentError .
43+
flake8 .
4444
4545
spellcheck:
4646
name: "Spellcheck everything"

‎devito/core/operator.py‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -263,6 +263,13 @@ def _check_kwargs(cls, **kwargs):
263263
"`npthreads` must be a positive integer"
264264
)
265265

266+
async_degree = oo['buf-async-degree']
267+
if async_degree is not None and \
268+
(type(async_degree) is not int or async_degree < 0):
269+
raise InvalidOperator(
270+
"`buf-async-degree` must be a non-negative integer"
271+
)
272+
266273
if oo['cire-maxpar'] not in (False, 'basic', 'compact'):
267274
raise InvalidOperator("Illegal `cire-maxpar` value")
268275

‎devito/passes/clusters/buffering.py‎

Lines changed: 29 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -44,11 +44,12 @@ def buffering(clusters, key, sregistry, options, **kwargs):
4444
Accepted: ['buf-async-degree', 'buf-reuse', 'npthreads'].
4545
* 'buf-async-degree': Specify the size of the buffer. By default, the
4646
buffer size is the minimal one, inferred from the memory accesses in
47-
the ``clusters`` themselves. An asynchronous degree equals to `k`
48-
means that the buffer will be enforced to size=`k` along the introduced
49-
ModuloDimensions. This might help relieving the synchronization
50-
overhead when asynchronous operations are used (these are however
51-
implemented by other passes).
47+
the ``clusters`` themselves. A positive asynchronous degree `k`
48+
requests `k` slots; values below the inferred minimum are ignored.
49+
Zero disables buffering. Read-buffer initialization remains limited to
50+
the minimum number of slots required by the memory accesses. A larger
51+
buffer might relieve synchronization overhead in asynchronous operations
52+
introduced by other passes.
5253
* 'buf-reuse': If True, the pass will try to reuse existing Buffers for
5354
different buffered Functions. By default, False.
5455
* 'npthreads': Number of pthreads for asynchronous tasks. The tasks are
@@ -559,11 +560,15 @@ def expand_halo_transfers(clusters, mapper):
559560
return processed
560561

561562

562-
def _include_halo(ispace, f):
563-
"""Extend `ispace` to include `f`'s HALO."""
563+
def _include_halo(ispace, f, dims=None):
564+
"""
565+
Extend `ispace` to include `f`'s HALO along `dims`.
566+
"""
567+
dims = dims or f.dimensions
568+
564569
ihalo = [
565570
Interval(i.dim, -f._size_halo[i.dim].left, f._size_halo[i.dim].right, i.stamp)
566-
for i in ispace if i.dim in f.dimensions
571+
for i in ispace if i.dim in dims
567572
]
568573

569574
return IterationSpace.union(ispace, IterationSpace(ihalo))
@@ -724,8 +729,9 @@ def write_to(self):
724729
# might be accessed through a stencil
725730
ispace = ispace.promote(lambda d: d.is_AbstractSub, mode='total')
726731

727-
# Analogous to the above, we need to include the halo region as well
728-
ispace = _include_halo(ispace, self.b)
732+
# Include the spatial halo without widening the temporal interval,
733+
# which may already be restricted by an earlier buffering round
734+
ispace = _include_halo(ispace, self.b, self.bdims)
729735

730736
return ispace
731737

@@ -872,6 +878,7 @@ def init_buffers(descriptors, options):
872878
Create the initializing Clusters for the given buffers.
873879
"""
874880
init_onwrite = options['buf-init-onwrite']
881+
async_degree = options['buf-async-degree']
875882

876883
init = []
877884
for b, v in descriptors.flat_items():
@@ -882,6 +889,7 @@ def init_buffers(descriptors, options):
882889
# multiple) buffering because it's completely unnecessary
883890
if v.is_double_buffering:
884891
continue
892+
885893
lhs = b.indexify()._subs(v.xd, v.first_idx.b)
886894
rhs = f.indexify()._subs(v.dim, v.first_idx.f)
887895

@@ -895,7 +903,17 @@ def init_buffers(descriptors, options):
895903
expr = Eq(lhs, rhs)
896904
expr = lower_exprs(expr)
897905

898-
ispace = v.write_to
906+
ispace = v.write_to.concrete
907+
if v.is_read and async_degree is not None:
908+
# The allocated capacity (`v.size`) may exceed the time-window width
909+
# that must be loaded before computation starts (`size` below). E.g.,
910+
# reads at u[t-1], u[t] and u[t+1] make `infer_buffer_size` return 3,
911+
# even if `buf-async-degree` gives us 4 slots (`v.size == 4`). Seed
912+
# only db0=0..2; the spare slot is filled as computation advances.
913+
# This preserves the stencil's data space and iteration bounds,
914+
# without requiring extra input time levels to fill the ring.
915+
size = infer_buffer_size(f, v.dim, v.clusters)
916+
ispace = ispace.translate(v.xd, 0, size - v.size)
899917

900918
guards = {}
901919
guards[None] = GuardBound(v.dim.root.symbolic_min, v.dim.root.symbolic_max)

‎tests/test_buffering.py‎

Lines changed: 79 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -212,7 +212,8 @@ def test_read_only_w_offset():
212212
assert np.all(v.data == v1.data)
213213

214214

215-
def test_read_only_backwards():
215+
@pytest.mark.parametrize('async_degree,expected_size', [(None, 3), (4, 4)])
216+
def test_read_only_backwards(async_degree, expected_size):
216217
nt = 10
217218
grid = Grid(shape=(2, 2))
218219

@@ -226,12 +227,14 @@ def test_read_only_backwards():
226227
eqns = [Eq(v.backward, v + u.backward + u + u.forward + 1.)]
227228

228229
op0 = Operator(eqns, opt='noop')
229-
op1 = Operator(eqns, opt='buffering')
230+
op1 = Operator(eqns, opt=('buffering',
231+
{'buf-async-degree': async_degree}))
230232

231233
# Check generated code
232234
assert len(retrieve_iteration_tree(op1)) == 4
233235
buffers = [i for i in FindSymbols().visit(op1) if i.is_Array and i._mem_heap]
234236
assert len(buffers) == 1
237+
assert buffers.pop().symbolic_shape[0] == expected_size
235238

236239
op0.apply(time_m=1)
237240
op1.apply(time_m=1, v=v1)
@@ -270,32 +273,83 @@ def test_read_only_backwards_unstructured():
270273
assert np.all(v.data == v1.data)
271274

272275

273-
@pytest.mark.parametrize('async_degree', [2, 4])
274-
def test_async_degree(async_degree):
276+
@pytest.mark.parametrize('async_degree', [1, 2, 4])
277+
@pytest.mark.parametrize('backward', [False, True],
278+
ids=['forward', 'backward'])
279+
def test_async_degree(async_degree, backward):
275280
nt = 10
276281
grid = Grid(shape=(4, 4))
277282

278283
u = TimeFunction(name='u', grid=grid, save=nt)
279284
u1 = TimeFunction(name='u', grid=grid, save=nt)
280285

281-
eqn = Eq(u.forward, u + 1)
286+
lhs = u.backward if backward else u.forward
287+
eqn = Eq(lhs, u + 1)
282288

283289
op0 = Operator(eqn, opt='noop')
284290
op1 = Operator(eqn, opt=('buffering', {'buf-async-degree': async_degree}))
285291

286292
# Check generated code
287293
assert len(retrieve_iteration_tree(op1)) == 3
288-
buffers = [i for i in FindSymbols().visit(op1) if i.is_Array and i._mem_heap]
294+
buffers = [i for i in FindSymbols().visit(op1)
295+
if i.is_Array and i._mem_heap]
289296
assert len(buffers) == 1
290-
assert buffers.pop().symbolic_shape[0] == async_degree
297+
assert buffers.pop().symbolic_shape[0] == max(2, async_degree)
291298

292-
op0.apply(time_M=nt-2)
293-
op1.apply(time_M=nt-2, u=u1)
299+
kwargs = {'time_m': 1} if backward else {'time_M': nt - 2}
300+
op0.apply(**kwargs)
301+
op1.apply(u=u1, **kwargs)
294302

295303
assert np.all(u.data == u1.data)
296304

297305

298-
def test_two_homogeneous_buffers():
306+
@pytest.mark.parametrize('backward,expected_bounds', [
307+
pytest.param(False, (0, 8), id='forward'),
308+
pytest.param(True, (1, 9), id='backward')
309+
])
310+
@pytest.mark.parametrize('async_degree', [0, 1, 4, 16])
311+
def test_async_degree_read_only(backward, expected_bounds, async_degree):
312+
nt = 10
313+
grid = Grid(shape=(4, 4))
314+
315+
u = TimeFunction(name='u', grid=grid, save=nt)
316+
v = TimeFunction(name='v', grid=grid)
317+
v1 = TimeFunction(name='v', grid=grid)
318+
319+
u.data[:] = np.arange(nt).reshape(nt, 1, 1)
320+
321+
lhs = v.backward if backward else v.forward
322+
eqn = Eq(lhs, v + u)
323+
324+
op0 = Operator(eqn, opt='noop', name='op0')
325+
op1 = Operator(eqn, opt=('buffering',
326+
{'buf-async-degree': async_degree}), name='op1')
327+
328+
buffers = [i for i in FindSymbols().visit(op1)
329+
if i.is_Array and i._mem_heap]
330+
assert len(buffers) == int(async_degree != 0)
331+
if async_degree:
332+
assert buffers[0].symbolic_shape[0] == async_degree
333+
334+
for op in [op0, op1]:
335+
args = op.arguments()
336+
assert (args['time_m'], args['time_M']) == expected_bounds
337+
338+
# Default bounds, either endpoint, a partial ring, and an empty interval
339+
time_m, time_M = expected_bounds
340+
for kwargs in [{}, {'time_m': time_m, 'time_M': time_m},
341+
{'time_m': time_M, 'time_M': time_M},
342+
{'time_m': 3, 'time_M': 4}, {'time_m': 1, 'time_M': 0}]:
343+
v.data[:] = 0
344+
v1.data[:] = 0
345+
op0.apply(**kwargs)
346+
op1.apply(v=v1, **kwargs)
347+
348+
assert np.all(v.data == v1.data)
349+
350+
351+
@pytest.mark.parametrize('async_degree', [None, 4])
352+
def test_two_homogeneous_buffers(async_degree):
299353
nt = 10
300354
grid = Grid(shape=(4, 4))
301355

@@ -308,8 +362,10 @@ def test_two_homogeneous_buffers():
308362
Eq(v.forward, u + v + u.backward + v.backward + 1.)]
309363

310364
op0 = Operator(eqns, opt='noop')
311-
op1 = Operator(eqns, opt='buffering')
312-
op2 = Operator(eqns, opt=('buffering', 'fuse'))
365+
op1 = Operator(eqns, opt=('buffering',
366+
{'buf-async-degree': async_degree}))
367+
op2 = Operator(eqns, opt=('buffering', 'fuse',
368+
{'buf-async-degree': async_degree}))
313369

314370
# Check generated code
315371
assert len(retrieve_iteration_tree(op1)) == 5
@@ -323,8 +379,16 @@ def test_two_homogeneous_buffers():
323379
assert np.all(u.data == u1.data)
324380
assert np.all(v.data == v1.data)
325381

382+
u1.data[:] = 0
383+
v1.data[:] = 0
384+
op2.apply(time_M=nt-2, u=u1, v=v1)
385+
386+
assert np.all(u.data == u1.data)
387+
assert np.all(v.data == v1.data)
388+
326389

327-
def test_two_heterogeneous_buffers():
390+
@pytest.mark.parametrize('async_degree', [None, 4])
391+
def test_two_heterogeneous_buffers(async_degree):
328392
nt = 10
329393
grid = Grid(shape=(4, 4))
330394

@@ -341,7 +405,8 @@ def test_two_heterogeneous_buffers():
341405
Eq(v.forward, u + v + v.backward)]
342406

343407
op0 = Operator(eqns, opt='noop')
344-
op1 = Operator(eqns, opt='buffering')
408+
op1 = Operator(eqns, opt=('buffering',
409+
{'buf-async-degree': async_degree}))
345410

346411
# Check generated code
347412
assert len(retrieve_iteration_tree(op1)) == 5

‎tests/test_gpu_common.py‎

Lines changed: 41 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -814,6 +814,7 @@ def test_streaming_conddim_forward(self, opt):
814814

815815
@pytest.mark.parametrize('opt', [
816816
('buffering', 'streaming', 'orchestrate'),
817+
('buffering', 'streaming', 'orchestrate', {'buf-async-degree': 4}),
817818
])
818819
def test_streaming_conddim_backward(self, opt):
819820
nt = 10
@@ -849,6 +850,42 @@ def test_streaming_conddim_backward(self, opt):
849850
# 3rd time u[1] = u[0]+u[1]+usave[2] = 0+7+2 = 9
850851
assert np.all(u.data[1] == 9)
851852

853+
@pytest.mark.parametrize('mode', [
854+
pytest.param(None, id='serial'),
855+
pytest.param(2, marks=pytest.mark.parallel, id='basic')
856+
])
857+
@pytest.mark.parametrize('backward,expected_bounds', [
858+
pytest.param(False, (0, 8), id='forward'),
859+
pytest.param(True, (1, 9), id='backward')
860+
])
861+
@pytest.mark.parametrize('async_degree', [4, 16])
862+
def test_streaming_async_degree(self, mode, backward, expected_bounds,
863+
async_degree):
864+
nt = 10
865+
grid = Grid(shape=(4, 4))
866+
867+
usave = TimeFunction(name='usave', grid=grid, save=nt)
868+
v = TimeFunction(name='v', grid=grid)
869+
v1 = TimeFunction(name='v', grid=grid)
870+
871+
usave.data._local[:] = np.arange(nt).reshape(nt, 1, 1)
872+
873+
lhs = v.backward if backward else v.forward
874+
eqn = Eq(lhs, v + usave)
875+
876+
op0 = Operator(eqn, opt=('noop', {'gpu-fit': usave}), name='op0')
877+
op1 = Operator(eqn, opt=('buffering', 'streaming', 'orchestrate',
878+
{'buf-async-degree': async_degree}), name='op1')
879+
880+
for op in [op0, op1]:
881+
args = op.arguments()
882+
assert (args['time_m'], args['time_M']) == expected_bounds
883+
884+
op0.apply()
885+
op1.apply(v=v1)
886+
887+
assert np.all(v.data == v1.data)
888+
852889
@pytest.mark.parametrize('opt,ntmps', [
853890
(('buffering', 'streaming', 'orchestrate'), 3),
854891
])
@@ -910,7 +947,8 @@ def test_streaming_multi_input_conddim_foward(self):
910947

911948
assert np.all(v.data == v1.data)
912949

913-
def test_streaming_multi_input_conddim_backward(self):
950+
@pytest.mark.parametrize('async_degree', [None, 5])
951+
def test_streaming_multi_input_conddim_backward(self, async_degree):
914952
nt = 10
915953
grid = Grid(shape=(4, 4))
916954
time_dim = grid.time_dim
@@ -931,7 +969,8 @@ def test_streaming_multi_input_conddim_backward(self):
931969
eqns = [Eq(v.backward, v + expr + 1.)]
932970

933971
op0 = Operator(eqns, opt=('noop', {'gpu-fit': u}))
934-
op1 = Operator(eqns, opt=('buffering', 'streaming', 'orchestrate'))
972+
op1 = Operator(eqns, opt=('buffering', 'streaming', 'orchestrate',
973+
{'buf-async-degree': async_degree}))
935974

936975
op0.apply(time_M=nt, dt=.01)
937976
op1.apply(time_M=nt, dt=.01, v=v1)

‎tests/test_operator.py‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -141,6 +141,11 @@ def test_opt_options(self):
141141
Operator(Eq(u, u + 1),
142142
opt=('advanced', {'npthreads': npthreads}))
143143

144+
for async_degree in (False, True, -1, 1.5):
145+
with pytest.raises(InvalidOperator, match='non-negative integer'):
146+
Operator(Eq(u, u + 1),
147+
opt=('advanced', {'buf-async-degree': async_degree}))
148+
144149
def test_compiler_uniqueness(self):
145150
grid = Grid(shape=(3, 3, 3))
146151

0 commit comments

Comments
 (0)