Skip to content

anystore.io

IOFormat

Bases: StrEnum

For use in typer cli

Source code in anystore/io/write.py
class IOFormat(StrEnum):
    """For use in typer cli"""

    csv = "csv"
    json = "json"

ModelWriter

Bases: Writer

A generic writer for pydantic objects to any out uri, either json or csv

Source code in anystore/io/write.py
class ModelWriter(Writer):
    """
    A generic writer for pydantic objects to any out uri, either json or csv
    """

    def write(self, row: BaseModel) -> None:
        data = row.model_dump(by_alias=True, mode="json")
        return super().write(data)

ProgressTask

A labelled bar on a SyncProgressBar.

Use it as a context manager to have it disappear when its work is done:

Example
with bar.task("chunk 1", total=len(keys)) as task:
    for key in keys:
        task.advance(size=len(store.get(key)))
Source code in anystore/io/progress.py
class ProgressTask:
    """
    A labelled bar on a [`SyncProgressBar`][anystore.io.progress.SyncProgressBar].

    Use it as a context manager to have it disappear when its work is done:

    Example:
        ```python
        with bar.task("chunk 1", total=len(keys)) as task:
            for key in keys:
                task.advance(size=len(store.get(key)))
        ```
    """

    def __init__(
        self, bar: "SyncProgressBar", task_id: TaskID, throughput: Throughput
    ) -> None:
        self.bar = bar
        self.task_id = task_id
        self.throughput = throughput

    def advance(self, items: int | None = 1, size: int | None = None) -> None:
        """
        Advance the task, optionally accounting transferred bytes.

        Args:
            items: Number of items completed (default: 1)
            size: Bytes transferred for them, added to this task's throughput
                and to the overall throughput of the bar
        """
        if size:
            self.throughput.add(size)
            # a task can be fed the bar's own throughput to sum up the others
            if self.throughput is not self.bar.throughput:
                self.bar.throughput.add(size)
        self.bar.progress.advance(self.task_id, items or 0)

    def update(self, **kwargs: Any) -> None:
        """Update the underlying rich task, e.g. `total` or `description`"""
        self.bar.progress.update(self.task_id, **kwargs)

    def remove(self) -> None:
        """Remove the bar from the display"""
        self.bar.progress.remove_task(self.task_id)

    def __enter__(self) -> Self:
        return self

    def __exit__(self, *args: Any) -> None:
        self.remove()

advance(items=1, size=None)

Advance the task, optionally accounting transferred bytes.

Parameters:

Name Type Description Default
items int | None

Number of items completed (default: 1)

1
size int | None

Bytes transferred for them, added to this task's throughput and to the overall throughput of the bar

None
Source code in anystore/io/progress.py
def advance(self, items: int | None = 1, size: int | None = None) -> None:
    """
    Advance the task, optionally accounting transferred bytes.

    Args:
        items: Number of items completed (default: 1)
        size: Bytes transferred for them, added to this task's throughput
            and to the overall throughput of the bar
    """
    if size:
        self.throughput.add(size)
        # a task can be fed the bar's own throughput to sum up the others
        if self.throughput is not self.bar.throughput:
            self.bar.throughput.add(size)
    self.bar.progress.advance(self.task_id, items or 0)

remove()

Remove the bar from the display

Source code in anystore/io/progress.py
def remove(self) -> None:
    """Remove the bar from the display"""
    self.bar.progress.remove_task(self.task_id)

update(**kwargs)

Update the underlying rich task, e.g. total or description

Source code in anystore/io/progress.py
def update(self, **kwargs: Any) -> None:
    """Update the underlying rich task, e.g. `total` or `description`"""
    self.bar.progress.update(self.task_id, **kwargs)

SmartHandler

Source code in anystore/io/handler.py
class SmartHandler:
    def __init__(
        self,
        uri: Uri,
        compression: CompressKind | str | None = None,
        **kwargs: Any,
    ) -> None:
        self.uri = uri
        self.is_buffer = self.uri == "-"
        kwargs["mode"] = kwargs.get("mode", DEFAULT_MODE)
        self.mode = kwargs["mode"]
        self.compression = compression
        # stdio and an injected handle need the codec applied here; a uri goes
        # through `Store.open`, which applies it at the funnel – doing both
        # would encode the stream twice
        self.sys_io = _get_sysio(binary_mode(self.mode) if compression else self.mode)
        self.kwargs = kwargs
        self.handler: IO | None = None

    def open(self) -> IO[Any]:
        try:
            if self.is_buffer:
                return self._wrap(self.sys_io)
            elif isinstance(self.uri, (BytesIO, StringIO, IOBase)):
                return self._wrap(self.uri)
            else:
                resource = UriResource(self.uri)
                mode = self.kwargs.pop("mode", DEFAULT_MODE)
                self.handler = resource.open(
                    mode, compression=self.compression, **self.kwargs
                ).__enter__()
                return self.handler
        except FileNotFoundError as e:
            raise DoesNotExist(str(e))

    def _wrap(self, io: IO[Any]) -> IO[Any]:
        """Apply the codec to a handle this class did not open.

        `own=False`, so closing the codec flushes its frame and leaves the
        handle for whoever owns it – stdout, or the caller who passed it in.
        The codec itself does become ours to close, so it is tracked.
        """
        if self.compression is None:
            return io
        self.handler = open_codec(io, self.compression, self.mode)
        return self.handler

    def close(self):
        # a tracked handler is always ours to close: either the store's stream
        # or a codec we layered over someone else's handle
        if self.handler is not None:
            self.handler.close()

    def __enter__(self):
        return self.open()

    def __exit__(self, *args, **kwargs) -> None:
        self.close()

