Class MergedOOCStream<T>
java.lang.Object
org.apache.sysds.runtime.ooc.stream.MergedOOCStream<T>
- All Implemented Interfaces:
OOCStream<T>,OOCStreamable<T>
-
Nested Class Summary
Nested classes/interfaces inherited from interface org.apache.sysds.runtime.instructions.ooc.OOCStream
OOCStream.GroupQueueCallback<T>, OOCStream.QueueCallback<T>, OOCStream.SimpleQueueCallback<T> -
Constructor Summary
ConstructorsConstructorDescriptionMergedOOCStream(List<OOCStream<T>> sources) MergedOOCStream(OOCStream<T>... sources) -
Method Summary
Modifier and TypeMethodDescriptionvoidvoidvoidvoidvoiddequeue()voidenqueue(OOCStream.QueueCallback<T> callback) voidgetData()booleanbooleanvoidvoidvoidvoidsetData(CacheableData<?> data) voidvoidsetIXTransform(BiFunction<Boolean, IndexRange, IndexRange> transform) voidsetSubscriber(Consumer<OOCStream.QueueCallback<T>> subscriber) Registers a new subscriber that consumes the stream.void
-
Constructor Details
-
MergedOOCStream
-
MergedOOCStream
-
-
Method Details
-
enqueue
-
enqueue
-
dequeue
-
dequeueCB
-
closeInput
public void closeInput()- Specified by:
closeInputin interfaceOOCStream<T>
-
propagateFailure
- Specified by:
propagateFailurein interfaceOOCStream<T>
-
hasStreamCache
public boolean hasStreamCache()- Specified by:
hasStreamCachein interfaceOOCStreamable<T>
-
getStreamCache
- Specified by:
getStreamCachein interfaceOOCStreamable<T>
-
setSubscriber
Description copied from interface:OOCStreamRegisters a new subscriber that consumes the stream. While there is no guarantee for any specific order, the closing item LocalTaskQueue.NO_MORE_TASKS is guaranteed to be invoked after every other item has finished processing. Thus, the NO_MORE_TASKS callback can be used to free dependent resources and close output streams.- Specified by:
setSubscriberin interfaceOOCStream<T>
-
getReadStream
- Specified by:
getReadStreamin interfaceOOCStreamable<T>
-
getWriteStream
- Specified by:
getWriteStreamin interfaceOOCStreamable<T>
-
isProcessed
public boolean isProcessed()- Specified by:
isProcessedin interfaceOOCStreamable<T>
-
getDataCharacteristics
- Specified by:
getDataCharacteristicsin interfaceOOCStreamable<T>
-
getData
- Specified by:
getDatain interfaceOOCStreamable<T>
-
setData
- Specified by:
setDatain interfaceOOCStreamable<T>
-
messageUpstream
- Specified by:
messageUpstreamin interfaceOOCStreamable<T>
-
messageDownstream
- Specified by:
messageDownstreamin interfaceOOCStreamable<T>
-
setUpstreamMessageRelay
- Specified by:
setUpstreamMessageRelayin interfaceOOCStreamable<T>
-
setDownstreamMessageRelay
- Specified by:
setDownstreamMessageRelayin interfaceOOCStreamable<T>
-
addUpstreamMessageRelay
- Specified by:
addUpstreamMessageRelayin interfaceOOCStreamable<T>
-
addDownstreamMessageRelay
- Specified by:
addDownstreamMessageRelayin interfaceOOCStreamable<T>
-
clearUpstreamMessageRelays
public void clearUpstreamMessageRelays()- Specified by:
clearUpstreamMessageRelaysin interfaceOOCStreamable<T>
-
clearDownstreamMessageRelays
public void clearDownstreamMessageRelays()- Specified by:
clearDownstreamMessageRelaysin interfaceOOCStreamable<T>
-
setIXTransform
- Specified by:
setIXTransformin interfaceOOCStreamable<T>
-