Streaming
Streaming¶
If the operation's input is a streaming blob, you can pass any AsyncIterator[bytes] or just a bytes object.
from capo_s3 import AsyncS3Client
s3_client = AsyncS3Client()
response = await s3_client.put_object("bucket_name", "key", body=b"some binary data")
Or, if you don't want to load the entire blob into memory, you can pass an AsyncIterator[bytes]:
from capo_s3 import AsyncS3Client
async def async_iterator():
yield b"capo"
response = await s3_client.put_object("bucket_name", "key", body=async_iterator(), content_length=4)
As you might have noticed, we also passed the content_length. That's an AWS requirement when using streaming inputs; it must always know the length of the blob before sending it to AWS.
Note that the stream can be any iterator of bytes; it need not be the file's content. You can stream any data you want, for example, directly from the HTTP response of another service, or from a database, etc.
The catch with a plain iterator is that it can be sent only once. If the request fails after the body was transmitted (a throttling error, a dropped connection), there is nothing left to resend, so the operation is not retried. To get retries for streamed uploads, pass a Body instead: it wraps a source that can be reopened, and every attempt streams a fresh copy. Body.from_path (sync client) and Body.async_from_path (async client) stream a file from disk and take the content_length from the file size, so you don't need to pass it:
from capo_s3 import AsyncS3Client, Body
s3_client = AsyncS3Client()
response = await s3_client.put_object("bucket_name", "key", body=Body.async_from_path("data.bin"))
For sources other than files, build a Body from an opener — a context manager that yields a (stream, length) pair each time it is entered. The SDK enters it before every attempt and exits it when the operation finishes:
from contextlib import asynccontextmanager
from capo_s3 import AsyncS3Client, Body
@asynccontextmanager
async def open_rows():
rows = await db.fetch_all() # re-read from your data source on every attempt
async def chunks():
for row in rows:
yield row
yield chunks(), sum(len(row) for row in rows)
response = await s3_client.put_object("bucket_name", "key", body=Body(open_rows))
The output as mentioned before also can be a stream, in such case, the operation will return a context manager that yield the response, ensuring that the resource is properly closed after the response is consumed.
from capo_s3 import AsyncS3Client
s3_client = AsyncS3Client()
async with s3_client.get_object("bucket_name", "key") as response:
async for chunk in response["body"]:
print(chunk)
The event streaming operations are similar, but instead of using AsyncIterator[bytes], they use AsyncIterator[Event], where Event is a TypedDict that represents the event type.
from capo_s3 import AsyncS3Client
s3_client = AsyncS3Client()
async def main():
async with s3_client.select_object_content(
"bucket_name",
"key",
expression="SELECT * FROM S3Object s WHERE s._1 > 100",
expression_type="SQL",
input_serialization={
"csv": {"file_header_info": "NONE"},
"compression_type": "NONE",
},
output_serialization={"csv": {}},
) as response:
async for event in response["payload"]:
if "End" in event:
print(event["End"])