SyncProgressBar

A live display of one or more bars, with byte throughput.

For a single bar, describe it up front and advance the display itself:

Example
with SyncProgressBar("download", total=len(keys)) as bar:
    for key in keys:
        bar.advance(size=len(store.get(key)))

For concurrent work, add a bar per unit of it with task – they may be advanced from any thread, and all of them live in this one display, as wrapping each bar in a display of its own would leave all but the first invisible (rich only renders the outermost live display):

Example
with SyncProgressBar() as bar:
    for prefix in prefixes:
        with bar.task(prefix, total=len(keys)) as task:
            for key in keys:
                task.advance(size=len(store.get(key)))

Parameters:

Name Type Description Default
description str | None

Label for the display's own bar – leave it out for the multi-task case, where each task brings its own

None
total int | None

Number of items for that bar, if known

None
console Console | None

Console to render on (default: a new one on the current sys.stderr)

None
columns Sequence[str | ProgressColumn] | None

Custom rich progress columns

None
transient bool | None

Remove the whole display when it's done (default: yes)

True
route_logging bool | None

Print log output through the same console while the display is running (default: yes)

True
disable bool | None

Render nothing at all, e.g. for a --quiet flag

False
window int | None

Seconds to average throughput rates over

DEFAULT_WINDOW
Source code in anystore/io/progress.py
class SyncProgressBar:
    """
    A live display of one or more bars, with byte throughput.

    For a single bar, describe it up front and advance the display itself:

    Example:
        ```python
        with SyncProgressBar("download", total=len(keys)) as bar:
            for key in keys:
                bar.advance(size=len(store.get(key)))
        ```

    For concurrent work, add a bar per unit of it with
    [`task`][anystore.io.progress.SyncProgressBar.task] – they may be advanced
    from any thread, and all of them live in this one display, as wrapping each
    bar in a display of its own would leave all but the first invisible (rich
    only renders the outermost live display):

    Example:
        ```python
        with SyncProgressBar() as bar:
            for prefix in prefixes:
                with bar.task(prefix, total=len(keys)) as task:
                    for key in keys:
                        task.advance(size=len(store.get(key)))
        ```

    Args:
        description: Label for the display's own bar – leave it out for the
            multi-task case, where each task brings its own
        total: Number of items for that bar, if known
        console: Console to render on (default: a new one on the current
            `sys.stderr`)
        columns: Custom rich progress columns
        transient: Remove the whole display when it's done (default: yes)
        route_logging: Print log output through the same console while the
            display is running (default: yes)
        disable: Render nothing at all, e.g. for a `--quiet` flag
        window: Seconds to average throughput rates over
    """

    def __init__(
        self,
        description: str | None = None,
        total: int | None = None,
        console: Console | None = None,
        columns: Sequence[str | ProgressColumn] | None = None,
        transient: bool | None = True,
        route_logging: bool | None = True,
        disable: bool | None = False,
        window: int | None = DEFAULT_WINDOW,
    ) -> None:
        self.throughput = Throughput(window)
        self.description = description
        self.total = total
        self.console = console
        self.columns = columns or DEFAULT_COLUMNS
        self.transient = bool(transient)
        self.route_logging = bool(route_logging)
        self.disable = bool(disable)
        self.window = window
        self._progress: Progress | None = None
        self._task: ProgressTask | None = None
        self._stack = ExitStack()

    @property
    def progress(self) -> Progress:
        """The underlying rich `Progress`, once the display is started"""
        if self._progress is None:
            raise RuntimeError("Progress display not started")
        return self._progress

    @property
    def default_task(self) -> ProgressTask:
        """The display's own bar, for the single-bar case"""
        if self._task is None:
            raise RuntimeError(
                "No bar of its own, give the `SyncProgressBar` a description "
                "or add a `task()`"
            )
        return self._task

    def advance(self, items: int | None = 1, size: int | None = None) -> None:
        """
        Advance the display's own bar, optionally accounting transferred bytes.

        Only for the single-bar case – with several tasks, advance those.

        Args:
            items: Number of items completed (default: 1)
            size: Bytes transferred for them
        """
        self.default_task.advance(items, size)

    def update(self, **kwargs: Any) -> None:
        """Update the display's own bar, e.g. `total` or `description`"""
        self.default_task.update(**kwargs)

    def task(
        self,
        description: str,
        total: int | None = None,
        throughput: Throughput | None = None,
    ) -> ProgressTask:
        """
        Add a labelled bar to the display.

        Args:
            description: The label, e.g. the key prefix being worked on
            total: Number of items, if known – an unknown total pulses instead
                of filling up
            throughput: Byte counter to display, pass the bar's own
                `throughput` for a task that sums up all the others (default: a
                fresh one for this task)

        Returns:
            The task handle, removable and usable as a context manager
        """
        throughput = throughput if throughput is not None else Throughput(self.window)
        task_id = self.progress.add_task(
            description, total=total, throughput=throughput
        )
        return ProgressTask(self, task_id, throughput)

    def start(self) -> Self:
        """Start rendering (and, unless turned off, routing log output)"""
        # bind the console late: it holds on to the current `sys.stderr`, which
        # `logging_through` is about to replace
        console = self.console or Console(file=sys.stderr)
        self.console = console
        self._progress = Progress(
            *self.columns,
            console=console,
            transient=self.transient,
            disable=self.disable,
            # `logging_through` routes stderr through this console already, and
            # rich's own redirect would fight it
            redirect_stdout=False,
            redirect_stderr=False,
        )
        if self.route_logging:
            self._stack.enter_context(logging_through(console))
        self._stack.enter_context(self._progress)
        if self.description is not None or self.total is not None:
            # one counter for a single bar: the task shares the bar's own
            self._task = self.task(
                self.description or "", total=self.total, throughput=self.throughput
            )
        return self

    def stop(self) -> None:
        """Stop rendering and restore log output"""
        self._progress = None
        self._task = None
        self._stack.close()

    def __enter__(self) -> Self:
        return self.start()

    def __exit__(self, *args: Any) -> None:
        self.stop()

