Interface OOCStream<T>

All Superinterfaces:
OOCStreamable<T>
All Known Implementing Classes:
FilteredOOCStream, MergedOOCStream, PlaybackStream, SourceOOCStream, SplittingOOCStream, SubOOCStream, SubscribableTaskQueue

public interface OOCStream<T> extends OOCStreamable<T>
  • Method Details

    • eos

    • enqueue

      void enqueue(T t)
    • enqueue

      void enqueue(OOCStream.QueueCallback<T> callback)
    • dequeue

      T dequeue()
    • dequeueCB

    • closeInput

      void closeInput()
    • propagateFailure

      void propagateFailure(DMLRuntimeException re)
    • setSubscriber

      void setSubscriber(Consumer<OOCStream.QueueCallback<T>> subscriber)
      Registers 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.