Tracks per-file processing progress for resumable execution. Persists a JSON file mapping each data file to its last processed index, allowing the pipeline to resume from the exact interruption point.
| 7 | |
| 8 | |
| 9 | class CheckpointManager: |
| 10 | """Tracks per-file processing progress for resumable execution. |
| 11 | |
| 12 | Persists a JSON file mapping each data file to its last processed |
| 13 | index, allowing the pipeline to resume from the exact interruption |
| 14 | point. |
| 15 | """ |
| 16 | |
| 17 | def __init__( |
| 18 | self, work_dir: str, checkpoint_name: str = "checkpoint.json", logger=None |
| 19 | ): |
| 20 | """Initialize checkpoint manager. |
| 21 | |
| 22 | Args: |
| 23 | work_dir: Working directory for checkpoint file. |
| 24 | checkpoint_name: Name of checkpoint file. |
| 25 | logger: Logger instance. |
| 26 | """ |
| 27 | self.work_dir = work_dir |
| 28 | self.checkpoint_file = os.path.join(work_dir, checkpoint_name) |
| 29 | self.logger = logger |
| 30 | self.checkpoints: Dict[str, dict] = {} |
| 31 | |
| 32 | os.makedirs(work_dir, exist_ok=True) |
| 33 | self._load() |
| 34 | |
| 35 | def _load(self) -> None: |
| 36 | """Load checkpoint from file.""" |
| 37 | if os.path.exists(self.checkpoint_file): |
| 38 | try: |
| 39 | with open(self.checkpoint_file, "r", encoding="utf-8") as f: |
| 40 | self.checkpoints = json.load(f) |
| 41 | if self.logger: |
| 42 | self.logger.info(f"Loaded checkpoint from {self.checkpoint_file}") |
| 43 | except Exception as e: |
| 44 | if self.logger: |
| 45 | self.logger.warning(f"Failed to load checkpoint: {e}") |
| 46 | self.checkpoints = {} |
| 47 | |
| 48 | def _save(self) -> None: |
| 49 | """Save checkpoint to file.""" |
| 50 | try: |
| 51 | # Write to temp file first, then rename for atomicity |
| 52 | temp_file = self.checkpoint_file + ".tmp" |
| 53 | with open(temp_file, "w", encoding="utf-8") as f: |
| 54 | json.dump(self.checkpoints, f, indent=2, ensure_ascii=False) |
| 55 | try: |
| 56 | os.replace(temp_file, self.checkpoint_file) |
| 57 | except OSError: |
| 58 | # Fallback for cross-device moves: copy and delete |
| 59 | import shutil |
| 60 | |
| 61 | shutil.copy2(temp_file, self.checkpoint_file) |
| 62 | os.remove(temp_file) |
| 63 | except Exception as e: |
| 64 | if self.logger: |
| 65 | self.logger.error(f"Failed to save checkpoint: {e}") |
| 66 |
no outgoing calls