default_task property

The display's own bar, for the single-bar case

progress property

The underlying rich Progress, once the display is started

advance(items=1, size=None)

Advance the display's own bar, optionally accounting transferred bytes.

Only for the single-bar case – with several tasks, advance those.

Parameters:

Name Type Description Default
items int | None

Number of items completed (default: 1)

1
size int | None

Bytes transferred for them

None
Source code in anystore/io/progress.py
def advance(self, items: int | None = 1, size: int | None = None) -> None:
    """
    Advance the display's own bar, optionally accounting transferred bytes.

    Only for the single-bar case – with several tasks, advance those.

    Args:
        items: Number of items completed (default: 1)
        size: Bytes transferred for them
    """
    self.default_task.advance(items, size)

start()

Start rendering (and, unless turned off, routing log output)

Source code in anystore/io/progress.py
def start(self) -> Self:
    """Start rendering (and, unless turned off, routing log output)"""
    # bind the console late: it holds on to the current `sys.stderr`, which
    # `logging_through` is about to replace
    console = self.console or Console(file=sys.stderr)
    self.console = console
    self._progress = Progress(
        *self.columns,
        console=console,
        transient=self.transient,
        disable=self.disable,
        # `logging_through` routes stderr through this console already, and
        # rich's own redirect would fight it
        redirect_stdout=False,
        redirect_stderr=False,
    )
    if self.route_logging:
        self._stack.enter_context(logging_through(console))
    self._stack.enter_context(self._progress)
    if self.description is not None or self.total is not None:
        # one counter for a single bar: the task shares the bar's own
        self._task = self.task(
            self.description or "", total=self.total, throughput=self.throughput
        )
    return self

stop()

Stop rendering and restore log output

Source code in anystore/io/progress.py
def stop(self) -> None:
    """Stop rendering and restore log output"""
    self._progress = None
    self._task = None
    self._stack.close()

task(description, total=None, throughput=None)

Add a labelled bar to the display.

Parameters:

Name Type Description Default
description str

The label, e.g. the key prefix being worked on

required
total int | None

Number of items, if known – an unknown total pulses instead of filling up

None
throughput Throughput | None

Byte counter to display, pass the bar's own throughput for a task that sums up all the others (default: a fresh one for this task)

None

Returns:

Type Description
ProgressTask

The task handle, removable and usable as a context manager

Source code in anystore/io/progress.py
def task(
    self,
    description: str,
    total: int | None = None,
    throughput: Throughput | None = None,
) -> ProgressTask:
    """
    Add a labelled bar to the display.

    Args:
        description: The label, e.g. the key prefix being worked on
        total: Number of items, if known – an unknown total pulses instead
            of filling up
        throughput: Byte counter to display, pass the bar's own
            `throughput` for a task that sums up all the others (default: a
            fresh one for this task)

    Returns:
        The task handle, removable and usable as a context manager
    """
    throughput = throughput if throughput is not None else Throughput(self.window)
    task_id = self.progress.add_task(
        description, total=total, throughput=throughput
    )
    return ProgressTask(self, task_id, throughput)

