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

65 statements  

1import logging 

2from collections import deque 

3 

4logger = logging.getLogger("fsspec") 

5 

6 

7class Transaction: 

8 """Filesystem transaction write context 

9 

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 """ 

14 

15 def __init__(self, fs, **kwargs): 

16 """ 

17 Parameters 

18 ---------- 

19 fs: FileSystem instance 

20 """ 

21 self.fs = fs 

22 self.files = deque() 

23 

24 def __enter__(self): 

25 self.start() 

26 return self 

27 

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 

36 

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 

41 

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 

70 

71 

72class FileActor: 

73 def __init__(self): 

74 self.files = [] 

75 

76 def commit(self): 

77 for f in self.files: 

78 f.commit() 

79 self.files.clear() 

80 

81 def discard(self): 

82 for f in self.files: 

83 f.discard() 

84 self.files.clear() 

85 

86 def append(self, f): 

87 self.files.append(f) 

88 

89 

90class DaskTransaction(Transaction): 

91 def __init__(self, fs): 

92 """ 

93 Parameters 

94 ---------- 

95 fs: FileSystem instance 

96 """ 

97 import distributed 

98 

99 super().__init__(fs) 

100 client = distributed.default_client() 

101 self.files = client.submit(FileActor, actor=True).result() 

102 

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