Streaming¶
Streaming lets you observe a program's output as it is produced, rather than
waiting for the final result. dspy4s streaming is synchronous: streamify turns
a program into a function that returns a ClosableIterator[StreamEvent], which
you consume with an ordinary loop.
Consuming a stream¶
Every stream is a sequence of StreamEvents. The type is sealed, so a single
match handles every case: tokens as they arrive, status messages, the final
prediction, and errors:
def consume(stream: ClosableIterator[StreamEvent]): Option[RawPrediction] =
var finalPrediction: Option[RawPrediction] = None
while stream.hasNext do
stream.next() match
case t: TokenEvent => println(s"Output token of field ${t.fieldName}: ${t.chunk}")
case s: StatusEvent => println(s.message)
case p: PredictionEvent => finalPrediction = Some(p.prediction)
case e: ErrorEvent => println(s"Error: ${e.error.message}")
finalPrediction
Streaming a program¶
Wrap a program with Streamify.streamify and declare which output fields to
stream with a StreamListener. Calling the result returns the event iterator:
def streamAnswer(question: String)(using RuntimeContext): Option[RawPrediction] =
val predict = DynamicPredict(layout = Signature.fromString("question -> answer").layout)
val streamPredict = Streamify.streamify(
program = predict,
streamListeners = Vector(StreamListener("answer"))
)
consume(streamPredict(DynamicValues.recordFromEntries(Vector("question" := question))))
The listener names the field (answer) you want streamed token by token. The
same approach works on a composite program: add a
listener per field, and disambiguate predictors that emit the same field name
with predictName.
Notes¶
- The language model must support streaming. The bundled OpenAI provider does.
- Streaming is synchronous, so there is no async iteration to manage; the
consumer is a plain
whileloop over the iterator.
Next: Coming from DSPy.