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:
DatasetClient.stream_items- Stream dataset items incrementally. Yields a raw streamingHttpResponse.KeyValueStoreClient.stream_record- Stream a key-value store record as raw data. Yields adictwith thekey,value, andcontent_typefields, wherevalueholds the raw streaming response, orNonewhen the record doesn't exist.LogClient.stream- Stream logs in real time. Yields a raw streamingHttpResponse, orNonewhen the log doesn't exist.
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:
- Async client
- Sync client
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)
from apify_client import ApifyClient
TOKEN = 'MY-APIFY-TOKEN'
def main() -> None:
apify_client = ApifyClient(TOKEN)
run_client = apify_client.run('MY-RUN-ID')
log_client = run_client.log()
with log_client.stream() as log_stream:
if log_stream:
for bytes_chunk in log_stream.iter_bytes():
print(bytes_chunk)
Streaming uploads
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.IOBasestream: an open file, anio.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
readmethod. The client can't tell how thatreadbehaves, so it callsread()once and sends the result as a single chunk, which holds the whole source in memory. - An iterable of
bytesorstrchunks: 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. Astr,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
HttpResponsefrom 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.
- Async client
- Sync client
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'
)
from pathlib import Path
from apify_client import ApifyClient
TOKEN = 'MY-APIFY-TOKEN'
def main() -> None:
apify_client = ApifyClient(TOKEN)
kvs_client = apify_client.key_value_store('MY-KVS-ID')
# The file is read in chunks as it uploads, so its size doesn't matter.
with Path('backup.tar.gz').open('rb') as backup:
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:
- Async client
- Sync client
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')
from collections.abc import Iterator
from apify_client import ApifyClient
TOKEN = 'MY-APIFY-TOKEN'
def csv_chunks() -> Iterator[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))
def main() -> None:
apify_client = ApifyClient(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.
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.IOBasestream, such as an open file or anio.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
readmethod outsideio.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.