Coverage for /pythoncovmergedfiles/medio/medio/usr/local/lib/python3.11/site-packages/fsspec/transaction.py: 26%
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
1import logging
2from collections import deque
4logger = logging.getLogger("fsspec")
7class Transaction:
8 """Filesystem transaction write context
10 Gathers files for deferred commit or discard, so that several write
11 operations can be finalized semi-atomically. This works by having this
12 instance as the ``.transaction`` attribute of the given filesystem
13 """
15 def __init__(self, fs, **kwargs):
16 """
17 Parameters
18 ----------
19 fs: FileSystem instance
20 """
21 self.fs = fs
22 self.files = deque()
24 def __enter__(self):
25 self.start()
26 return self
28 def __exit__(self, exc_type, exc_val, exc_tb):
29 """End transaction and commit, if exit is not due to exception"""
30 # only commit if there was no exception
31 self.complete(commit=exc_type is None)
32 if self.fs:
33 self.fs._intrans = False
34 self.fs._transaction = None
35 self.fs = None
37 def start(self):
38 """Start a transaction on this FileSystem"""
39 self.files = deque() # clean up after previous failed completions
40 self.fs._intrans = True
42 def complete(self, commit=True):
43 """Finish transaction: commit or discard all deferred files"""
44 f = None
45 try:
46 while self.files:
47 f = self.files.popleft()
48 if commit:
49 f.commit()
50 else:
51 f.discard()
52 f = None
53 finally:
54 # the file being processed when the error was raised is already
55 # off the queue; put it back so its temporary file is cleaned up
56 if f is not None:
57 self.files.appendleft(f)
58 # A failed commit or discard must still end the transaction.
59 # Leaving _intrans set would defer every later write on this
60 # filesystem into a temporary file that nothing ever commits,
61 # and instances are cached, so that would persist process-wide.
62 while self.files:
63 try:
64 self.files.popleft().discard()
65 except Exception:
66 logger.debug("Discarding deferred file failed", exc_info=True)
67 self.fs._intrans = False
68 self.fs._transaction = None
69 self.fs = None
72class FileActor:
73 def __init__(self):
74 self.files = []
76 def commit(self):
77 for f in self.files:
78 f.commit()
79 self.files.clear()
81 def discard(self):
82 for f in self.files:
83 f.discard()
84 self.files.clear()
86 def append(self, f):
87 self.files.append(f)
90class DaskTransaction(Transaction):
91 def __init__(self, fs):
92 """
93 Parameters
94 ----------
95 fs: FileSystem instance
96 """
97 import distributed
99 super().__init__(fs)
100 client = distributed.default_client()
101 self.files = client.submit(FileActor, actor=True).result()
103 def complete(self, commit=True):
104 """Finish transaction: commit or discard all deferred files"""
105 try:
106 if commit:
107 self.files.commit().result()
108 else:
109 self.files.discard().result()
110 finally:
111 self.fs._intrans = False
112 self.fs._transaction = None
113 self.fs = None