|
17 | 17 | */ |
18 | 18 | package org.apache.giraph.block_app.library; |
19 | 19 |
|
| 20 | +import java.util.ArrayList; |
20 | 21 | import java.util.Iterator; |
| 22 | +import java.util.List; |
21 | 23 |
|
22 | 24 | import org.apache.giraph.block_app.framework.api.BlockMasterApi; |
23 | 25 | import org.apache.giraph.block_app.framework.api.BlockWorkerReceiveApi; |
|
26 | 28 | import org.apache.giraph.block_app.framework.piece.Piece; |
27 | 29 | import org.apache.giraph.block_app.framework.piece.global_comm.ReducerAndBroadcastWrapperHandle; |
28 | 30 | import org.apache.giraph.block_app.framework.piece.global_comm.ReducerHandle; |
| 31 | +import org.apache.giraph.block_app.framework.piece.global_comm.array.BroadcastArrayHandle; |
29 | 32 | import org.apache.giraph.block_app.framework.piece.interfaces.VertexReceiver; |
30 | 33 | import org.apache.giraph.block_app.framework.piece.interfaces.VertexSender; |
31 | 34 | import org.apache.giraph.block_app.library.internal.SendMessagePiece; |
32 | 35 | import org.apache.giraph.block_app.library.internal.SendMessageWithCombinerPiece; |
| 36 | +import org.apache.giraph.block_app.reducers.array.ArrayOfHandles; |
33 | 37 | import org.apache.giraph.combiner.MessageCombiner; |
34 | 38 | import org.apache.giraph.function.Consumer; |
35 | 39 | import org.apache.giraph.function.PairConsumer; |
| 40 | +import org.apache.giraph.function.Supplier; |
36 | 41 | import org.apache.giraph.function.vertex.ConsumerWithVertex; |
37 | 42 | import org.apache.giraph.function.vertex.SupplierFromVertex; |
38 | 43 | import org.apache.giraph.graph.Vertex; |
@@ -319,6 +324,91 @@ public String toString() { |
319 | 324 | }; |
320 | 325 | } |
321 | 326 |
|
| 327 | + /** |
| 328 | + * Like reduceAndBroadcast, but uses array of handles for reducers and |
| 329 | + * broadcasts, to make it feasible and performant when values are large. |
| 330 | + * Each supplied value to reduce will be reduced in the handle defined by |
| 331 | + * handleHashSupplier%numHandles |
| 332 | + * |
| 333 | + * @param <S> Single value type, objects passed on workers |
| 334 | + * @param <R> Reduced value type |
| 335 | + * @param <I> Vertex id type |
| 336 | + * @param <V> Vertex value type |
| 337 | + * @param <E> Edge value type |
| 338 | + */ |
| 339 | + public static |
| 340 | + <S, R extends Writable, I extends WritableComparable, V extends Writable, |
| 341 | + E extends Writable> |
| 342 | + Piece<I, V, E, NoMessage, Object> reduceAndBroadcastWithArrayOfHandles( |
| 343 | + final String name, |
| 344 | + final int numHandles, |
| 345 | + final Supplier<ReduceOperation<S, R>> reduceOp, |
| 346 | + final SupplierFromVertex<I, V, E, Long> handleHashSupplier, |
| 347 | + final SupplierFromVertex<I, V, E, S> valueSupplier, |
| 348 | + final ConsumerWithVertex<I, V, E, R> reducedValueConsumer) { |
| 349 | + return new Piece<I, V, E, NoMessage, Object>() { |
| 350 | + protected ArrayOfHandles.ArrayOfReducers<S, R> reducers; |
| 351 | + protected BroadcastArrayHandle<R> broadcasts; |
| 352 | + |
| 353 | + private int getHandleIndex(Vertex<I, V, E> vertex) { |
| 354 | + return (int) Math.abs(handleHashSupplier.get(vertex) % numHandles); |
| 355 | + } |
| 356 | + |
| 357 | + @Override |
| 358 | + public void registerReducers( |
| 359 | + final CreateReducersApi reduceApi, Object executionStage) { |
| 360 | + reducers = new ArrayOfHandles.ArrayOfReducers<>( |
| 361 | + numHandles, |
| 362 | + new Supplier<ReducerHandle<S, R>>() { |
| 363 | + @Override |
| 364 | + public ReducerHandle<S, R> get() { |
| 365 | + return reduceApi.createLocalReducer(reduceOp.get()); |
| 366 | + } |
| 367 | + }); |
| 368 | + } |
| 369 | + |
| 370 | + @Override |
| 371 | + public VertexSender<I, V, E> getVertexSender( |
| 372 | + BlockWorkerSendApi<I, V, E, NoMessage> workerApi, |
| 373 | + Object executionStage) { |
| 374 | + return new InnerVertexSender() { |
| 375 | + @Override |
| 376 | + public void vertexSend(Vertex<I, V, E> vertex) { |
| 377 | + reducers.get(getHandleIndex(vertex)).reduce( |
| 378 | + valueSupplier.get(vertex)); |
| 379 | + } |
| 380 | + }; |
| 381 | + } |
| 382 | + |
| 383 | + @Override |
| 384 | + public void masterCompute(BlockMasterApi master, Object executionStage) { |
| 385 | + broadcasts = reducers.broadcastValue(master); |
| 386 | + } |
| 387 | + |
| 388 | + @Override |
| 389 | + public VertexReceiver<I, V, E, NoMessage> getVertexReceiver( |
| 390 | + BlockWorkerReceiveApi<I> workerApi, Object executionStage) { |
| 391 | + final List<R> values = new ArrayList<>(); |
| 392 | + for (int i = 0; i < numHandles; i++) { |
| 393 | + values.add(broadcasts.get(i).getBroadcast(workerApi)); |
| 394 | + } |
| 395 | + return new InnerVertexReceiver() { |
| 396 | + @Override |
| 397 | + public void vertexReceive( |
| 398 | + Vertex<I, V, E> vertex, Iterable<NoMessage> messages) { |
| 399 | + reducedValueConsumer.apply( |
| 400 | + vertex, values.get(getHandleIndex(vertex))); |
| 401 | + } |
| 402 | + }; |
| 403 | + } |
| 404 | + |
| 405 | + @Override |
| 406 | + public String toString() { |
| 407 | + return name; |
| 408 | + } |
| 409 | + }; |
| 410 | + } |
| 411 | + |
322 | 412 | /** |
323 | 413 | * Creates Piece that for each vertex, sends message provided by |
324 | 414 | * messageSupplier to all targets provided by targetsSupplier. |
|
0 commit comments