update(**kwargs)

Update the display's own bar, e.g. total or description

Source code in anystore/io/progress.py
def update(self, **kwargs: Any) -> None:
    """Update the display's own bar, e.g. `total` or `description`"""
    self.default_task.update(**kwargs)

Throughput

Thread-safe byte counter exposing a rolling transfer rate.

Parameters:

Name Type Description Default
window int | None

Seconds to average the rate over

DEFAULT_WINDOW
Source code in anystore/io/progress.py
class Throughput:
    """
    Thread-safe byte counter exposing a rolling transfer rate.

    Args:
        window: Seconds to average the rate over
    """

    def __init__(self, window: int | None = DEFAULT_WINDOW) -> None:
        self.window = window or DEFAULT_WINDOW
        self.total = 0
        self._events: deque[tuple[float, int]] = deque()
        self._started = time.monotonic()
        self._lock = threading.Lock()

    def add(self, size: int) -> None:
        """Account `size` transferred bytes."""
        with self._lock:
            now = time.monotonic()
            self.total += size
            self._events.append((now, size))
            self._expire(now)

    @property
    def rate(self) -> float:
        """Bytes per second over the last `window` seconds."""
        with self._lock:
            now = time.monotonic()
            self._expire(now)
            span = min(now - self._started, self.window)
            if span <= 0:
                return 0.0
            return sum(size for _, size in self._events) / span

    def _expire(self, now: float) -> None:
        while self._events and now - self._events[0][0] > self.window:
            self._events.popleft()

rate property

Bytes per second over the last window seconds.

add(size)

Account size transferred bytes.

Source code in anystore/io/progress.py
def add(self, size: int) -> None:
    """Account `size` transferred bytes."""
    with self._lock:
        now = time.monotonic()
        self.total += size
        self._events.append((now, size))
        self._expire(now)

ThroughputColumn

Bases: ProgressColumn

Renders the byte throughput of a task's Throughput field.

The item count says how far along a task is; this says how fast bytes are actually moving.

Source code in anystore/io/progress.py
class ThroughputColumn(ProgressColumn):
    """Renders the byte throughput of a task's `Throughput` field.

    The item count says how far along a task is; this says how fast bytes are
    actually moving.
    """

    def render(self, task: Task) -> Text:
        throughput: Throughput | None = task.fields.get("throughput")
        rate = throughput.rate if throughput is not None else 0.0
        return Text(format_bytes(rate, "B/s"), style="progress.data.speed")

Writer

A generic writer for python dict objects to any out uri, either json or csv

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to write to, "-" for stdout, or an already open handle

required
mode str | None

open mode, default wb (forced to text for csv)

DEFAULT_WRITE_MODE
output_format Formats | None

csv or json (default: json)

'json'
fieldnames list[str] | None

csv header, inferred from the first row when omitted

None
clean bool | None

Apply clean_dict

False
compression CompressKind | str | None

Codec to compress the output with ("gz", "bz2", "xz", "zst", "lz4")

None
lazy bool | None

Defer creating the target to the first write, so a run that writes nothing leaves no file behind. Off by default – an empty csv carrying just its header is a legitimate thing to want.

False
**kwargs

pass through storage-specific options

{}
Source code in anystore/io/write.py
class Writer:
    """
    A generic writer for python dict objects to any out uri, either json or csv

    Args:
        uri: string or path-like key uri to write to, `"-"` for stdout, or an
            already open handle
        mode: open mode, default `wb` (forced to text for csv)
        output_format: csv or json (default: json)
        fieldnames: csv header, inferred from the first row when omitted
        clean: Apply [clean_dict][anystore.util.data.clean_dict]
        compression: Codec to compress the output with ("gz", "bz2", "xz",
            "zst", "lz4")
        lazy: Defer creating the target to the first `write`, so a run that
            writes nothing leaves no file behind. Off by default – an empty
            csv carrying just its header is a legitimate thing to want.
        **kwargs: pass through storage-specific options
    """

    def __init__(
        self,
        uri: Uri,
        mode: str | None = DEFAULT_WRITE_MODE,
        output_format: Formats | None = "json",
        fieldnames: list[str] | None = None,
        clean: bool | None = False,
        compression: CompressKind | str | None = None,
        lazy: bool | None = False,
        **kwargs,
    ) -> None:
        if output_format not in (FORMAT_JSON, FORMAT_CSV):
            raise ValueError("Invalid output format, only csv or json allowed")
        mode = mode or DEFAULT_WRITE_MODE
        self.mode = mode.replace("b", "") if output_format == "csv" else mode
        self.handler = SmartHandler(
            uri, mode=self.mode, compression=compression, **kwargs
        )
        self.fieldnames = fieldnames
        self.output_format = output_format
        self.clean = clean
        self.lazy = lazy
        self.csv_writer: csv.DictWriter | None = None
        self._io: IO[Any] | None = None

    def open(self) -> IO[Any]:
        """The open target, created on first ask.

        A csv whose header is known up front writes it here, so a writer that
        never sees a row still leaves a complete – if empty – table instead of
        a zero-byte file.
        """
        if self._io is None:
            self._io = self.handler.open()
            if self.output_format == FORMAT_CSV and self.fieldnames:
                self._start_csv(self.fieldnames)
        return self._io

    def _start_csv(self, fieldnames: Iterable[str]) -> None:
        """Bind the csv writer to the open target and write its header."""
        self.csv_writer = csv.DictWriter(self.io, fieldnames)
        self.csv_writer.writeheader()

    def close(self) -> None:
        """Flush and close the target, if one was ever opened."""
        self.handler.close()

    @property
    def io(self) -> IO[Any]:
        return self.open()

    def __enter__(self) -> Self:
        if not self.lazy:
            self.open()
        return self

    def __exit__(self, *args) -> None:
        self.close()

    def write(self, row: SDict) -> None:
        if self.output_format == FORMAT_CSV:
            # opening starts the writer itself when the header is known up front
            self.open()
            if self.csv_writer is None:
                self._start_csv(row.keys())

        if self.output_format == FORMAT_JSON:
            if self.clean:
                row = clean_dict(row)
            line = orjson.dumps(
                row,
                default=_default_serializer,
                option=orjson.OPT_APPEND_NEWLINE | orjson.OPT_NAIVE_UTC,
            )
            if "b" not in self.mode:
                line = line.decode()
            self.io.write(line)
        elif self.csv_writer:
            self.csv_writer.writerow(row)

