-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathstate_manager.py
More file actions
67 lines (56 loc) · 1.8 KB
/
Copy pathstate_manager.py
File metadata and controls
67 lines (56 loc) · 1.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
"""
Persist migration progress to a JSON file so the run can be resumed.
"""
import json
import logging
import pathlib
log = logging.getLogger(__name__)
_DEFAULTS = {
"last_completed_offset": 0,
"processed_count": 0,
"skipped_count": 0,
"error_count": 0,
"completed": False,
}
class StateManager:
def __init__(self, path):
self.path = pathlib.Path(path)
self.state = self._load()
def _load(self):
if self.path.exists():
with open(self.path) as f:
saved = json.load(f)
return {**_DEFAULTS, **saved}
return dict(_DEFAULTS)
def save(self):
with open(self.path, "w") as f:
json.dump(self.state, f, indent=2)
@property
def resume_offset(self):
return self.state["last_completed_offset"]
@property
def is_complete(self):
return self.state["completed"]
def record_batch(self, new_offset, processed, skipped, errors):
self.state["last_completed_offset"] = new_offset
self.state["processed_count"] += processed
self.state["skipped_count"] += skipped
self.state["error_count"] += errors
self.save()
log.info(
"Batch done. offset=%d processed=%d skipped=%d errors=%d | "
"totals: processed=%d skipped=%d errors=%d",
new_offset, processed, skipped, errors,
self.state["processed_count"],
self.state["skipped_count"],
self.state["error_count"],
)
def mark_complete(self):
self.state["completed"] = True
self.save()
log.info(
"Migration complete. processed=%d skipped=%d errors=%d",
self.state["processed_count"],
self.state["skipped_count"],
self.state["error_count"],
)