Skip to main content
Version: v5

Streaming responses

A streaming command returns rows as they arrive instead of a finished result set. The provider's fetcher produces an async iterator of typed rows, the Python Interface hands back an OBBStream handle over that iterator, and the REST API sends the rows as a text/event-stream response. OBBStream is a sibling of OBBject, not a subclass, and a command declares one or the other through its return annotation.

No V5 package currently ships a streaming command. The examples on this page use a placeholder extension that registers a router named my_stream and a provider named my_provider with a model called MyTicks.

Fetcher contract​

A fetcher becomes a streaming fetcher by overriding stream_data instead of transform_data. Its second type parameter is AsyncIterator[<row model>], and the engine unwraps that to find the row type for docs and schemas.

aextract_data should return an async iterator without consuming it. Defining an inner async generator and returning it, as below, means the upstream connection (a websocket, SSE feed, or polling loop) opens only when the caller starts iterating. stream_data receives that iterator as data and must be an async def generator that yields validated rows.

import asyncio
from collections.abc import AsyncIterator
from typing import Any

from openbb_core.provider.abstract.data import Data
from openbb_core.provider.abstract.fetcher import Fetcher
from openbb_core.provider.abstract.query_params import QueryParams
from pydantic import Field


class MyTicksQueryParams(QueryParams):
"""Tick stream query."""

symbol: str = Field(description="Symbol to stream.")
limit: int = Field(default=5, description="Number of ticks to emit.")


class MyTicksData(Data):
"""One tick."""

symbol: str = Field(description="Symbol.")
sequence: int = Field(description="Tick sequence number.")
price: float = Field(description="Last price.")


class MyTicksFetcher(Fetcher[MyTicksQueryParams, AsyncIterator[MyTicksData]]):
"""Tick stream fetcher."""

require_credentials = False

@staticmethod
def transform_query(params: dict[str, Any]) -> MyTicksQueryParams:
"""Validate the query."""
return MyTicksQueryParams(**params)

@staticmethod
async def aextract_data(
query: MyTicksQueryParams,
credentials: dict[str, str] | None,
**kwargs: Any,
) -> AsyncIterator[dict]:
"""Return an async iterator of raw messages."""

async def messages():
for sequence in range(query.limit):
await asyncio.sleep(0.1)
yield {"symbol": query.symbol, "sequence": sequence, "price": 100 + sequence}

return messages()

@staticmethod
async def stream_data(
query: MyTicksQueryParams,
data: AsyncIterator[dict],
**kwargs: Any,
) -> AsyncIterator[MyTicksData]:
"""Validate each raw message."""
async for message in data:
yield MyTicksData.model_validate(message)

MyTicksFetcher.is_streaming is True because stream_data is overridden. await MyTicksFetcher.fetch_data({"symbol": "ABC"}) returns the row iterator, and MyTicksFetcher.test({"symbol": "ABC", "limit": 2}) consumes the first row and checks that it is a MyTicksData with every declared field.

Router contract​

A streaming command is annotated -> OBBStream and returns await OBBStream.from_query(Query(**locals())). Otherwise it looks like any other provider-backed command.

from openbb_core.app.model.command_context import CommandContext
from openbb_core.app.model.example import APIEx
from openbb_core.app.model.stream import OBBStream
from openbb_core.app.provider_interface import (
ExtraParams,
ProviderChoices,
StandardParams,
)
from openbb_core.app.query import Query
from openbb_core.app.router import Router

router = Router(prefix="", description="Streaming example.")


@router.command(
model="MyTicks",
examples=[APIEx(parameters={"symbol": "ABC", "provider": "my_provider"})],
)
async def ticks(
cc: CommandContext,
provider_choices: ProviderChoices,
standard_params: StandardParams,
extra_params: ExtraParams,
) -> OBBStream:
"""Stream ticks for a symbol."""
return await OBBStream.from_query(Query(**locals()))

For an OBBStream return annotation, @router.command registers the route without a response model and documents a 200 response with text/event-stream content in the OpenAPI schema. OBBStream.from_query copies the provider name onto the handle and, if the result is an AnnotatedResult, stores its metadata in extra["results_metadata"].

Using a stream from Python​

The command returns immediately with an OBBStream. Nothing is fetched until the stream is iterated or started. Iterate it inside async code:

import asyncio

from openbb import obb


async def main():
async for tick in obb.my_stream.ticks(symbol="ABC", limit=3):
print(tick.sequence, tick.price)


asyncio.run(main())

start() drives the stream on a daemon thread with its own event loop and returns the handle, so a notebook or script keeps running while rows arrive. Without arguments, each row is written to standard output. stop() cancels the stream and waits up to timeout seconds (default 5.0) for the thread to exit:

from openbb import obb

stream = obb.my_stream.ticks(symbol="ABC", limit=100)
stream.start()
stream.stop()

The output argument of start() also accepts a file path (opened for writing and closed when the stream ends), an open file or pipe, or a socket; anything without a write or sendall method raises TypeError. The handler argument receives each row before it is written. It can be sync or async. Returning None drops the row, and returning any other value writes that value instead. wait() blocks until the stream finishes, or until its optional timeout, and still responds to Ctrl+C while waiting.

from openbb import obb


def even_only(tick):
return tick if tick.sequence % 2 == 0 else None


stream = obb.my_stream.ticks(symbol="ABC", limit=10)
stream.start(output="ticks.ndjson", handler=even_only)
stream.wait()

Rows are serialized one per line. Pydantic rows are written with model_dump_json(), other objects with json.dumps(..., default=str), and strings or bytes pass through unchanged, so a fetcher can yield pre-formatted SSE frames.

The underlying iterator is single-use. After a stream has been iterated or started, calling start() again raises RuntimeError; call the command again for a new stream. The first call to start() in a process registers an exit hook that stops every running stream. If that call happens on the main thread, it also installs SIGINT and SIGTERM handlers that stop the streams and then defer to the handler that was installed before.

The handle carries id (a UUIDv7 string), provider, warnings raised while the command was opening the stream, extra (including extra["metadata"] when the metadata preference is on), and media_type. The started and running properties report whether iteration has begun and whether the background thread is alive.

OBBStream has none of the to_* conversion methods, so leave output_type at its default "OBBject" when calling streaming commands. A command whose handler returns a Starlette StreamingResponse directly, such as a wrapped FastAPI route, is also returned as an OBBStream over the response body.

Using a stream from the REST API​

The REST route returns a Starlette StreamingResponse. The body is the same line-by-line serialization described above, and the envelope fields that do not fit in the body are sent as headers:

HeaderValue
X-OpenBB-Stream-IdThe stream id
X-OpenBB-ProviderThe provider name, when set
X-OpenBB-WarningThe warnings as a JSON array, when any were raised
curl -N "http://127.0.0.1:8000/api/v1/my_stream/ticks?symbol=ABC&limit=3"

The media type is text/event-stream unless the wrapped source declares another. A command that returns a Starlette Response of its own is forwarded unchanged.