python-trio / python-trio/trio

Providing a way to wrap blocking iterators

Open
#1,344 3 comments 1 reaction 0 assignees View on GitHub

Nobody has claimed this yet.

threads
Dominant language
Python
Stars
7.3k
Forks
431
Avg merge
2d 17h
Merged PRs (30d)
6

Description

It is currently non-trivial to integrate blocking iterators into trio, see the issues #501 and #1308 about trio.Path.iterdir for instance.

The main problem is that there are many kind of iterators:

  • fast vs slow
  • small/finite vs infinite/large
  • ordered vs sporadic

Delegating to a thread is not free: in particular, there are two costy operations:

  • spawning a thread
  • switching contexts (i.e giving the control back to the trio event loop)

On my machine for instance, the full round trip is 200 us. This means that running trio.to_thread.run_sync for producing each item in range(5000) would add an overhead of about one second.

I've wrote a small benchmark comparing several approaches and I came up with one that (I think) comes close to ticking all the boxes:

import trio
import time
import outcome

async def to_thread_iter_sync(fn, *args, cancellable=False, limiter=None):
    """Convert a blocking iteration into an async iteration using a thread.

    In order to attenuate the overhead of spawning threads and switching
    contexts, values from the blocking iteration are batched for a time one
    order of magnitude greater than the spawn time of a thread.
    """

    def run_batch(items_iter, start_time):
        now = time.monotonic()
        spawn_time = now - start_time
        deadline = now + 10 * spawn_time

        if items_iter is None:
            items_iter = iter(fn(*args))

        batch = []
        while True:

            try:
                item = next(items_iter)
            except Exception as exc:
                batch.append(outcome.Error(exc))
                break
            else:
                batch.append(outcome.Value(item))

            if time.monotonic() > deadline:
                break

        return items_iter, batch

    items_iter = None
    while True:

        items_iter, batch = await trio.to_thread.run_sync(
            run_batch,
            items_iter,
            time.monotonic(),
            cancellable=cancellable,
            limiter=limiter
        )

        for result in batch:
            try:
                yield result.unwrap()
            except StopIteration:
                return

I'd like to submit a PR to add this function as trio.to_thread.iter_sync, if you think the idea is worth considering.

Contributor guide

Open the contributing guide

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start with the proposed implementation in this issue, the existing trio.to_thread.run_sync entry point, and the benchmark linked from issue #1308. Review how batching, cancellation, limiters, and iterator termination should work for fast, slow, finite, infinite, ordered, and sporadic iterators. Done means the design is accepted and the resulting trio.to_thread.iter_sync API is implemented and validated.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
backend
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.