| 27 | |
| 28 | |
| 29 | class DeleteImportItemJob: |
| 30 | def __init__(self, in_file="", state_file="", error_file="", batch_size=1000, dry_run=False): |
| 31 | self.in_file = in_file |
| 32 | self.state_file = state_file |
| 33 | self.error_file = error_file |
| 34 | |
| 35 | self.batch_size = batch_size |
| 36 | self.dry_run = dry_run |
| 37 | |
| 38 | self.start_line = 1 |
| 39 | state_path = Path(state_file) |
| 40 | if state_path.exists(): |
| 41 | with state_path.open("r") as f: |
| 42 | line = f.readline() |
| 43 | if line: |
| 44 | self.start_line = int(line) |
| 45 | |
| 46 | def run(self): |
| 47 | with open(self.in_file) as f: |
| 48 | # Seek to start line |
| 49 | for _ in range(1, self.start_line): |
| 50 | f.readline() |
| 51 | |
| 52 | # Delete a batch of records |
| 53 | lines_processed = 0 |
| 54 | affected_records = 0 |
| 55 | num_deleted = 0 |
| 56 | for _ in range(self.batch_size): |
| 57 | line = f.readline() |
| 58 | if not line: |
| 59 | break |
| 60 | fields = line.strip().split("\t") |
| 61 | ia_ids = fields[2:] |
| 62 | try: |
| 63 | result = ImportItem.delete_items(ia_ids, _test=self.dry_run) |
| 64 | if self.dry_run: |
| 65 | # Result is string "DELETE FROM ..." |
| 66 | print(result) |
| 67 | else: |
| 68 | # Result is number of records deleted |
| 69 | num_deleted += result |
| 70 | except Exception as e: |
| 71 | print(f"Error when deleting: {e}") |
| 72 | if not self.dry_run: |
| 73 | write_to(self.error_file, line, mode="a+") |
| 74 | lines_processed += 1 |
| 75 | affected_records += int(fields[0]) |
| 76 | |
| 77 | # Write next line number to state file: |
| 78 | if not self.dry_run: |
| 79 | write_to(self.state_file, f"{self.start_line + lines_processed}") |
| 80 | |
| 81 | return { |
| 82 | "lines_processed": lines_processed, |
| 83 | "num_deleted": num_deleted, |
| 84 | "affected_records": affected_records, |
| 85 | } |
| 86 | |