Skip to content

Separate an S3 core from the fsspec adapter in pyathena.filesystem #1053

Description

@laughingman7743

Use case

pyathena/filesystem/ implements S3 operations inside the shape of fsspec's AbstractFileSystem/AbstractBufferedFile API, inherited from the s3fs design it started from. fsspec's conventions (string names, glob semantics, the DirCache, the buffered-file lifecycle) leak into the S3 logic, and each S3 concept that does not fit them becomes another override or workaround. #979 is a current example: a version-qualified path is globbed because fsspec's copy()/get() call glob.has_magic() on the source string and name the destinations after the source strings, so supporting versions means overriding fsspec's pairing.

Inventory of the fsspec-shaped code on master 28ec68d (s3.py unless noted):

  • Paths and names. Two scheme parsers (PATTERN_PATH and fsspec's _strip_protocol(), 10 calls). The rstrip("/") in _strip_protocol() is why dir/ and dir meet in info()/open(). Cache keys are built with fsspec's _parent() (:466, :828, :2735, :2880). Directory entries strip trailing slashes to match fsspec names (:577, :653, :797, :829). Model S3 paths with one value class instead of parsing strings at each call site #1049 adds the S3Path model; the names themselves are still fsspec's.
  • Listing and cache. The whole cache is fsspec's DirCache, holding layouts fsspec never uses: single objects under a path, listings under (path, delimiter), buckets under "", and per-parameter lookups under (path, "lookups") with their own expiry (:65-67, :556-596, :2796-2824). _ls_from_cache() reimplements fsspec's (:2849-2899), _evict_cache() works around a DirCache.pop() race (:2737-2751), and invalidate_cache() walks parents with _parent() without super() (:2692-2735). ls() returns names or objects, and a file path falls back to HeadObject (:599-639).
  • Globbing and expansion. find()/_find()/_find_levels()/_extract_parent_directories() (:832-1009) follow fsspec's find() semantics: synthesized directories, the path itself included "as in fsspec", fsspec's maxdepth numbering, and dict vs list returns. glob/expand_path/walk/du are inherited and treat ? as a wildcard. _expand_delete_paths() routes versioned paths around expand_path() (:1102-1136).
  • Copy, move and delete pairing. _move_paths() repeats fsspec copy()'s pairing line by line (expand_path, trailing_sep/isdir, has_magic, other_paths; :1505-1585) before adding S3 rules (null versions, overlap checks). cp_file() takes unused fsspec arguments (:1606-1635), and _copy_file() skips the directories that recursive copy() passes (:1669-1673). fsspec's onerror typo is popped in both sync and async (:1652-1658, s3_async.py:425-427).
  • File objects. S3File.__init__() validates before super().__init__() because fsspec's __del__ would commit a half-built file. It also sets the details and size so that the base class does not look up the latest version (:3140-3310). Multipart logic lives in _initiate_upload()/_upload_chunk()/commit()/discard(), shaped by fsspec's flush() contract (:3354-3551). _close_without_commit(), _write_and_close() and _write_file_and_close() exist because a with block on an fsspec file commits on error. _fetch_range() clamps the ranges fsspec caches request (:3608-3650).
  • Transactions and the async bridge. pipe_file() reroutes through open() in a transaction (:1922-1929). AioS3FileSystem keeps its own transaction copies of pipe_file/put_file (s3_async.py:168-280), wraps the sync filesystem with asyncio.to_thread, and hand-writes the sync wrappers that fsspec does not mirror (mv, rmdir, touch). S3AioExecutor exists so that the sync buffered file can run parts on the loop.
  • Results. S3Object is a MutableMapping so that fsspec can read info()["size"]/["name"]/["type"]. _directory_object() and the root and bucket objects are synthesized fsspec entries (:322-337, :391-403, :502-515, :748-760).

Parts that are already S3-native and could form a core: _get_object, _put_object, the multipart helpers (:2937-3089), _finish_multipart_upload, the DeleteObjects batching (:1168-1298), _get_copy_ranges, _get_operation_kwargs/_call, and the typed response classes in s3_object.py.

Proposed change

Restructure the filesystem as a typed S3 core wrapped by a thin fsspec adapter, so that fsspec's conventions stay at the edge:

  • Core (no fsspec imports): S3Path (Model S3 paths with one value class instead of parsing strings at each call site #1049); typed entries for objects, prefixes, versions and buckets (instead of fsspec-shaped dicts); operations that take S3Path and return typed results, such as head, list (by prefix and delimiter, by level, all versions), get, put, multipart upload/copy with abort on failure, copy, delete in batches, tags, ACL and metadata; and a cache keyed by S3Path and (bucket, prefix, delimiter), with its own invalidation and expiry.
  • Adapter: S3FileSystem, AioS3FileSystem and S3File keep their public API (see below). They translate fsspec strings into S3Path, own the fsspec semantics (glob and expand_path, find's directory synthesis and depth numbering, copy/move/get destination pairing, info() dicts, detail=True/False, the DirCache interface, transactions, callbacks), and drive the core's multipart writer from the buffered-file lifecycle.

The public surface the adapter must keep:

  • register_s3_filesystem(), and the constructors with their s3fs-compatible arguments (version_aware, s3_additional_kwargs, default_block_size, ...);
  • every method documented in docs/filesystem.md and the API reference, including ls(versions=True), object_version_info(), multipart listing, tags/xattrs/ACLs, transactions and the async _x methods;
  • parse_path()/PATTERN_PATH, and the mapping and attribute access of the S3Objects that info()/ls(detail=True) return;
  • the internal users in pyathena/s3fs/, pyathena/pandas/, pyathena/aio/s3fs/ and Polars through the fsspec registration.

A possible order, one reviewable PR per step, each without a behavior change unless agreed in its own issue:

  1. S3Path (Model S3 paths with one value class instead of parsing strings at each call site #1049).
  2. Typed listing entries and a core listing/head API; the adapter builds S3Objects and names from them.
  3. A core copy/move/delete API that takes pairs of S3Path. The adapter alone does fsspec's pairing, one shared implementation for copy(), get() and mv() in sync and async (this would also absorb Keys containing '?' are rejected, and pinned or listed versions are not addressable #979's pairing fix).
  4. A core multipart writer; S3File delegates to it.
  5. The cache behind a core interface, still stored in DirCache for fsspec compatibility.

Open questions:

  • Does the core become public API (for example pyathena.filesystem.core), or stay private until it settles?
  • Should the async adapter call an async core, or keep running the sync core in threads as it does today?
  • Should the cache keep living in fsspec's DirCache (which fsspec users configure with use_listings_cache/listings_expiry_time), or move into the core with those options mapped?

Out of scope: new S3 features and behavior changes, which go through their own issues.

Validation plan (if implementing)

Each step keeps the existing offline and S3 integration tests in tests/pyathena/filesystem/ passing unchanged, and adds unit tests for the core API it introduces. Steps that move behavior between layers are checked against the resolved fsspec version and the documented behavior in docs/filesystem.md.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions