Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/fsspec/core.py: 16%
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
Shortcuts on this page
r m x toggle line displays
j k next/prev highlighted chunk
0 (zero) top of page
1 (one) first highlighted chunk
1from __future__ import annotations
3import io
4import logging
5import posixpath
6import re
7from glob import has_magic
8from pathlib import Path
10# for backwards compat, we export cache things from here too
11from fsspec.caching import ( # noqa: F401
12 BaseCache,
13 BlockCache,
14 BytesCache,
15 MMapCache,
16 ReadAheadCache,
17 caches,
18)
19from fsspec.compression import compr
20from fsspec.config import conf
21from fsspec.registry import available_protocols, filesystem, get_filesystem_class
22from fsspec.utils import (
23 _unstrip_protocol,
24 build_name_function,
25 infer_compression,
26 stringify_path,
27)
29logger = logging.getLogger("fsspec")
32class OpenFile:
33 """
34 File-like object to be used in a context
36 Can layer (buffered) text-mode and compression over any file-system, which
37 are typically binary-only.
39 These instances are safe to serialize, as the low-level file object
40 is not created until invoked using ``with``.
42 Parameters
43 ----------
44 fs: FileSystem
45 The file system to use for opening the file. Should be a subclass or duck-type
46 with ``fsspec.spec.AbstractFileSystem``
47 path: str
48 Location to open
49 mode: str like 'rb', optional
50 Mode of the opened file
51 compression: str or None, optional
52 Compression to apply
53 encoding: str or None, optional
54 The encoding to use if opened in text mode.
55 errors: str or None, optional
56 How to handle encoding errors if opened in text mode.
57 newline: None or str
58 Passed to TextIOWrapper in text mode, how to handle line endings.
59 autoopen: bool
60 If True, calls open() immediately. Mostly used by pickle
61 pos: int
62 If given and autoopen is True, seek to this location immediately
63 """
65 def __init__(
66 self,
67 fs,
68 path,
69 mode="rb",
70 compression=None,
71 encoding=None,
72 errors=None,
73 newline=None,
74 ):
75 self.fs = fs
76 self.path = path
77 self.mode = mode
78 self.compression = get_compression(path, compression)
79 self.encoding = encoding
80 self.errors = errors
81 self.newline = newline
82 self.fobjects = []
84 def __reduce__(self):
85 return (
86 OpenFile,
87 (
88 self.fs,
89 self.path,
90 self.mode,
91 self.compression,
92 self.encoding,
93 self.errors,
94 self.newline,
95 ),
96 )
98 def __repr__(self):
99 return f"<OpenFile '{self.path}'>"
101 def __enter__(self):
102 mode = self.mode.replace("t", "").replace("b", "") + "b"
104 try:
105 f = self.fs.open(self.path, mode=mode)
106 except FileNotFoundError as e:
107 if has_magic(self.path):
108 raise FileNotFoundError(
109 "%s not found. The URL contains glob characters: you maybe needed\n"
110 "to pass expand=True in fsspec.open() or the storage_options of \n"
111 "your library. You can also set the config value 'open_expand'\n"
112 "before import, or fsspec.core.DEFAULT_EXPAND at runtime, to True.",
113 self.path,
114 ) from e
115 raise
117 self.fobjects = [f]
119 if self.compression is not None:
120 compress = compr[self.compression]
121 f = compress(f, mode=mode[0])
122 self.fobjects.append(f)
124 if "b" not in self.mode:
125 # assume, for example, that 'r' is equivalent to 'rt' as in builtin
126 f = PickleableTextIOWrapper(
127 f, encoding=self.encoding, errors=self.errors, newline=self.newline
128 )
129 self.fobjects.append(f)
131 return self.fobjects[-1]
133 def __exit__(self, *args):
134 self.close()
136 @property
137 def full_name(self):
138 return _unstrip_protocol(self.path, self.fs)
140 def open(self):
141 """Materialise this as a real open file without context
143 The OpenFile object should be explicitly closed to avoid enclosed file
144 instances persisting. You must, therefore, keep a reference to the OpenFile
145 during the life of the file-like it generates.
146 """
147 return self.__enter__()
149 def close(self):
150 """Close all encapsulated file objects"""
151 for f in reversed(self.fobjects):
152 if "r" not in self.mode and not f.closed:
153 f.flush()
154 f.close()
155 self.fobjects.clear()
158class OpenFiles(list):
159 """List of OpenFile instances
161 Can be used in a single context, which opens and closes all of the
162 contained files. Normal list access to get the elements works as
163 normal.
165 A special case is made for caching filesystems - the files will
166 be down/uploaded together at the start or end of the context, and
167 this may happen concurrently, if the target filesystem supports it.
168 """
170 def __init__(self, *args, mode="rb", fs=None):
171 self.mode = mode
172 self.fs = fs
173 self.files = []
174 super().__init__(*args)
176 def __enter__(self):
177 if self.fs is None:
178 raise ValueError("Context has already been used")
180 fs = self.fs
181 while True:
182 if hasattr(fs, "open_many"):
183 # check for concurrent cache download; or set up for upload
184 self.files = fs.open_many(self)
185 return self.files
186 if hasattr(fs, "fs") and fs.fs is not None:
187 fs = fs.fs
188 else:
189 break
190 return [s.__enter__() for s in self]
192 def __exit__(self, *args):
193 fs = self.fs
194 [s.__exit__(*args) for s in self]
195 if "r" in self.mode:
196 # open_many() returns files without populating the OpenFile objects.
197 for f in self.files:
198 f.close()
199 else:
200 while True:
201 if hasattr(fs, "open_many"):
202 # check for concurrent cache upload
203 fs.commit_many(self.files)
204 return
205 if hasattr(fs, "fs") and fs.fs is not None:
206 fs = fs.fs
207 else:
208 break
210 def __getitem__(self, item):
211 out = super().__getitem__(item)
212 if isinstance(item, slice):
213 return OpenFiles(out, mode=self.mode, fs=self.fs)
214 return out
216 def __repr__(self):
217 return f"<List of {len(self)} OpenFile instances>"
220def open_files(
221 urlpath,
222 mode="rb",
223 compression=None,
224 encoding="utf8",
225 errors=None,
226 name_function=None,
227 num=1,
228 protocol=None,
229 newline=None,
230 auto_mkdir=True,
231 expand=True,
232 **kwargs,
233):
234 """Given a path or paths, return a list of ``OpenFile`` objects.
236 For writing, a str path must contain the "*" character, which will be filled
237 in by increasing numbers, e.g., "part*" -> "part1", "part2" if num=2.
239 For either reading or writing, can instead provide explicit list of paths.
241 Parameters
242 ----------
243 urlpath: string or list
244 Absolute or relative filepath(s). Prefix with a protocol like ``s3://``
245 to read from alternative filesystems. To read from multiple files you
246 can pass a globstring or a list of paths, with the caveat that they
247 must all have the same protocol.
248 mode: 'rb', 'wt', etc.
249 compression: string or None
250 If given, open file using compression codec. Can either be a compression
251 name (a key in ``fsspec.compression.compr``) or "infer" to guess the
252 compression from the filename suffix.
253 encoding: str
254 For text mode only
255 errors: None or str
256 Passed to TextIOWrapper in text mode
257 name_function: function or None
258 if opening a set of files for writing, those files do not yet exist,
259 so we need to generate their names by formatting the urlpath for
260 each sequence number
261 num: int [1]
262 if writing mode, number of files we expect to create (passed to
263 name+function)
264 protocol: str or None
265 If given, overrides the protocol found in the URL.
266 newline: bytes or None
267 Used for line terminator in text mode. If None, uses system default;
268 if blank, uses no translation.
269 auto_mkdir: bool (True)
270 If in write mode, this will ensure the target directory exists before
271 writing, by calling ``fs.mkdirs(exist_ok=True)``.
272 expand: bool
273 **kwargs: dict
274 Extra options that make sense to a particular storage connection, e.g.
275 host, port, username, password, etc.
277 Examples
278 --------
279 >>> files = open_files('2015-*-*.csv') # doctest: +SKIP
280 >>> files = open_files(
281 ... 's3://bucket/2015-*-*.csv.gz', compression='gzip'
282 ... ) # doctest: +SKIP
284 Returns
285 -------
286 An ``OpenFiles`` instance, which is a list of ``OpenFile`` objects that can
287 be used as a single context
289 Notes
290 -----
291 For a full list of the available protocols and the implementations that
292 they map across to see the latest online documentation:
294 - For implementations built into ``fsspec`` see
295 https://filesystem-spec.readthedocs.io/en/latest/api.html#built-in-implementations
296 - For implementations in separate packages see
297 https://filesystem-spec.readthedocs.io/en/latest/api.html#other-known-implementations
298 """
299 fs, fs_token, paths = get_fs_token_paths(
300 urlpath,
301 mode,
302 num=num,
303 name_function=name_function,
304 storage_options=kwargs,
305 protocol=protocol,
306 expand=expand,
307 )
308 if fs.protocol == "file":
309 fs.auto_mkdir = auto_mkdir
310 elif "r" not in mode and auto_mkdir:
311 parents = {fs._parent(path) for path in paths}
312 for parent in parents:
313 try:
314 fs.makedirs(parent, exist_ok=True)
315 except PermissionError:
316 pass
317 return OpenFiles(
318 [
319 OpenFile(
320 fs,
321 path,
322 mode=mode,
323 compression=compression,
324 encoding=encoding,
325 errors=errors,
326 newline=newline,
327 )
328 for path in paths
329 ],
330 mode=mode,
331 fs=fs,
332 )
335def _un_chain(path, kwargs):
336 # Avoid a circular import
337 from fsspec.implementations.chained import ChainedFileSystem
339 if "::" in path:
340 x = re.compile(".*[^a-z]+.*") # test for non protocol-like single word
341 known_protocols = set(available_protocols())
342 bits = []
344 # split on '::', then ensure each bit has a protocol
345 for p in path.split("::"):
346 if p in known_protocols:
347 bits.append(p + "://")
348 elif "://" in p or x.match(p):
349 bits.append(p)
350 else:
351 bits.append(p + "://")
352 else:
353 bits = [path]
355 # [[url, protocol, kwargs], ...]
356 out = []
357 previous_bit = None
358 kwargs = kwargs.copy()
360 for bit in reversed(bits):
361 protocol = kwargs.pop("protocol", None) or split_protocol(bit)[0] or "file"
362 cls = get_filesystem_class(protocol)
363 extra_kwargs = cls._get_kwargs_from_urls(bit)
364 kws = kwargs.pop(protocol, {})
366 if bit is bits[0]:
367 kws.update(kwargs)
369 kw = dict(
370 **{k: v for k, v in extra_kwargs.items() if k not in kws or v != kws[k]},
371 **kws,
372 )
373 bit = cls._strip_protocol(bit)
375 if (
376 "target_protocol" not in kw
377 and issubclass(cls, ChainedFileSystem)
378 and not bit
379 ):
380 # replace bit if we are chaining and no path given
381 bit = previous_bit
383 out.append((bit, protocol, kw))
384 previous_bit = bit
386 out.reverse()
387 return out
390def url_to_fs(url, **kwargs):
391 """
392 Turn fully-qualified and potentially chained URL into filesystem instance
394 Parameters
395 ----------
396 url : str
397 The fsspec-compatible URL
398 **kwargs: dict
399 Extra options that make sense to a particular storage connection, e.g.
400 host, port, username, password, etc.
402 Returns
403 -------
404 filesystem : FileSystem
405 The new filesystem discovered from ``url`` and created with
406 ``**kwargs``.
407 urlpath : str
408 The file-systems-specific URL for ``url``.
409 """
410 url = stringify_path(url)
411 # non-FS arguments that appear in fsspec.open()
412 # inspect could keep this in sync with open()'s signature
413 known_kwargs = {
414 "compression",
415 "encoding",
416 "errors",
417 "expand",
418 "mode",
419 "name_function",
420 "newline",
421 "num",
422 }
423 kwargs = {k: v for k, v in kwargs.items() if k not in known_kwargs}
424 chain = _un_chain(url, kwargs)
425 inkwargs = {}
426 # Reverse iterate the chain, creating a nested target_* structure
427 for i, ch in enumerate(reversed(chain)):
428 urls, protocol, kw = ch
429 if i == len(chain) - 1:
430 inkwargs = dict(**kw, **inkwargs)
431 continue
432 inkwargs["target_options"] = dict(**kw, **inkwargs)
433 inkwargs["target_protocol"] = protocol
434 inkwargs["fo"] = urls
435 urlpath, protocol, _ = chain[0]
436 fs = filesystem(protocol, **inkwargs)
437 return fs, urlpath
440DEFAULT_EXPAND = conf.get("open_expand", False)
443def open(
444 urlpath,
445 mode="rb",
446 compression=None,
447 encoding="utf8",
448 errors=None,
449 protocol=None,
450 newline=None,
451 expand=None,
452 **kwargs,
453):
454 """Given a path or paths, return one ``OpenFile`` object.
456 Parameters
457 ----------
458 urlpath: string or list
459 Absolute or relative filepath. Prefix with a protocol like ``s3://``
460 to read from alternative filesystems. Should not include glob
461 character(s).
462 mode: 'rb', 'wt', etc.
463 compression: string or None
464 If given, open file using compression codec. Can either be a compression
465 name (a key in ``fsspec.compression.compr``) or "infer" to guess the
466 compression from the filename suffix.
467 encoding: str
468 For text mode only
469 errors: None or str
470 Passed to TextIOWrapper in text mode
471 protocol: str or None
472 If given, overrides the protocol found in the URL.
473 newline: bytes or None
474 Used for line terminator in text mode. If None, uses system default;
475 if blank, uses no translation.
476 expand: bool or None
477 Whether to regard file paths containing special glob characters as needing
478 expansion (finding the first match) or absolute. Setting False allows using
479 paths which do embed such characters. If None (default), this argument
480 takes its value from the DEFAULT_EXPAND module variable, which takes
481 its initial value from the "open_expand" config value at startup, which will
482 be False if not set.
483 **kwargs: dict
484 Extra options that make sense to a particular storage connection, e.g.
485 host, port, username, password, etc.
487 Examples
488 --------
489 >>> openfile = open('2015-01-01.csv') # doctest: +SKIP
490 >>> openfile = open(
491 ... 's3://bucket/2015-01-01.csv.gz', compression='gzip'
492 ... ) # doctest: +SKIP
493 >>> with openfile as f:
494 ... df = pd.read_csv(f) # doctest: +SKIP
495 ...
497 Returns
498 -------
499 ``OpenFile`` object.
501 Notes
502 -----
503 For a full list of the available protocols and the implementations that
504 they map across to see the latest online documentation:
506 - For implementations built into ``fsspec`` see
507 https://filesystem-spec.readthedocs.io/en/latest/api.html#built-in-implementations
508 - For implementations in separate packages see
509 https://filesystem-spec.readthedocs.io/en/latest/api.html#other-known-implementations
510 """
511 expand = DEFAULT_EXPAND if expand is None else expand
512 out = open_files(
513 urlpath=[urlpath],
514 mode=mode,
515 compression=compression,
516 encoding=encoding,
517 errors=errors,
518 protocol=protocol,
519 newline=newline,
520 expand=expand,
521 **kwargs,
522 )
523 if not out:
524 raise FileNotFoundError(urlpath)
525 return out[0]
528def open_local(
529 url: str | list[str] | Path | list[Path],
530 mode: str = "rb",
531 **storage_options: dict,
532) -> str | list[str]:
533 """Open file(s) which can be resolved to local
535 For files which either are local, or get downloaded upon open
536 (e.g., by file caching)
538 Parameters
539 ----------
540 url: str or list(str)
541 mode: str
542 Must be read mode
543 storage_options:
544 passed on to FS for or used by open_files (e.g., compression)
545 """
546 if "r" not in mode:
547 raise ValueError("Can only ensure local files when reading")
548 of = open_files(url, mode=mode, **storage_options)
549 if not getattr(of[0].fs, "local_file", False):
550 raise ValueError(
551 "open_local can only be used on a filesystem which"
552 " has attribute local_file=True"
553 )
554 with of as files:
555 paths = [f.name for f in files]
556 if (isinstance(url, str) and not has_magic(url)) or isinstance(url, Path):
557 return paths[0]
558 return paths
561def get_compression(urlpath, compression):
562 if compression == "infer":
563 compression = infer_compression(urlpath)
564 if compression is not None and compression not in compr:
565 raise ValueError(f"Compression type {compression} not supported")
566 return compression
569def split_protocol(urlpath):
570 """Return protocol, path pair"""
571 urlpath = stringify_path(urlpath)
572 if "://" in urlpath:
573 protocol, path = urlpath.split("://", 1)
574 if len(protocol) > 1:
575 # excludes Windows paths
576 return protocol, path
577 if urlpath.startswith("data:"):
578 return urlpath.split(":", 1)
579 return None, urlpath
582def strip_protocol(urlpath):
583 """Return only path part of full URL, according to appropriate backend"""
584 protocol, _ = split_protocol(urlpath)
585 cls = get_filesystem_class(protocol)
586 return cls._strip_protocol(urlpath)
589def expand_paths_if_needed(paths, mode, num, fs, name_function):
590 """Expand paths if they have a ``*`` in them (write mode) or any of ``*?[]``
591 in them (read mode).
593 :param paths: list of paths
594 mode: str
595 Mode in which to open files.
596 num: int
597 If opening in writing mode, number of files we expect to create.
598 fs: filesystem object
599 name_function: callable
600 If opening in writing mode, this callable is used to generate path
601 names. Names are generated for each partition by
602 ``urlpath.replace('*', name_function(partition_index))``.
603 :return: list of paths
604 """
605 expanded_paths = []
606 paths = list(paths)
608 if "w" in mode or "x" in mode: # write mode
609 if sum(1 for p in paths if "*" in p) > 1:
610 raise ValueError(
611 "When writing data, only one filename mask can be specified."
612 )
613 num = max(num, len(paths))
615 for curr_path in paths:
616 if "*" in curr_path:
617 # expand using name_function
618 expanded_paths.extend(_expand_paths(curr_path, name_function, num))
619 else:
620 expanded_paths.append(curr_path)
621 # if we generated more paths that asked for, trim the list
622 if len(expanded_paths) > num:
623 expanded_paths = expanded_paths[:num]
625 else: # read mode
626 for curr_path in paths:
627 if has_magic(curr_path):
628 # expand using glob
629 expanded_paths.extend(fs.glob(curr_path))
630 else:
631 expanded_paths.append(curr_path)
633 return expanded_paths
636def get_fs_token_paths(
637 urlpath,
638 mode="rb",
639 num=1,
640 name_function=None,
641 storage_options=None,
642 protocol=None,
643 expand=True,
644):
645 """Filesystem, deterministic token, and paths from a urlpath and options.
647 Parameters
648 ----------
649 urlpath: string or iterable
650 Absolute or relative filepath, URL (may include protocols like
651 ``s3://``), or globstring pointing to data.
652 mode: str, optional
653 Mode in which to open files.
654 num: int, optional
655 If opening in writing mode, number of files we expect to create.
656 name_function: callable, optional
657 If opening in writing mode, this callable is used to generate path
658 names. Names are generated for each partition by
659 ``urlpath.replace('*', name_function(partition_index))``.
660 storage_options: dict, optional
661 Additional keywords to pass to the filesystem class.
662 protocol: str or None
663 To override the protocol specifier in the URL
664 expand: bool
665 Expand string paths for writing, assuming the path is a directory
666 """
667 if isinstance(urlpath, (list, tuple, set)):
668 if not urlpath:
669 raise ValueError("empty urlpath sequence")
670 urlpath0 = stringify_path(next(iter(urlpath)))
671 else:
672 urlpath0 = stringify_path(urlpath)
673 storage_options = storage_options or {}
674 if protocol:
675 storage_options["protocol"] = protocol
676 chain = _un_chain(urlpath0, storage_options or {})
677 inkwargs = {}
678 # Reverse iterate the chain, creating a nested target_* structure
679 for i, ch in enumerate(reversed(chain)):
680 urls, nested_protocol, kw = ch
681 if i == len(chain) - 1:
682 inkwargs = dict(**kw, **inkwargs)
683 continue
684 inkwargs["target_options"] = dict(**kw, **inkwargs)
685 inkwargs["target_protocol"] = nested_protocol
686 inkwargs["fo"] = urls
687 paths, protocol, _ = chain[0]
688 fs = filesystem(protocol, **inkwargs)
689 if isinstance(urlpath, (list, tuple, set)):
690 pchains = [
691 _un_chain(stringify_path(u), storage_options or {})[0] for u in urlpath
692 ]
693 if len({pc[1] for pc in pchains}) > 1:
694 raise ValueError("Protocol mismatch getting fs from %s", urlpath)
695 paths = [pc[0] for pc in pchains]
696 else:
697 paths = fs._strip_protocol(paths)
698 if isinstance(paths, (list, tuple, set)):
699 if expand:
700 paths = expand_paths_if_needed(paths, mode, num, fs, name_function)
701 elif not isinstance(paths, list):
702 paths = list(paths)
703 else:
704 if ("w" in mode or "x" in mode) and expand:
705 paths = _expand_paths(paths, name_function, num)
706 elif "*" in paths:
707 paths = [f for f in sorted(fs.glob(paths)) if not fs.isdir(f)]
708 else:
709 paths = [paths]
711 return fs, fs._fs_token, paths
714def _expand_paths(path, name_function, num):
715 if isinstance(path, str):
716 if path.count("*") > 1:
717 raise ValueError("Output path spec must contain exactly one '*'.")
718 elif "*" not in path:
719 # paths are "/"-separated on every filesystem, including local ones
720 path = posixpath.join(path, "*.part")
722 if name_function is None:
723 name_function = build_name_function(num - 1)
725 paths = [path.replace("*", name_function(i)) for i in range(num)]
726 if paths != sorted(paths):
727 logger.warning(
728 "In order to preserve order between partitions"
729 " paths created with ``name_function`` should "
730 "sort to partition order"
731 )
732 elif isinstance(path, (tuple, list)):
733 assert len(path) == num
734 paths = list(path)
735 else:
736 raise ValueError(
737 "Path should be either\n"
738 "1. A list of paths: ['foo.json', 'bar.json', ...]\n"
739 "2. A directory: 'foo/\n"
740 "3. A path with a '*' in it: 'foo.*.json'"
741 )
742 return paths
745class PickleableTextIOWrapper(io.TextIOWrapper):
746 """TextIOWrapper cannot be pickled. This solves it.
748 Requires that ``buffer`` be pickleable, which all instances of
749 AbstractBufferedFile are.
750 """
752 def __init__(
753 self,
754 buffer,
755 encoding=None,
756 errors=None,
757 newline=None,
758 line_buffering=False,
759 write_through=False,
760 ):
761 self.args = buffer, encoding, errors, newline, line_buffering, write_through
762 super().__init__(*self.args)
764 def __reduce__(self):
765 return PickleableTextIOWrapper, self.args