You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
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:
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.
Use case
pyathena/filesystem/implements S3 operations inside the shape of fsspec'sAbstractFileSystem/AbstractBufferedFileAPI, inherited from the s3fs design it started from. fsspec's conventions (string names, glob semantics, theDirCache, 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'scopy()/get()callglob.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.pyunless noted):PATTERN_PATHand fsspec's_strip_protocol(), 10 calls). Therstrip("/")in_strip_protocol()is whydir/anddirmeet ininfo()/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 theS3Pathmodel; the names themselves are still 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 aDirCache.pop()race (:2737-2751), andinvalidate_cache()walks parents with_parent()withoutsuper()(:2692-2735).ls()returns names or objects, and a file path falls back to HeadObject (:599-639).find()/_find()/_find_levels()/_extract_parent_directories()(:832-1009) follow fsspec'sfind()semantics: synthesized directories, the path itself included "as in fsspec", fsspec's maxdepth numbering, and dict vs list returns.glob/expand_path/walk/duare inherited and treat?as a wildcard._expand_delete_paths()routes versioned paths aroundexpand_path()(:1102-1136)._move_paths()repeats fsspeccopy()'s pairing line by line (expand_path,trailing_sep/isdir,has_magic,other_paths; :1505-1585) before adding S3 rules (nullversions, overlap checks).cp_file()takes unused fsspec arguments (:1606-1635), and_copy_file()skips the directories that recursivecopy()passes (:1669-1673). fsspec'sonerrortypo is popped in both sync and async (:1652-1658,s3_async.py:425-427).S3File.__init__()validates beforesuper().__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'sflush()contract (:3354-3551)._close_without_commit(),_write_and_close()and_write_file_and_close()exist because awithblock on an fsspec file commits on error._fetch_range()clamps the ranges fsspec caches request (:3608-3650).pipe_file()reroutes throughopen()in a transaction (:1922-1929).AioS3FileSystemkeeps its own transaction copies ofpipe_file/put_file(s3_async.py:168-280), wraps the sync filesystem withasyncio.to_thread, and hand-writes the sync wrappers that fsspec does not mirror (mv,rmdir,touch).S3AioExecutorexists so that the sync buffered file can run parts on the loop.S3Objectis aMutableMappingso that fsspec can readinfo()["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 ins3_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:
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 takeS3Pathand 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 byS3Pathand(bucket, prefix, delimiter), with its own invalidation and expiry.S3FileSystem,AioS3FileSystemandS3Filekeep their public API (see below). They translate fsspec strings intoS3Path, own the fsspec semantics (glob andexpand_path,find's directory synthesis and depth numbering, copy/move/get destination pairing,info()dicts,detail=True/False, theDirCacheinterface, 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, ...);docs/filesystem.mdand the API reference, includingls(versions=True),object_version_info(), multipart listing, tags/xattrs/ACLs, transactions and the async_xmethods;parse_path()/PATTERN_PATH, and the mapping and attribute access of theS3Objects thatinfo()/ls(detail=True)return;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:
S3Path(Model S3 paths with one value class instead of parsing strings at each call site #1049).S3Objects and names from them.S3Path. The adapter alone does fsspec's pairing, one shared implementation forcopy(),get()andmv()in sync and async (this would also absorb Keys containing '?' are rejected, and pinned or listed versions are not addressable #979's pairing fix).S3Filedelegates to it.DirCachefor fsspec compatibility.Open questions:
pyathena.filesystem.core), or stay private until it settles?DirCache(which fsspec users configure withuse_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 indocs/filesystem.md.