close()

Flush and close the target, if one was ever opened.

Source code in anystore/io/write.py
def close(self) -> None:
    """Flush and close the target, if one was ever opened."""
    self.handler.close()

open()

The open target, created on first ask.

A csv whose header is known up front writes it here, so a writer that never sees a row still leaves a complete – if empty – table instead of a zero-byte file.

Source code in anystore/io/write.py
def open(self) -> IO[Any]:
    """The open target, created on first ask.

    A csv whose header is known up front writes it here, so a writer that
    never sees a row still leaves a complete – if empty – table instead of
    a zero-byte file.
    """
    if self._io is None:
        self._io = self.handler.open()
        if self.output_format == FORMAT_CSV and self.fieldnames:
            self._start_csv(self.fieldnames)
    return self._io

logged_items(items, action, chunk_size=10000, item_name=None, logger=None, total=None, **log_kwargs)

Log process of iterating items for io operations.

Example
from anystore.io import logged_items

items = [...]
for item in logged_items(items, "Read", uri="/tmp/foo.csv"):
    yield item

Parameters:

Name Type Description Default
items Iterable[T]

Sequence of any items

required
action str

Action name to log

required
chunk_size int | None

Log on every chunk_size

10000
item_name str | None

Name of item

None
logger Logger | BoundLogger | None

Specific logger to use

None

Yields:

Type Description
T

The input items

Source code in anystore/io/logging.py
def logged_items(
    items: Iterable[T],
    action: str,
    chunk_size: int | None = 10_000,
    item_name: str | None = None,
    logger: logging.Logger | BoundLogger | None = None,
    total: int | None = None,
    **log_kwargs,
) -> Generator[T, None, None]:
    """
    Log process of iterating items for io operations.

    Example:
        ```python
        from anystore.io import logged_items

        items = [...]
        for item in logged_items(items, "Read", uri="/tmp/foo.csv"):
            yield item
        ```

    Args:
        items: Sequence of any items
        action: Action name to log
        chunk_size: Log on every chunk_size
        item_name: Name of item
        logger: Specific logger to use

    Yields:
        The input items
    """
    log_ = logger or log
    chunk_size = chunk_size or 10_000
    ix = 0
    item_name = item_name or "Item"
    if total:
        log_.info(f"{action} {total} `{item_name}s` ...", **log_kwargs)
        yield from tqdm(items, total=total, unit=item_name)
        ix = total
    else:
        for ix, item in enumerate(items, 1):
            if ix == 1:
                item_name = item_name or item.__class__.__name__.title()
            if ix % chunk_size == 0:
                item_name = item_name or item.__class__.__name__.title()
                log_.info(f"{action} `{item_name}` {ix} ...", **log_kwargs)
            yield item
    if ix:
        log_.info(f"{action} {ix} `{item_name}s`: Done.", **log_kwargs)

logging_through(console)

Route log output through a rich console for the duration of the block, so log lines don't cut into whatever that console is rendering.

anystore's log handlers resolve sys.stderr on every emit, so swapping it is enough to catch them – as well as warnings and any other library writing to stderr.

Parameters:

Name Type Description Default
console Console

The console to print through

