From 3a8bd852d5c91987c6a86d9756a48e7d26b82cc8 Mon Sep 17 00:00:00 2001 From: ruthra kumar Date: Fri, 31 Jul 2026 10:14:50 +0530 Subject: [PATCH] refactor: rebuild pcv on mapreduce (parallelization) --- .../period_closing_voucher.py | 180 +++++++++++++++++- 1 file changed, 175 insertions(+), 5 deletions(-) diff --git a/erpnext/accounts/doctype/period_closing_voucher/period_closing_voucher.py b/erpnext/accounts/doctype/period_closing_voucher/period_closing_voucher.py index 1f4a60e6f14..9d424f9b893 100644 --- a/erpnext/accounts/doctype/period_closing_voucher/period_closing_voucher.py +++ b/erpnext/accounts/doctype/period_closing_voucher/period_closing_voucher.py @@ -3,11 +3,22 @@ import copy +from datetime import timedelta import frappe -from frappe import _ -from frappe.query_builder.functions import Max, Sum -from frappe.utils import add_days, flt, fmt_money, formatdate, get_link_to_form, getdate +from frappe import _, qb +from frappe.query_builder.functions import Max, Min, Sum +from frappe.utils import ( + add_days, + ceil, + cint, + flt, + fmt_money, + formatdate, + get_datetime, + get_link_to_form, + getdate, +) from erpnext import is_perpetual_inventory_enabled from erpnext.accounts.doctype.account_closing_balance.account_closing_balance import ( @@ -265,8 +276,14 @@ class PeriodClosingVoucher(AccountsController): if frappe.get_single_value("Accounts Settings", "use_legacy_controller_for_pcv"): self.make_gl_entries() else: - ppcv = frappe.get_doc({"doctype": "Process Period Closing Voucher", "parent_pcv": self.name}) - ppcv.save().submit() + from frappe.utils.background_jobs import mapreduce + + data = self.get_data_for_mapreduce() + mapreduce( + "erpnext.accounts.doctype.period_closing_voucher.period_closing_voucher.mapper", + "erpnext.accounts.doctype.period_closing_voucher.period_closing_voucher.reducer", + data, + ) def on_cancel(self): self.ignore_linked_doctypes = ( @@ -594,6 +611,91 @@ class PeriodClosingVoucher(AccountsController): {"voucher_type": "Period Closing Voucher", "voucher_no": self.name, "is_cancelled": 0}, ) + def get_data_for_mapreduce(self): + return self.generate_tasks_for_normal_balance() + self.generate_tasks_for_opening_balance() + + def get_period_range_for_tasks(self, start_date, end_date, step_size, report_type, balance_type): + start_date = getdate(start_date) + end_date = getdate(end_date) + + # split period into date ranges + curr_date = getdate(start_date) + date_splits = [] + while True: + next_date = getdate(add_days(curr_date, step_size)) + if next_date < end_date: + date_splits.append( + { + "from_date": str(curr_date), + "to_date": str(next_date), + "pcv": self.name, + "report_type": report_type, + "balance_type": balance_type, + } + ) + curr_date = getdate(add_days(next_date, 1)) + else: + date_splits.append( + { + "from_date": str(curr_date), + "to_date": str(end_date), + "pcv": self.name, + "report_type": report_type, + "balance_type": balance_type, + } + ) + break + + return date_splits + + def generate_tasks_for_normal_balance(self): + # estimation can be wrong by a factor of 2 + estimated_count = ( + cint( + frappe.db.sql( + f"explain select count(*) from `tabGL Entry` where is_cancelled = 0 and posting_date between {self.period_start_date} and {self.period_end_date};", + as_dict=True, + )[0].rows + ) + * 2 + ) + job_count = ( + 1 if estimated_count / 2000000 else ceil(estimated_count / 2000000) + ) # conservative chunk size + days = (getdate(self.period_end_date) - getdate(self.period_start_date)).days + step_size = 1 if days / job_count < 1 else ceil(days / job_count) + return self.get_period_range_for_tasks( + self.period_start_date, self.period_end_date, step_size, "Balance Sheet", "Normal Balance" + ) + self.get_period_range_for_tasks( + self.period_start_date, self.period_end_date, step_size, "Profit and Loss", "Normal Balance" + ) + + def generate_tasks_for_opening_balance(self): + tasks = [] + if self.is_first_period_closing_voucher(): + gl = qb.DocType("GL Entry") + min = qb.from_(gl).select(Min(gl.posting_date)).run()[0][0] + max = qb.from_(gl).select(Max(gl.posting_date)).run()[0][0] + + # estimation can be wrong by a factor of 2 + estimated_count = ( + cint( + frappe.db.sql( + f"explain select count(*) from `tabGL Entry` where is_cancelled = 0 and is_opening = 0 and posting_date between {min} and {max};", + as_dict=True, + )[0].rows + ) + * 2 + ) + job_count = ( + 1 if estimated_count / 2000000 else ceil(estimated_count / 2000000) + ) # conservative chunk size + days = (getdate(self.period_end_date) - getdate(self.period_start_date)).days + step_size = 1 if days / job_count < 1 else ceil(days / job_count) + tasks = self.get_period_range_for_tasks(min, max, step_size, "Balance Sheet", "Opening Balance") + + return tasks + def process_gl_and_closing_entries(doc): from erpnext.accounts.general_ledger import make_gl_entries @@ -673,3 +775,71 @@ def get_previous_closed_period_in_current_year(fiscal_year, company): order_by="period_end_date desc", ) return prev_closed_period_end_date + + +def mapper(val): + start_date = val.from_date + end_date = val.to_date + pcv = val.pcv + report_type = val.report_type + balance_type = val.balance_type + company = frappe.db.get_value("Period Closing Voucher", pcv, "company") + dimensions = get_dimensions() + + accounts = frappe.db.get_all( + "Account", filters={"company": company, "report_type": report_type}, pluck="name" + ) + + # summarize + gle = qb.DocType("GL Entry") + query = qb.from_(gle).select(gle.account) + for dim in dimensions: + query = query.select(gle[dim]) + query = query.select( + Sum(gle.debit).as_("debit"), + Sum(gle.credit).as_("credit"), + Sum(gle.debit_in_account_currency).as_("debit_in_account_currency"), + Sum(gle.credit_in_account_currency).as_("credit_in_account_currency"), + # account_currency is constant per grouped account -> Max() keeps the GROUP BY postgres-valid + Max(gle.account_currency).as_("account_currency"), + ).where( + (gle.company.eq(company)) + & (gle.is_cancelled.eq(0)) + & (gle.posting_date.between(start_date, end_date)) + & (gle.account.isin(accounts)) + ) + + if balance_type == "Opening Balance": + query = query.where(gle.is_opening.eq("Yes")) + else: + # Keep balances aligned with legacy PCV logic (non-opening transactions only) + query = query.where(gle.is_opening.eq("No")) + + query = query.groupby(gle.account) + for dim in dimensions: + query = query.groupby(gle[dim]) + + res = query.run(as_dict=True) + return res + + +def reducer(final, partial_res): + if final is None: + final = [] + + gl_entries = [] + if partial_res: + for x in partial_res: + gl_entries.append(frappe._dict(x)) + + return final + gl_entries + + +def get_dimensions(): + from erpnext.accounts.doctype.accounting_dimension.accounting_dimension import ( + get_accounting_dimensions, + ) + + default_dimensions = ["cost_center", "finance_book", "project"] + dimensions = default_dimensions + get_accounting_dimensions() + return dimensions