Skip to content

Import a local directory into a SharePoint library through a bounded async pipeline.

A producer walks --source-dir and feeds one (remote_name, local_path) per file onto a bounded :class:asyncio.Queue; a pool of worker tasks drains it, reading each file on the executor and uploading it. The bounded queue is the point: instead of firing one asyncio.gather over every file (which would read them all into memory and open an unbounded number of requests), workers pull only what they can process, so disk reads, HTTP uploads and the event loop overlap while memory and connection usage stay flat.

Each worker uploads through its own ctx.clone(site_url) because a context owns a single pending-query queue; clones share credentials and the connection pool, so the whole run still uses one token and one set of sockets. Subfolders are flattened into the file name (a/b.txt -> a__b.txt) so no worker has to race to create the same remote folder.

Permissions

Requires delegated permission Sites.FullControl.All (writes files).

Reference

View source

from __future__ import annotations

import argparse
import asyncio
import sys
import time
from pathlib import Path

from office365.sharepoint.client_context import ClientContext
from office365.sharepoint.lists.templates.type import ListTemplateType
from tests.settings import cert_path, cert_thumbprint, client_id, site_url, tenant


async def read_bytes(path: Path) -> bytes:
    """Read a file on the default executor so the I/O never blocks the loop."""
    loop = asyncio.get_running_loop()
    return await loop.run_in_executor(None, path.read_bytes)


def resolve_source(raw: str) -> Path:
    """Validate and return the source directory (filesystem work off the async path)."""
    source = Path(raw)
    if not source.is_dir():
        raise NotADirectoryError(raw)
    return source


async def worker(
    name: str,
    queue: "asyncio.Queue[tuple[str, Path] | None]",
    ctx: ClientContext,
    list_title: str,
    stats: dict[str, int],
    failures: list[tuple[str, str]],
) -> None:
    """Drain one worker slot: read a file and upload it to the library root."""
    while True:
        item = await queue.get()
        try:
            if item is None:
                return
            remote_name, local_path = item
            content = await read_bytes(local_path)
            clone = ctx.clone(site_url)
            folder = clone.web.lists.get_by_title(list_title).root_folder
            folder.upload_file(remote_name, content)
            await clone.execute_query_async()
            stats["files"] += 1
            stats["bytes"] += len(content)
            print(f"  [{name}] {remote_name} ({len(content):,} bytes)", flush=True)
        except Exception as ex:  # one bad file must not stop the pipeline
            failures.append((str(item[1]) if item else "?", str(ex)))
        finally:
            queue.task_done()


async def produce(queue: "asyncio.Queue[tuple[str, Path] | None]", source: Path) -> None:
    """Walk the source tree and offer each file, blocking when the queue is full."""
    loop = asyncio.get_running_loop()
    files = await loop.run_in_executor(None, lambda: sorted(p for p in source.rglob("*") if p.is_file()))
    for path in files:
        remote_name = path.relative_to(source).as_posix().replace("/", "__")
        await queue.put((remote_name, path))


async def main() -> None:
    parser = argparse.ArgumentParser(description="Import a directory into SharePoint via an async pipeline")
    parser.add_argument("--source-dir", required=True, help="local directory to upload")
    parser.add_argument("--list-title", default="Documents", help="target document library")
    parser.add_argument("--workers", type=int, default=4, help="upload workers (default: 4)")
    parser.add_argument("--queue-size", type=int, default=8, help="bounded queue depth (default: 8)")
    args = parser.parse_args()

    try:
        source = resolve_source(args.source_dir)
    except NotADirectoryError:
        parser.error(f"--source-dir is not a directory: {args.source_dir}")

    worker_count = max(1, args.workers)
    queue: asyncio.Queue[tuple[str, Path] | None] = asyncio.Queue(maxsize=max(worker_count, args.queue_size))
    stats = {"files": 0, "bytes": 0}
    failures: list[tuple[str, str]] = []

    ctx = ClientContext(site_url).with_client_certificate(
        tenant, client_id=client_id, thumbprint=cert_thumbprint, cert_path=cert_path
    )
    await ctx.web.ensure_list(args.list_title, template_type=ListTemplateType.DocumentLibrary).execute_query_async()

    started = time.perf_counter()
    workers = [
        asyncio.create_task(worker(f"w{i + 1}", queue, ctx, args.list_title, stats, failures))
        for i in range(worker_count)
    ]
    await produce(queue, source)
    for _ in workers:
        await queue.put(None)
    await asyncio.gather(*workers)
    elapsed = time.perf_counter() - started

    mb = stats["bytes"] / (1024 * 1024)
    throughput = mb / elapsed if elapsed else 0.0
    print(f"\nUploaded {stats['files']} files ({mb:.1f} MB) in {elapsed:.1f}s ({throughput:.1f} MB/s)")
    for path, error in failures:
        print(f"  failed: {path}: {error}")
    if failures:
        sys.exit(1)


if __name__ == "__main__":
    asyncio.run(main())

← Back to Async / await