required
Source code in anystore/io/progress.py
@contextmanager
def logging_through(console: Console) -> Iterator[None]:
    """
    Route log output through a rich console for the duration of the block, so
    log lines don't cut into whatever that console is rendering.

    anystore's log handlers resolve `sys.stderr` on every emit, so swapping it
    is enough to catch them – as well as warnings and any other library writing
    to stderr.

    Args:
        console: The console to print through
    """
    stderr = sys.stderr
    sys.stderr = _ConsoleStream(console)  # type: ignore[assignment]
    try:
        yield
    finally:
        sys.stderr = stderr

open_virtual(uri, algorithm=None, **kwargs)

Wrapper for UriResource.local_open

Parameters:

Name Type Description Default
uri 'TUri'

string or path-like key uri to open

required
algorithm str | None

Checksum algorithm from hashlib (default: "sha1")

None
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/read.py
def open_virtual(
    uri: "TUri", algorithm: str | None = None, **kwargs: Any
) -> ContextManager[VirtualIO]:
    """Wrapper for [UriResource.local_open][anystore.store.resource.UriResource.local_open]

    Args:
        uri: string or path-like key uri to open
        algorithm: Checksum algorithm from `hashlib` (default: "sha1")
        **kwargs: pass through storage-specific options
    """
    return UriResource(uri, **kwargs).local_open(algorithm=algorithm)

smart_open(uri, mode=DEFAULT_MODE, compression=None, **kwargs)

IO context similar to pythons built-in open().

Example
from anystore import smart_open

with smart_open("s3://mybucket/foo.csv") as fh:
    return fh.read()

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
mode str | None

open mode, default rb for byte reading.

DEFAULT_MODE
compression CompressKind | str | None

Codec to (de-)compress the stream with ("gz", "bz2", "xz", "zst", "lz4")

None
**kwargs Any

pass through storage-specific options

{}

Yields:

Type Description
IO[Any]

A generic file-handler like context object

Source code in anystore/io/handler.py
@contextlib.contextmanager
def smart_open(
    uri: Uri,
    mode: str | None = DEFAULT_MODE,
    compression: CompressKind | str | None = None,
    **kwargs: Any,
) -> Generator[IO[Any], None, None]:
    """
    IO context similar to pythons built-in `open()`.

    Example:
        ```python
        from anystore import smart_open

        with smart_open("s3://mybucket/foo.csv") as fh:
            return fh.read()
        ```

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        mode: open mode, default `rb` for byte reading.
        compression: Codec to (de-)compress the stream with ("gz", "bz2",
            "xz", "zst", "lz4")
        **kwargs: pass through storage-specific options

    Yields:
        A generic file-handler like context object
    """
    handler = SmartHandler(uri, mode=mode, compression=compression, **kwargs)
    try:
        yield handler.open()
    except FileNotFoundError as e:
        raise DoesNotExist(str(e))
    finally:
        handler.close()

smart_read(uri, mode=DEFAULT_MODE, **kwargs)

Return content for a given file-like key directly.

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
mode str | None

open mode, default rb for byte reading.

DEFAULT_MODE
**kwargs Any

pass through storage-specific options

{}

Returns:

Type Description
AnyStr

str or byte content, depending on mode

Source code in anystore/io/read.py
def smart_read(uri: Uri, mode: str | None = DEFAULT_MODE, **kwargs: Any) -> AnyStr:
    """
    Return content for a given file-like key directly.

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        mode: open mode, default `rb` for byte reading.
        **kwargs: pass through storage-specific options

    Returns:
        `str` or `byte` content, depending on `mode`
    """
    with smart_open(uri, mode, **kwargs) as fh:
        return fh.read()

smart_stream(uri, mode=DEFAULT_MODE, **kwargs)

Stream content line by line.

Example
import orjson
from anystore import smart_stream

while data := smart_stream("s3://mybucket/data.json"):
    yield orjson.loads(data)

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
mode str | None

open mode, default rb for byte reading.

DEFAULT_MODE
**kwargs Any

pass through storage-specific options

{}

Yields:

Type Description
AnyStr

A generator of str or byte content, depending on mode

Source code in anystore/io/read.py
def smart_stream(
    uri: Uri, mode: str | None = DEFAULT_MODE, **kwargs: Any
) -> Generator[AnyStr, None, None]:
    """
    Stream content line by line.

    Example:
        ```python
        import orjson
        from anystore import smart_stream

        while data := smart_stream("s3://mybucket/data.json"):
            yield orjson.loads(data)
        ```

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        mode: open mode, default `rb` for byte reading.
        **kwargs: pass through storage-specific options

    Yields:
        A generator of `str` or `byte` content, depending on `mode`
    """
    with smart_open(uri, mode, **kwargs) as fh:
        yield from iter_lines(fh)

smart_stream_csv(uri, **kwargs)

Stream csv as python objects.

Example
from anystore import smart_stream_csv

for data in smart_stream_csv("s3://mybucket/data.csv"):
    yield data.get("foo")

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
**kwargs Any

pass through storage-specific options

{}

Yields:

Type Description
SDictGenerator

A generator of dicts loaded via csv.DictReader

