Skip to main content

Crate reqwest_streams

Crate reqwest_streams 

Source
Expand description

Streaming responses support for reqwest for different formats:

  • JSON array stream format
  • JSON Lines (NL/NewLines) format
  • CSV stream format
  • Protobuf len-prefixed stream format
  • Apache Arrow IPC stream format

This type of responses are useful when you are reading huge stream of objects from some source (such as database, file, etc) and want to avoid huge memory allocations to store on the server side.

§Features

Note: The default features do not include any formats.

  • json: JSON array and JSON Lines (JSONL) stream formats
  • csv: CSV stream format
  • protobuf: Protobuf len-prefixed stream format
  • arrow: Apache Arrow IPC stream format
  • tracing: report progress and errors through tracing

§Example

use futures::stream::BoxStream as _;
use reqwest_streams::JsonStreamResponse as _;
use serde::Deserialize;

#[derive(Debug, Clone, Deserialize)]
struct MyTestStructure {
    some_test_field: String
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {

    let _stream = reqwest::get("http://localhost:8080/json-array")
        .await?
        .json_array_stream::<MyTestStructure>(1024);

    Ok(())
}

More and complete examples available on the github in the examples directory.

§Need server support?

There is the same functionality:

§Observing stream errors

An error that happens mid-stream is yielded as an item, so a consumer that stops at the first one silently gets a truncated result. Use ReqwestStreamOptions::on_error to observe them, or enable the tracing feature to have them logged at ERROR on the reqwest_streams target.

§Observing stream progress

Nothing at the call site can tell you how much of a response actually arrived, because it is read long after the call returned. With the tracing feature every stream reports its totals once at INFO when it ends, on a reqwest_streams::response_stream span:

INFO reqwest_streams::response_stream{format="json_array" status=200 items=1000 bytes=28001 elapsed_ms=11239 outcome="completed"}: Finished streaming an HTTP body

The outcome tells apart the three ways a stream can end: completed, aborted (the consumer stopped reading early) and failed, which reports at ERROR instead. Raise the filter to reqwest_streams=debug for a progress line about once a second, and to reqwest_streams=trace for one per body chunk.

The same accounting is available without tracing, for metrics, via ReqwestStreamOptions::on_progress.

Modules§

error
Error types for streaming responses.

Structs§

ReqwestStreamOptions
Options shared by every streaming format.
ReqwestStreamProgress
A snapshot of how much of a streamed response has been read so far.

Enums§

ReqwestStreamOutcome
How a streamed response ended, or that it is still going.

Traits§

ArrowIpcStreamResponsearrow
Extension trait for reqwest::Response that provides streaming support for the Apache Arrow IPC format.
CsvStreamResponsecsv
Extension trait for reqwest::Response that provides streaming support for the CSV format.
JsonStreamResponsejson
Extension trait for reqwest::Response that provides streaming support for the JSON array and JSON Lines (NL/NewLines) formats.
ProtobufStreamResponseprotobuf
Extension trait for reqwest::Response that provides streaming support for the Protobuf format.

Type Aliases§

ReqwestStreamErrorHandler
A callback invoked for every error produced while reading a streamed response.
ReqwestStreamProgressHandler
A callback invoked with progress snapshots while reading a streamed response.
StreamBodyResult
Alias for the Result type returned by streaming responses.