From 86fe0c1f4b5df1a2f542f153238e5f7e380ef2dc Mon Sep 17 00:00:00 2001 From: rohitwaghchaure Date: Wed, 16 Sep 2026 15:23:08 +0530 Subject: [PATCH] perf(stock): chunk the serial and batch entry backfill patch (#59076) * perf(stock): chunk the serial and batch entry backfill patch * fix(stock): support postgres in the serial and batch entry backfill patch --- .../v16_0/update_serial_batch_entries.py | 148 ++++++++++++++++-- 1 file changed, 133 insertions(+), 15 deletions(-) diff --git a/erpnext/patches/v16_0/update_serial_batch_entries.py b/erpnext/patches/v16_0/update_serial_batch_entries.py index a2391edd57f..ea91b827079 100644 --- a/erpnext/patches/v16_0/update_serial_batch_entries.py +++ b/erpnext/patches/v16_0/update_serial_batch_entries.py @@ -1,19 +1,137 @@ +import time + import frappe +CHILD_TABLE = "tabSerial and Batch Entry" + +BUNDLE_TABLE = "tabSerial and Batch Bundle" + +CHUNK_SIZE = 50_000 + +# Denormalised columns copied from the bundle onto every entry row. +COLUMNS = ( + "posting_datetime", + "voucher_type", + "voucher_no", + "voucher_detail_no", + "type_of_transaction", + "is_cancelled", + "item_code", +) + def execute(): - if frappe.db.has_table("Serial and Batch Entry"): - frappe.db.sql( - """ - UPDATE `tabSerial and Batch Entry` SABE, `tabSerial and Batch Bundle` SABB - SET - SABE.posting_datetime = SABB.posting_datetime, - SABE.voucher_type = SABB.voucher_type, - SABE.voucher_no = SABB.voucher_no, - SABE.voucher_detail_no = SABB.voucher_detail_no, - SABE.type_of_transaction = SABB.type_of_transaction, - SABE.is_cancelled = SABB.is_cancelled, - SABE.item_code = SABB.item_code - WHERE SABE.parent = SABB.name - """ - ) + if not frappe.db.has_table("Serial and Batch Entry"): + return + + # Only ever used to give the log a denominator. It is a cached information_schema + # estimate, so it can read 0 for a table that has rows -- gating the backfill on it + # would silently skip the whole migration. The loop below decides when it is done. + total = frappe.db.estimate_count("Serial and Batch Entry") + + last_name = "" + done = 0 + started_at = time.monotonic() + + while True: + upper = get_chunk_end(last_name) + + update_chunk(last_name, upper) + + # Commit per chunk. Doing every row in one transaction grows the undo log until + # each read has to walk it, which is what made this run for hours on large sites. + frappe.db.commit() + + if not upper: + # The tail is whatever was left after the last boundary, so it has to be + # counted rather than assumed. Only ever scans a sub-chunk range. + done += frappe.db.count("Serial and Batch Entry", {"name": (">", last_name)}) + log_progress(done, total, started_at) + break + + # A bounded chunk is exactly CHUNK_SIZE rows by construction. + done += CHUNK_SIZE + last_name = upper + log_progress(done, total, started_at) + + +def get_chunk_end(last_name): + """Return the name that closes the next chunk, or None when the tail is left. + + Keyset pagination, so each chunk is a sequential range scan on the clustered + index rather than a deep OFFSET over the whole table. + """ + entry = frappe.qb.DocType("Serial and Batch Entry") + + boundary = ( + frappe.qb.from_(entry) + .select(entry.name) + .where(entry.name > last_name) + .orderby(entry.name) + .limit(1) + .offset(CHUNK_SIZE - 1) + ).run(pluck=True) + + return boundary[0] if boundary else None + + +def update_chunk(last_name, upper): + """Copy the bundle's values onto one chunk of entries. + + Raw SQL because the query builder cannot express this statement. Every one of + COLUMNS exists on both tables, and pypika renders the assignment target without + its table (`_set_sql` forces `with_namespace=False`), so a joined UPDATE fails + with "Column 'voucher_no' in field list is ambiguous". Aliasing the bundle in a + derived table clears the ambiguity but makes MariaDB materialise the whole bundle + table for every chunk and drive the join from it, and a correlated subquery per + column costs one lookup per column per row instead of one per row. + + The two dialects spell a joined UPDATE differently and Frappe does not translate + between them, so each gets its own statement. + """ + condition = "AND SABE.name <= %(upper)s" if upper else "" + + frappe.db.multisql( + { + "mariadb": get_mariadb_query(condition), + "postgres": get_postgres_query(condition), + }, + {"last_name": last_name, "upper": upper}, + ) + + +def get_mariadb_query(condition): + set_clause = ",\n\t\t\t\t".join(f"SABE.{column} = SABB.{column}" for column in COLUMNS) + + return f""" + UPDATE `{CHILD_TABLE}` SABE + INNER JOIN `{BUNDLE_TABLE}` SABB + ON SABE.parent = SABB.name + SET + {set_clause} + WHERE SABE.name > %(last_name)s {condition} + """ + + +def get_postgres_query(condition): + # Postgres joins through FROM rather than JOIN, and rejects the table alias on the + # assignment target, so the SET columns are bare here. + set_clause = ",\n\t\t\t\t".join(f"{column} = SABB.{column}" for column in COLUMNS) + + return f""" + UPDATE `{CHILD_TABLE}` SABE + SET + {set_clause} + FROM `{BUNDLE_TABLE}` SABB + WHERE SABE.parent = SABB.name + AND SABE.name > %(last_name)s {condition} + """ + + +def log_progress(done, total, started_at): + elapsed = time.monotonic() - started_at + rate = done / elapsed if elapsed else 0 + print( + f"Serial and Batch Entry: {done:,} rows of ~{total:,} ({rate:,.0f} rows/sec)", + flush=True, + )