Skip to main content

Crate reqwest_streams

Crate reqwest_streams 

Source
Expand description

HTTP body streaming support for reqwest, in both directions, 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(())
}

§Streaming uploads

The same formats work the other way round: give a POST or PUT a stream of items and it is encoded into the request body as it is sent, without ever holding the whole thing in memory.

use futures::stream;
use reqwest_streams::JsonStreamRequest as _;
use serde::Serialize;

#[derive(Serialize)]
struct MyTestStructure {
    some_test_field: String
}

let items = stream::iter(vec![MyTestStructure { some_test_field: "value".into() }]);

reqwest::Client::new()
    .post("http://localhost:8080/ingest")
    .json_array_stream_body(items)
    .send()
    .await?;

The Content-Type is set from the format. Use ReqwestStreamBody directly when you need options, a body for multipart, or a request built by hand.

Read ReqwestStreamBody’s caveats before using this in anger. Streaming a request body is much less universally supported than streaming a response: the body cannot be retried, a redirect silently sends an empty body, and buffering reverse proxies defeat the streaming entirely.

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

§Need server support?

axum-streams is the other half of the pair, and covers both directions too. Since its 0.29 it can also receive a streamed request body, so an upload sent with JsonStreamRequest and friends is decoded on the server by its StreamBodyFrom extractor. Both crates encode and decode through the same http-streams-core.

§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 an http_streams_core::stream span:

INFO http_streams_core::stream{format="json_array" direction="response" side="client" status=200 items=1000 bytes=28001 errors=0 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,http_streams_core=debug for a progress line about once a second, and to http_streams_core=trace for one per body chunk. Naming both targets keeps anything this crate logs itself visible alongside the shared accounting.

The target is http_streams_core rather than reqwest_streams because the accounting is shared with axum-streams; the direction and side span fields tell the cases apart.

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

Re-exports§

pub use http_streams_core;

Modules§

error
Error types for streaming responses.

Structs§

ReqwestStreamBodyarrow or csv or json or protobuf
A request body that streams a sequence of items.
ReqwestStreamBodyOptionsarrow or csv or json or protobuf
Options for a streamed request body.
ReqwestStreamOptionsarrow or csv or json or protobuf
Options shared by every streaming format.
ReqwestStreamProgressarrow or csv or json or protobuf
A snapshot of one stream’s accounting.

Enums§

ReqwestStreamOutcomearrow or csv or json or protobuf
How a stream ended.

Traits§

ArrowIpcStreamRequestarrow
Extension trait for reqwest::RequestBuilder that streams an Arrow IPC request body.
ArrowIpcStreamResponsearrow
Extension trait for reqwest::Response that provides streaming support for the Apache Arrow IPC format.
CsvStreamRequestcsv
Extension trait for reqwest::RequestBuilder that streams a CSV request body.
CsvStreamResponsecsv
Extension trait for reqwest::Response that provides streaming support for the CSV format.
JsonStreamRequestjson
Extension trait for reqwest::RequestBuilder that streams a JSON request body.
JsonStreamResponsejson
Extension trait for reqwest::Response that provides streaming support for the JSON array and JSON Lines (NL/NewLines) formats.
ProtobufStreamRequestprotobuf
Extension trait for reqwest::RequestBuilder that streams a protobuf request body.
ProtobufStreamResponseprotobuf
Extension trait for reqwest::Response that provides streaming support for the Protobuf format.
StreamBodyRequestarrow or csv or json or protobuf
Sets a streamed body and its Content-Type on a request in one step.

Type Aliases§

ReqwestStreamErrorHandlerarrow or csv or json or protobuf
Called for every error reported by a stream.
ReqwestStreamProgressHandlerarrow or csv or json or protobuf
Called for every progress report, interim and terminal.
StreamBodyResult
Alias for the Result type returned by streaming responses.