|
| 1 | +from django.apps import apps |
| 2 | +from django.core.exceptions import FieldDoesNotExist |
| 3 | +from django.core.management.base import BaseCommand |
| 4 | +from django.core.management.base import CommandError |
| 5 | +from django.db import transaction |
| 6 | +from django.db.models import F |
| 7 | + |
| 8 | + |
| 9 | +class Command(BaseCommand): |
| 10 | + help = ( |
| 11 | + "Idempotent, resumable online backfill of one column into another, in batches." |
| 12 | + ) |
| 13 | + |
| 14 | + def add_arguments(self, parser): |
| 15 | + parser.add_argument("--model", required=True, help="app_label.ModelName") |
| 16 | + parser.add_argument("--source-field", required=True) |
| 17 | + parser.add_argument("--target-field", required=True) |
| 18 | + parser.add_argument("--batch-size", type=int, default=10000) |
| 19 | + parser.add_argument("--start-id", default=None, help="resume from this pk") |
| 20 | + parser.add_argument( |
| 21 | + "--progress-check", |
| 22 | + action="store_true", |
| 23 | + help="report unbackfilled rows, exit nonzero if any", |
| 24 | + ) |
| 25 | + |
| 26 | + def _resolve_model_fields(self, model_label, source, target): |
| 27 | + try: |
| 28 | + model = apps.get_model(model_label) |
| 29 | + except (LookupError, ValueError) as e: |
| 30 | + raise CommandError("Bad --model {!r}: {}".format(model_label, e)) |
| 31 | + try: |
| 32 | + model._meta.get_field(source) |
| 33 | + model._meta.get_field(target) |
| 34 | + except FieldDoesNotExist as e: |
| 35 | + raise CommandError(str(e)) |
| 36 | + return model |
| 37 | + |
| 38 | + def _batch_end_pk(self, queryset, pk_name, start_pk, batch_size): |
| 39 | + """Last pk of the batch of `batch_size` rows starting at `start_pk`. |
| 40 | +
|
| 41 | + Returns None when fewer than `batch_size` rows remain at/after |
| 42 | + `start_pk` — the final, short batch. Keyset paging by pk, so it works |
| 43 | + for any pk type (int or UUID). |
| 44 | + """ |
| 45 | + return ( |
| 46 | + queryset.filter(pk__gte=start_pk) |
| 47 | + .order_by(pk_name) |
| 48 | + .values_list("pk", flat=True)[batch_size - 1 : batch_size] |
| 49 | + .first() |
| 50 | + ) |
| 51 | + |
| 52 | + def handle(self, *args, **options): |
| 53 | + if options["batch_size"] < 1: |
| 54 | + raise CommandError("--batch-size must be >= 1") |
| 55 | + source = options["source_field"] |
| 56 | + target = options["target_field"] |
| 57 | + model = self._resolve_model_fields(options["model"], source, target) |
| 58 | + |
| 59 | + pk_name = model._meta.pk.name |
| 60 | + batch_size = options["batch_size"] |
| 61 | + only_unfilled = {target + "__isnull": True, source + "__isnull": False} |
| 62 | + unfilled = model.objects.filter(**only_unfilled) |
| 63 | + unfilled_pks = unfilled.order_by(pk_name).values_list("pk", flat=True) |
| 64 | + |
| 65 | + if options["progress_check"]: |
| 66 | + # exists(), not count() — the target table can have millions of rows. |
| 67 | + if unfilled.exists(): |
| 68 | + raise CommandError("backfill incomplete: rows still pending") |
| 69 | + self.stdout.write("Backfill complete: no rows pending.") |
| 70 | + return |
| 71 | + |
| 72 | + # Start at the first unfilled pk (>= --start-id if given); re-runs and |
| 73 | + # resumes skip straight past an already-filled prefix. |
| 74 | + batch_start = unfilled_pks |
| 75 | + if options["start_id"] is not None: |
| 76 | + batch_start = batch_start.filter(pk__gte=options["start_id"]) |
| 77 | + batch_start = batch_start.first() |
| 78 | + |
| 79 | + total = 0 |
| 80 | + while batch_start is not None: |
| 81 | + batch_end = self._batch_end_pk( |
| 82 | + model.objects, pk_name, batch_start, batch_size |
| 83 | + ) |
| 84 | + if batch_end is None: |
| 85 | + window = {"pk__gte": batch_start} |
| 86 | + else: |
| 87 | + window = {"pk__gte": batch_start, "pk__lte": batch_end} |
| 88 | + with transaction.atomic(): |
| 89 | + total += model.objects.filter(**window, **only_unfilled).update( |
| 90 | + **{target: F(source)} |
| 91 | + ) |
| 92 | + self.stdout.write( |
| 93 | + "backfilled through pk={} (updated {} so far)".format( |
| 94 | + batch_start if batch_end is None else batch_end, total |
| 95 | + ) |
| 96 | + ) |
| 97 | + if batch_end is None: |
| 98 | + break |
| 99 | + batch_start = unfilled_pks.filter(pk__gt=batch_end).first() |
| 100 | + self.stdout.write("Done. {} rows updated.".format(total)) |
0 commit comments