Source code in anystore/io/read.py
def smart_stream_csv(uri: Uri, **kwargs: Any) -> SDictGenerator:
    """
    Stream csv as python objects.

    Example:
        ```python
        from anystore import smart_stream_csv

        for data in smart_stream_csv("s3://mybucket/data.csv"):
            yield data.get("foo")
        ```

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        **kwargs: pass through storage-specific options

    Yields:
        A generator of `dict`s loaded via `csv.DictReader`
    """
    kwargs["mode"] = "r"
    with smart_open(uri, **kwargs) as f:
        yield from csv.DictReader(f)

smart_stream_csv_models(uri, model, **kwargs)

Stream csv as pydantic objects

Source code in anystore/io/read.py
def smart_stream_csv_models(uri: Uri, model: Type[M], **kwargs: Any) -> MGenerator:
    """
    Stream csv as pydantic objects
    """
    for row in logged_items(
        smart_stream_csv(uri, **kwargs),
        "Read",
        uri=uri,
        item_name=model.__name__,
    ):
        yield model(**row)

smart_stream_data(uri, input_format, **kwargs)

Stream data objects loaded as dict from json or csv sources

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
input_format str

csv or json

required
**kwargs Any

pass through storage-specific options

{}

Yields:

Type Description
SDictGenerator

A generator of dicts loaded via orjson

Source code in anystore/io/read.py
def smart_stream_data(uri: Uri, input_format: str, **kwargs: Any) -> SDictGenerator:
    """
    Stream data objects loaded as dict from json or csv sources

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        input_format: csv or json
        **kwargs: pass through storage-specific options

    Yields:
        A generator of `dict`s loaded via `orjson`
    """
    if input_format == "csv":
        yield from smart_stream_csv(uri, **kwargs)
    else:
        yield from smart_stream_json(uri, **kwargs)

smart_stream_json(uri, mode=DEFAULT_MODE, **kwargs)

Stream line-based json as python objects.

Example
from anystore import smart_stream_json

for data in smart_stream_json("s3://mybucket/data.json"):
    yield data.get("foo")

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
mode str | None

open mode, default rb for byte reading.

DEFAULT_MODE
**kwargs Any

pass through storage-specific options

{}

Yields:

Type Description
SDictGenerator

A generator of dicts loaded via orjson

Source code in anystore/io/read.py
def smart_stream_json(
    uri: Uri, mode: str | None = DEFAULT_MODE, **kwargs: Any
) -> SDictGenerator:
    """
    Stream line-based json as python objects.

    Example:
        ```python
        from anystore import smart_stream_json

        for data in smart_stream_json("s3://mybucket/data.json"):
            yield data.get("foo")
        ```

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        mode: open mode, default `rb` for byte reading.
        **kwargs: pass through storage-specific options

    Yields:
        A generator of `dict`s loaded via `orjson`
    """
    for line in smart_stream(uri, mode, **kwargs):
        yield orjson.loads(line)

smart_stream_json_models(uri, model, **kwargs)

Stream json as pydantic objects

Source code in anystore/io/read.py
def smart_stream_json_models(uri: Uri, model: Type[M], **kwargs: Any) -> MGenerator:
    """
    Stream json as pydantic objects
    """
    for row in logged_items(
        smart_stream_json(uri, **kwargs),
        "Read",
        uri=uri,
        item_name=model.__name__,
    ):
        yield model(**row)

smart_stream_models(uri, model, input_format, **kwargs)

Stream json as pydantic objects

Source code in anystore/io/read.py
def smart_stream_models(
    uri: Uri, model: Type[M], input_format: str, **kwargs: Any
) -> MGenerator:
    """
    Stream json as pydantic objects
    """
    if input_format == "csv":
        yield from smart_stream_csv_models(uri, model, **kwargs)
    elif input_format == "json":
        yield from smart_stream_json_models(uri, model, **kwargs)
    else:
        raise ValueError("Invalid format, only csv or json allowed")

smart_write(uri, content, mode=DEFAULT_WRITE_MODE, **kwargs)

Write content to a given file-like key directly.

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
content bytes | str

str or bytes content to write.

required
mode str | None

open mode, default wb for byte writing.

DEFAULT_WRITE_MODE
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/write.py
def smart_write(
    uri: Uri, content: bytes | str, mode: str | None = DEFAULT_WRITE_MODE, **kwargs: Any
) -> None:
    """
    Write content to a given file-like key directly.

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        content: `str` or `bytes` content to write.
        mode: open mode, default `wb` for byte writing.
        **kwargs: pass through storage-specific options
    """
    if uri == "-":
        if isinstance(content, str):
            content = content.encode()
    with smart_open(uri, mode, **kwargs) as fh:
        fh.write(content)

smart_write_csv(uri, items, mode=DEFAULT_WRITE_MODE, **kwargs)

Write python data to csv

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
items Iterable[SDict]

Iterable of dictionaries

required
mode str | None

open mode, default wb for byte writing.

