144 lines
4.5 KiB
Java
Raw Normal View History

2019-06-06 15:21:15 +03:00
/*******************************************************************************
* Copyright (c) 2015-2018 Skymind, Inc.
*
* This program and the accompanying materials are made available under the
* terms of the Apache License, Version 2.0 which is available at
* https://www.apache.org/licenses/LICENSE-2.0.
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
* License for the specific language governing permissions and limitations
* under the License.
*
* SPDX-License-Identifier: Apache-2.0
******************************************************************************/
Refactor packages to fix split package issues (#411) * Refactor nd4j-common: org.nd4j.* -> org.nd4j.common.* Signed-off-by: Alex Black <blacka101@gmail.com> * Fix CUDA (missed nd4j-common package refactoring changes) Signed-off-by: Alex Black <blacka101@gmail.com> * nd4j-kryo: org.nd4j -> org.nd4j.kryo Signed-off-by: Alex Black <blacka101@gmail.com> * Fix nd4j-common for deeplearning4j-cuda Signed-off-by: Alex Black <blacka101@gmail.com> * nd4j-grppc-client: org.nd4j.graph -> org.nd4j.remote.grpc Signed-off-by: Alex Black <blacka101@gmail.com> * deeplearning4j-common: org.deeplearning4.* -> org.deeplearning4j.common.* Signed-off-by: Alex Black <blacka101@gmail.com> * deeplearning4j-core: org.deeplearning4j.* -> org.deeplearning.core.* Signed-off-by: Alex Black <blacka101@gmail.com> * deeplearning4j-cuda: org.deeplearning4j.nn.layers.* -> org.deeplearning4j.cuda.* Signed-off-by: Alex Black <blacka101@gmail.com> * Import fixes Signed-off-by: Alex Black <blacka101@gmail.com> * deeplearning4j-nlp-*: org.deeplearning4.text.* -> org.deeplearning4j.nlp.(language).* Signed-off-by: Alex Black <blacka101@gmail.com> * deeplearning4j-ui-model: org.deeplearning4j.ui -> org.deeplearning4j.ui.model Signed-off-by: Alex Black <blacka101@gmail.com> * datavec-spark-inference-{server/model/client}: org.datavec.spark.transform -> org.datavec.spark.inference.{server/model/client} Signed-off-by: Alex Black <blacka101@gmail.com> * datavec-jdbc: org.datavec.api -> org.datavec.jdbc Signed-off-by: Alex Black <blacka101@gmail.com> * Delete org.deeplearning4j.datasets.iterator.impl.MultiDataSetIteratorAdapter in favor of (essentially identical) org.nd4j.linalg.dataset.adapter.MultiDataSetIteratorAdapter Signed-off-by: Alex Black <blacka101@gmail.com> * ND4S fixes Signed-off-by: Alex Black <blacka101@gmail.com> * Fixes Signed-off-by: Alex Black <blacka101@gmail.com> * nd4j-common-tests: org.nd4j.* -> org.nd4j.common.tests Signed-off-by: Alex Black <blacka101@gmail.com> * Trigger CI Signed-off-by: Alex Black <blacka101@gmail.com> * Fixes Signed-off-by: Alex Black <blacka101@gmail.com> * #8878 Ignore CUDA tests on modules with 'nd4j-native under cuda' issue Signed-off-by: Alex Black <blacka101@gmail.com> * Fix bad imports in tests Signed-off-by: Alex Black <blacka101@gmail.com> * Add ignore on test (already failing) due to #8882 Signed-off-by: Alex Black <blacka101@gmail.com> * Import fixes Signed-off-by: Alex Black <blacka101@gmail.com> * Additional import fixes Signed-off-by: Alex Black <blacka101@gmail.com>
2020-04-29 11:19:26 +10:00
package org.deeplearning4j.core.parallelism;
2019-06-06 15:21:15 +03:00
import lombok.Getter;
import lombok.NonNull;
import lombok.extern.slf4j.Slf4j;
import java.util.Iterator;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.atomic.AtomicBoolean;
/**
* Asynchronous Iterator for better performance of iterators in dl4j-nn & dl4j-nlp
*
* @author raver119@gmail.com
*/
@Slf4j
public class AsyncIterator<T extends Object> implements Iterator<T> {
@Getter
protected BlockingQueue<T> buffer;
protected ReaderThread<T> thread;
protected Iterator<T> iterator;
@Getter
protected T terminator = (T) new Object();
protected T nextElement;
protected AtomicBoolean shouldWork = new AtomicBoolean(true);
public AsyncIterator(@NonNull Iterator<T> iterator, int bufferSize) {
this.buffer = new LinkedBlockingQueue<>(bufferSize);
this.iterator = iterator;
thread = new ReaderThread<>(iterator, this.buffer, terminator);
thread.start();
}
public AsyncIterator(@NonNull Iterator<T> iterator) {
this(iterator, 1024);
}
@Override
public boolean hasNext() {
try {
if (nextElement != null && nextElement != terminator) {
return true;
}
// if on previous run we've got terminator - just return false
if (nextElement == terminator)
return false;
nextElement = buffer.take();
// same on this run
return (nextElement != terminator);
2019-06-06 15:21:15 +03:00
} catch (Exception e) {
log.error("Premature end of loop!");
return false;
}
}
@Override
public T next() {
T temp = nextElement;
nextElement = temp == terminator ? terminator : null;
return temp;
}
@Override
public void remove() {
// no-op
}
public void shutdown() {
if (shouldWork.get()) {
shouldWork.set(false);
thread.interrupt();
try {
// Shutdown() should be a synchronous operation since the iterator is reset after shutdown() is
// called in AsyncLabelAwareIterator.reset().
thread.join();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
nextElement = terminator;
buffer.clear();
}
}
private class ReaderThread<T> extends Thread implements Runnable {
private BlockingQueue<T> buffer;
private Iterator<T> iterator;
private T terminator;
public ReaderThread(Iterator<T> iterator, BlockingQueue<T> buffer, T terminator) {
this.buffer = buffer;
this.iterator = iterator;
this.terminator = terminator;
setDaemon(true);
setName("AsyncIterator Reader thread");
}
@Override
public void run() {
//log.info("AsyncReader [{}] started", Thread.currentThread().getId());
try {
while (iterator.hasNext() && shouldWork.get()) {
T smth = iterator.next();
if (smth != null)
buffer.put(smth);
}
buffer.put(terminator);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
// do nothing
shouldWork.set(false);
} catch (Exception e) {
// TODO: pass that forward
throw new RuntimeException(e);
} finally {
//log.info("AsyncReader [{}] stopped", Thread.currentThread().getId());
}
}
}
}