Skip to main content
Version: Next

Streaming resources

Large resources don't have to pass through memory whole. The client streams downloads of dataset items, key-value store records, and logs as they arrive, and it streams uploads of key-value store records and Actor inputs from a file or another stream as they're read. Both directions work with the synchronous and the asynchronous client. A download can also feed an upload directly, so one Actor's output becomes another Actor's input without passing through memory.

Streaming downloads

Dataset items, key-value store records, and logs can be streamed directly from the Apify API, so you process them incrementally instead of downloading them whole. The supported methods are:

All three methods are context managers. Consume the streamed data within a with block, so the connection is closed automatically.

The following example shows how to stream the logs of an Actor run incrementally:

from apify_client import ApifyClientAsync

TOKEN = 'MY-APIFY-TOKEN'


async def main() -> None:
apify_client = ApifyClientAsync(TOKEN)
run_client = apify_client.run('MY-RUN-ID')
log_client = run_client.log()

async with log_client.stream() as log_stream:
if log_stream:
async for bytes_chunk in log_stream.aiter_bytes():
print(bytes_chunk)

Streaming uploads

Experimental feature

Streaming uploads are experimental. Their behavior and interface may change in future releases.

KeyValueStoreClient.set_record and the run_input of ActorClient.start, ActorClient.call, and RunClient.metamorph accept a value the client streams to the API in chunks:

  • An io.IOBase stream: an open file, an io.BytesIO, a member of a ZIP archive, or a pipe. The client reads it in chunks of 64 KiB, or 64 Ki characters for a text-mode file, from its current position and doesn't close it.
  • Any other object with a read method. The client can't tell how that read behaves, so it calls read() once and sends the result as a single chunk, which holds the whole source in memory.
  • An iterable of bytes or str chunks: an iterator such as a generator, or an object that only implements __iter__, for example a class that yields chunks from a generator method. The iterable decides the chunk sizes. An iterator is consumed by the upload, so passing the same one to a second upload sends an empty body. A generator is closed once the upload ends, even when the upload fails before reaching its end. A str, bytes, list, tuple, dict, or other container, and a pydantic model, are never streamed even though they can be iterated - they are the values the client uploads whole or serializes as JSON.
  • A streaming HttpResponse from one of the download methods, whose body is forwarded chunk by chunk.

The asynchronous client also accepts an async iterable - an async generator or an object that only implements __aiter__ - and an object with an async def read, such as a file opened with aiofiles, which is read whole like any other non-io.IOBase object. The synchronous client rejects those with a TypeError. A synchronous file or iterator works with both clients. The asynchronous client reads it in a worker thread, so the event loop stays free.

Apart from a non-io.IOBase object with a read method, only the chunk being sent is held in memory, so the process uploading a file doesn't need memory for the whole file. What the API accepts doesn't change: an Actor input is still capped at 9 MB, and a streamed body over that size is rejected with a 413 during the upload.

The following example uploads a file from disk. Both variants pass a plain file, which the asynchronous client reads in a worker thread.

import asyncio
from pathlib import Path

from apify_client import ApifyClientAsync

TOKEN = 'MY-APIFY-TOKEN'


async def main() -> None:
apify_client = ApifyClientAsync(TOKEN)
kvs_client = apify_client.key_value_store('MY-KVS-ID')

# The file is read in chunks in a worker thread as it uploads, so its size
# doesn't matter.
backup = await asyncio.to_thread(Path('backup.tar.gz').open, 'rb')
with backup:
await kvs_client.set_record(
'backup.tar.gz', backup, content_type='application/gzip'
)

With a generator, you produce the data while it uploads, for example from pages of a database query:

from collections.abc import AsyncIterator

from apify_client import ApifyClientAsync

TOKEN = 'MY-APIFY-TOKEN'


async def csv_chunks() -> AsyncIterator[str]:
"""Build the CSV in pieces, for example from pages of a database query."""
yield 'id,value\n'
for start in range(0, 1_000_000, 10_000):
yield ''.join(f'{i},{i * i}\n' for i in range(start, start + 10_000))


async def main() -> None:
apify_client = ApifyClientAsync(TOKEN)
kvs_client = apify_client.key_value_store('MY-KVS-ID')

# Each chunk is encoded and sent as it's produced. An iterator can't be
# rewound, so a failed upload isn't retried.
await kvs_client.set_record('report.csv', csv_chunks(), content_type='text/csv')

Content type

A file the standard library opened in text mode is uploaded as text/plain; charset=utf-8, encoded to UTF-8 chunk by chunk. Any other streamed value is uploaded as application/octet-stream unless you pass content_type, an aiofiles text file included - the client can't tell it from a binary one. An iterator or a response carries no type of its own, so set content_type explicitly, as the examples do.

Compression

A streamed body is never compressed, because the client sees each chunk only as it's sent. To upload compressed data, compress it yourself and declare the encoding with content_encoding. The client then forwards the bytes and the header as they are. For details, see Pre-compressed bodies.

Retries

The client retries a failed request only when it can send the body again:

  • A seekable io.IOBase stream, such as an open file or an io.BytesIO, is sought back to where it started before each retry, so the upload is retried like any other request.
  • An iterator, a streaming response, a pipe, an object with a read method outside io.IOBase, or any asynchronous source is consumed by the attempt that sends it. The request gets a single attempt, and a failure raises immediately.
  • A failure inside the source itself, such as a disk error while reading the file, is raised as that error and isn't retried, whatever the source.

If an upload from a source that can't be rewound needs retries, download the data to a temporary file first and upload the file. The object tempfile.NamedTemporaryFile() returns, which is also what tempfile.TemporaryFile() returns on Windows, isn't an io.IOBase stream, so pass its file attribute to have it read in chunks and retried. In the asynchronous client, prefer a file opened with open() over one opened with aiofiles: the plain file is read in chunks in a worker thread and can be sought back for a retry, while the aiofiles file is read whole in one call and gets a single attempt.

Time limits

The API ends a request whose body doesn't arrive in full within about 5 minutes with a 408 Request Timeout. The client doesn't retry the error and logs a warning that names the limit.

The client's own timeout applies to an upload as well. With the default Impit HTTP client, the timeout covers the whole request, so set_record gives up on an upload after its default of 360 seconds, and ActorClient.start after 30 seconds. For a slow upload that still fits within the API's limit, pass a longer timeout.

To upload data that a source produces slowly, such as the results of a long-running query, write it to a temporary file first and upload the file. The upload then runs at network speed, and a failed attempt can be retried. For details, see Retries.

Custom HTTP clients

A custom transport receives a streamed body as an iterator of bytes chunks, or as an async iterator in the asynchronous client, so any HTTP library that sends iterable bodies can send it. The StreamedRequestBody class documents the rest of the contract. For details, see The transport contract.