DEFAULT_WRITE_MODE
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/write.py
def smart_write_csv(
    uri: Uri,
    items: Iterable[SDict],
    mode: str | None = DEFAULT_WRITE_MODE,
    **kwargs: Any,
) -> None:
    """
    Write python data to csv

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        items: Iterable of dictionaries
        mode: open mode, default `wb` for byte writing.
        **kwargs: pass through storage-specific options
    """
    with Writer(uri, mode, output_format="csv", **kwargs) as writer:
        for item in items:
            writer.write(item)

smart_write_data(uri, items, mode=DEFAULT_WRITE_MODE, output_format='json', **kwargs)

Write python data to json or csv

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
items Iterable[SDict]

Iterable of dictionaries

required
mode str | None

open mode, default wb for byte writing.

DEFAULT_WRITE_MODE
output_format Formats | None

csv or json (default: json)

'json'
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/write.py
def smart_write_data(
    uri: Uri,
    items: Iterable[SDict],
    mode: str | None = DEFAULT_WRITE_MODE,
    output_format: Formats | None = "json",
    **kwargs: Any,
) -> None:
    """
    Write python data to json or csv

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        items: Iterable of dictionaries
        mode: open mode, default `wb` for byte writing.
        output_format: csv or json (default: json)
        **kwargs: pass through storage-specific options
    """
    with Writer(uri, mode, output_format=output_format, **kwargs) as writer:
        for item in items:
            writer.write(item)

smart_write_json(uri, items, mode=DEFAULT_WRITE_MODE, **kwargs)

Write python data to json

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
items Iterable[SDict]

Iterable of dictionaries

required
mode str | None

open mode, default wb for byte writing.

DEFAULT_WRITE_MODE
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/write.py
def smart_write_json(
    uri: Uri,
    items: Iterable[SDict],
    mode: str | None = DEFAULT_WRITE_MODE,
    **kwargs: Any,
) -> None:
    """
    Write python data to json

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        items: Iterable of dictionaries
        mode: open mode, default `wb` for byte writing.
        **kwargs: pass through storage-specific options
    """
    with Writer(uri, mode, output_format="json", **kwargs) as writer:
        for item in items:
            writer.write(item)

smart_write_model(uri, obj, mode=DEFAULT_WRITE_MODE, output_format='json', clean=False, **kwargs)

Write a single pydantic object to the target

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
obj BaseModel

Pydantic object

required
mode str | None

open mode, default wb for byte writing.

DEFAULT_WRITE_MODE
clean bool | None

Apply clean_dict

False
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/write.py
def smart_write_model(
    uri: Uri,
    obj: BaseModel,
    mode: str | None = DEFAULT_WRITE_MODE,
    output_format: Formats | None = "json",
    clean: bool | None = False,
    **kwargs: Any,
) -> None:
    """
    Write a single pydantic object to the target

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        obj: Pydantic object
        mode: open mode, default `wb` for byte writing.
        clean: Apply [clean_dict][anystore.util.data.clean_dict]
        **kwargs: pass through storage-specific options
    """
    with ModelWriter(uri, mode, output_format, clean=clean, **kwargs) as writer:
        writer.write(obj)

smart_write_models(uri, objects, mode=DEFAULT_WRITE_MODE, output_format='json', clean=False, **kwargs)

Write pydantic objects to json lines or csv

Parameters:

Name Type Description Default
uri Uri

string or path-like key uri to open, e.g. ./local/data.txt or s3://mybucket/foo

required
objects Iterable[BaseModel]

Iterable of pydantic objects

required
mode str | None

open mode, default wb for byte writing.

DEFAULT_WRITE_MODE
clean bool | None

Apply clean_dict

False
**kwargs Any

pass through storage-specific options

{}
Source code in anystore/io/write.py
def smart_write_models(
    uri: Uri,
    objects: Iterable[BaseModel],
    mode: str | None = DEFAULT_WRITE_MODE,
    output_format: Formats | None = "json",
    clean: bool | None = False,
    **kwargs: Any,
) -> None:
    """
    Write pydantic objects to json lines or csv

    Args:
        uri: string or path-like key uri to open, e.g. `./local/data.txt` or
            `s3://mybucket/foo`
        objects: Iterable of pydantic objects
        mode: open mode, default `wb` for byte writing.
        clean: Apply [clean_dict][anystore.util.data.clean_dict]
        **kwargs: pass through storage-specific options
    """
    with ModelWriter(uri, mode, output_format, clean=clean, **kwargs) as writer:
        for obj in objects:
            writer.write(obj)

stream_bytes(key, source, target, **kwargs)

Stream binary content for key from source to target store. Returns streamed bytes count

Source code in anystore/logic/io.py
def stream_bytes(key: str, source: "Store", target: "Store", **kwargs: Any) -> int:
    """Stream binary content for *key* from *source* to *target* store. Returns
    streamed bytes count"""
    with source.open(key, "rb", **kwargs) as i:
        with target.open(key, "wb") as o:
            return stream(i, o)