diff --git a/config/management/commands/recover_wechat_income_stats.py b/config/management/commands/recover_wechat_income_stats.py index 4532579..ade460c 100644 --- a/config/management/commands/recover_wechat_income_stats.py +++ b/config/management/commands/recover_wechat_income_stats.py @@ -1,38 +1,39 @@ """ 恢复被 repair_wechat_income 搞乱的收入统计。 -repair 会对已入账订单重复加钱(200 变 900 就是这个)。 -本命令:删掉当日错误日志 → 按真实微信支付重建 → 重算日统计和 szjilu。 +repair 会对已入账订单重复加钱(200 变 900)。 +本命令:删当日错误日志 → 按真实微信支付重建 → 扣回 szjilu 多记部分。 -先不加 --yes 可预览重建后金额。 +用法: + 1. 先预览:python manage.py recover_wechat_income_stats + 2. 看「重建后」金额是否对(应接近真实收入,如 200 左右) + 3. 确认后:python manage.py recover_wechat_income_stats --yes + 4. 若多日错乱:python manage.py recover_wechat_income_stats --days 7 --yes + 5. 最后全量对齐日统计表:python manage.py recover_wechat_income_stats --yes --rebuild-all-daily """ from datetime import date, timedelta from django.core.management.base import BaseCommand +from django.db.models import Sum -from jituan.services.szjilu_accounting import recover_income_stats_for_date +from config.models import DailyIncomeStat +from jituan.services.szjilu_accounting import ( + rebuild_all_daily_income_from_logs, + recover_income_stats_for_date, +) class Command(BaseCommand): - help = '恢复 repair_wechat_income 造成的重复记账(先预览,再加 --yes 执行)' + help = '恢复 repair_wechat_income 重复记账(先预览,--yes 执行)' def add_arguments(self, parser): + parser.add_argument('--date', type=str, default='', help='日期 YYYY-MM-DD,默认今天') + parser.add_argument('--days', type=int, default=1, help='向前恢复几天') + parser.add_argument('--yes', action='store_true', help='确认执行') parser.add_argument( - '--date', - type=str, - default='', - help='恢复日期 YYYY-MM-DD,默认今天', - ) - parser.add_argument( - '--days', - type=int, - default=1, - help='从指定日期起向前恢复几天,默认 1', - ) - parser.add_argument( - '--yes', + '--rebuild-all-daily', action='store_true', - help='确认执行(不加则仅预览,不改数据库)', + help='恢复后按 platform_income_log 全量重写 daily_income_stat(纠正累计总收入)', ) def handle(self, *args, **options): @@ -46,40 +47,46 @@ class Command(BaseCommand): if dry_run: self.stdout.write(self.style.WARNING( - '【预览模式】未改数据库。确认金额无误后加 --yes 执行。\n' + '【预览】未改数据库。下面「重建后」应是正确金额。\n' )) else: - self.stdout.write(self.style.WARNING('【执行模式】正在恢复数据...\n')) + self.stdout.write(self.style.WARNING('【执行】正在恢复...\n')) for i in range(days): d = start_day + timedelta(days=i) r = recover_income_stats_for_date(d, dry_run=dry_run) - self.stdout.write(f"=== {r['date']} ===") + self.stdout.write(f'=== {r["date"]} ===') self.stdout.write( - f" 当前日统计(错乱): {r['before_daily_amount']:.2f} 元 / {r['before_daily_count']} 笔" + f' 错乱前日统计: {r["before_daily_amount"]:.2f} 元 / {r["before_daily_count"]} 笔' ) self.stdout.write( - f" 当前幂等日志: {r['before_log_amount']:.2f} 元 / {r['before_log_count']} 笔" + f' 当日幂等日志: {r["before_log_amount"]:.2f} 元 / {r["before_log_count"]} 笔' ) self.stdout.write(self.style.SUCCESS( - f" 重建后(真实支付): {r['after_amount']:.2f} 元 / {r['after_count']} 笔" + f' 重建后(正确值): {r["after_amount"]:.2f} 元 / {r["after_count"]} 笔' )) - if r['items'] and dry_run: - self.stdout.write(' 明细(前20笔):') - for ref, cid, amt, kind in r['items'][:20]: + if r['items']: + self.stdout.write(' 明细:') + for ref, cid, amt, kind in r['items']: self.stdout.write(f' {ref} {amt}元 club={cid} ({kind})') - if len(r['items']) > 20: - self.stdout.write(f' ... 共 {len(r["items"])} 笔') - if dry_run: - self.stdout.write(self.style.WARNING( - '\n预览完成。若重建后金额正确,执行:\n' - ' python manage.py recover_wechat_income_stats --yes\n' - '指定日期:\n' - ' python manage.py recover_wechat_income_stats --date 2026-06-24 --yes' - )) - else: + if not dry_run and options['rebuild_all_daily']: + n = rebuild_all_daily_income_from_logs() + self.stdout.write(self.style.SUCCESS(f'已按日志全量重写 daily_income_stat: {n} 行')) + + if not dry_run: + agg = DailyIncomeStat.objects.filter(date=end_day).aggregate( + amt=Sum('total_amount'), + cnt=Sum('total_count'), + ) self.stdout.write(self.style.SUCCESS( - '\n恢复完成。请刷新后台财务页核对今日收入。' - '\n勿再执行 repair_wechat_income(已禁用)。' + f'\n数据库今日日统计现为: {float(agg["amt"] or 0):.2f} 元 / {int(agg["cnt"] or 0)} 笔' + )) + self.stdout.write(self.style.SUCCESS('恢复完成,请刷新后台财务页。')) + else: + self.stdout.write(self.style.WARNING( + '\n预览完成。重建后金额对的话执行:\n' + ' python manage.py recover_wechat_income_stats --yes\n' + '累计总收入仍不对再加:\n' + ' python manage.py recover_wechat_income_stats --yes --rebuild-all-daily' )) diff --git a/jituan/services/szjilu_accounting.py b/jituan/services/szjilu_accounting.py index de6b210..a39ab2a 100644 --- a/jituan/services/szjilu_accounting.py +++ b/jituan/services/szjilu_accounting.py @@ -204,24 +204,29 @@ _ORDER_INCOME_STATUSES = [1, 2, 3, 4, 6, 7, 8] _CZ_WECHAT_LEIXING = [1, 2, 3, 4, 5] -def collect_canonical_wechat_income_for_date(stat_date): - """ - 按真实支付时间(UpdateTime)收集当日微信收入,每笔 biz_ref 唯一。 - 返回 [(biz_ref, club_id, amount, kind), ...] +def collect_canonical_wechat_income_for_date(stat_date, existing_refs=None): """ + 收集当日真实微信收入(每笔 biz_ref 全局唯一,不重复计已入账订单)。 + - 订单:当日创建且已付款,或当日完成微信支付(有微信交易号) + - 充值:当日创建且 zhuangtai=3 + - 不含退款单(Status=5) + """ from datetime import datetime, timedelta + from django.db.models import Q + from orders.models import Order from products.models import Czjilu day_start = datetime.combine(stat_date, datetime.min.time()) day_end = day_start + timedelta(days=1) + skip_refs = existing_refs or set() seen = {} items = [] def _add(ref, club_id, amount, kind): ref = (ref or '').strip() - if not ref or ref in seen: + if not ref or ref in seen or ref in skip_refs: return jine = Decimal(str(amount or 0)) if jine <= 0: @@ -230,16 +235,25 @@ def collect_canonical_wechat_income_for_date(stat_date): seen[ref] = True items.append((ref, cid, jine, kind)) - for od in Order.query.filter( - UpdateTime__gte=day_start, - UpdateTime__lt=day_end, + order_qs = Order.query.filter( Status__in=_ORDER_INCOME_STATUSES, - ): + ).filter( + Q(CreateTime__gte=day_start, CreateTime__lt=day_end) + | Q( + UpdateTime__gte=day_start, + UpdateTime__lt=day_end, + WechatTransactionID__isnull=False, + ), + ) + for od in order_qs: + txn = (getattr(od, 'WechatTransactionID', None) or '').strip() + if not txn: + continue _add(f'order:{od.OrderID}', getattr(od, 'ClubID', None), od.Amount, 'order') for cz in Czjilu.query.filter( - UpdateTime__gte=day_start, - UpdateTime__lt=day_end, + CreateTime__gte=day_start, + CreateTime__lt=day_end, zhuangtai=3, leixing__in=_CZ_WECHAT_LEIXING, ): @@ -248,6 +262,53 @@ def collect_canonical_wechat_income_for_date(stat_date): return items +def _existing_refs_outside_date(stat_date): + """该日之外已入账的 biz_ref(删当日错误日志后仍保留的)。""" + try: + return set( + PlatformIncomeLog.objects.exclude( + CreateTime__date=stat_date, + ).values_list('biz_ref', flat=True) + ) + except _PLATFORM_LOG_ERRORS: + return set() + + +def rebuild_all_daily_income_from_logs(): + """按 platform_income_log 全量重写 daily_income_stat(日志为唯一准绳)。""" + from django.db.models.functions import TruncDate + + try: + grouped = ( + PlatformIncomeLog.objects.annotate(log_date=TruncDate('CreateTime')) + .values('log_date', 'club_id') + .annotate(total_amount=Sum('amount'), total_count=Count('id')) + ) + except _PLATFORM_LOG_ERRORS as exc: + logger.error('platform_income_log 不可用: %s', exc) + return 0 + + with transaction.atomic(): + DailyIncomeStat.objects.all().delete() + n = 0 + for row in grouped: + d = row['log_date'] + if not d: + continue + cid = normalize_szjilu_club_id(row['club_id']) + DailyIncomeStat.objects.create( + date=d, + club_id=cid, + year=d.year, + month=d.month, + day=d.day, + total_amount=row['total_amount'] or _ZERO, + total_count=int(row['total_count'] or 0), + ) + n += 1 + return n + + def recover_income_stats_for_date(stat_date, dry_run=True): """ 撤销 repair_wechat_income 造成的重复记账: @@ -283,7 +344,8 @@ def recover_income_stats_for_date(stat_date, dry_run=True): except _PLATFORM_LOG_ERRORS: pass - items = collect_canonical_wechat_income_for_date(stat_date) + existing_refs = _existing_refs_outside_date(stat_date) + items = collect_canonical_wechat_income_for_date(stat_date, existing_refs=existing_refs) rebuild_amt = sum((x[2] for x in items), _ZERO) rebuild_cnt = len(items) @@ -316,6 +378,9 @@ def recover_income_stats_for_date(stat_date, dry_run=True): DailyIncomeStat.objects.filter(date=stat_date).delete() + existing_refs = _existing_refs_outside_date(stat_date) + items = collect_canonical_wechat_income_for_date(stat_date, existing_refs=existing_refs) + for ref, cid, jine, _kind in items: log = PlatformIncomeLog.objects.create( biz_ref=ref,