fix: 恢复 repair 重复记账 + 禁用错误补账命令

This commit is contained in:
XingQue
2026-06-24 21:18:57 +08:00
parent a6068f7e2b
commit 600188a99a
4 changed files with 433 additions and 112 deletions

View File

@@ -1,11 +1,13 @@
"""各俱乐部全局收支流水szjilu统一记账。"""
import logging
from datetime import date
from decimal import Decimal
from django.db import IntegrityError, transaction
from django.db.models import Count, Sum
from django.db.utils import OperationalError, ProgrammingError
from config.models import PlatformIncomeLog, Szjilu
from config.models import DailyIncomeStat, PlatformIncomeLog, Szjilu
from jituan.constants import CLUB_ID_DEFAULT
logger = logging.getLogger(__name__)
@@ -26,6 +28,16 @@ def normalize_szjilu_club_id(club_id):
return (club_id or '').strip() or CLUB_ID_DEFAULT
def income_log_exists(biz_ref):
ref = (biz_ref or '').strip()
if not ref:
return False
try:
return PlatformIncomeLog.objects.filter(biz_ref=ref).exists()
except _PLATFORM_LOG_ERRORS:
return False
def _get_szjilu_for_update(club_id):
cid = normalize_szjilu_club_id(club_id)
szjilu, _ = Szjilu.objects.select_for_update().get_or_create(
@@ -47,7 +59,7 @@ def apply_szjilu_income(amount, club_id=None):
def _apply_income_stats(amount, club_id):
"""写入 daily_income_stat + szjilu核心统计,必须成功)。"""
"""写入 daily_income_stat + szjilu仅在新 platform_income_log 落库后调用)。"""
from orders.utils import update_daily_income
cid = normalize_szjilu_club_id(club_id)
@@ -58,13 +70,15 @@ def _apply_income_stats(amount, club_id):
def record_wechat_income_once(amount, club_id=None, biz_ref=None):
"""
微信入账:同一 biz_ref 只记一次 daily_income_stat + szjilu
返回 True=本次新入账False=已记过(重复回调)。
幂等日志与日统计在同一事务内提交,避免「已支付但未入账」。
微信入账:同一 biz_ref 只记一次。
返回 True=本次新入账False=已记过(重复回调,绝不重复加钱)。
"""
ref = (biz_ref or '').strip()
if not ref:
raise ValueError('biz_ref 不能为空')
if income_log_exists(ref):
return False
cid = normalize_szjilu_club_id(club_id)
jine = Decimal(str(amount))
@@ -77,34 +91,275 @@ def record_wechat_income_once(amount, club_id=None, biz_ref=None):
)
_apply_income_stats(jine, cid)
except IntegrityError:
logger.info('平台入账已存在,跳过重复记账 biz_ref=%s', ref)
return False
except _PLATFORM_LOG_ERRORS as exc:
logger.warning(
'platform_income_log 不可用(%s)跳过幂等直接记入收入 biz_ref=%s',
'platform_income_log 不可用(%s)无法幂等入账 biz_ref=%s,请 migrate config 后执行 sync_daily_income_from_log',
exc, ref,
)
_apply_income_stats(jine, cid)
return True
return False
except Exception as exc:
logger.error(
'platform_income_log 入账异常,仍尝试记入收入 biz_ref=%s err=%s',
'platform_income_log 入账失败 biz_ref=%s err=%s',
ref, exc, exc_info=True,
)
_apply_income_stats(jine, cid)
return True
return False
logger.info('平台入账记账成功 biz_ref=%s club=%s amount=%s', ref, cid, jine)
return True
def repair_wechat_income_if_missing(amount, club_id=None, biz_ref=None):
"""已支付订单/充值单重复回调时安全补记收支(幂等)"""
"""已支付单重复回调时补记;日志已存在则绝不重复加钱"""
return record_wechat_income_once(amount, club_id, biz_ref)
def sync_daily_income_stat_for_date(stat_date, reset_day=False):
"""
用 platform_income_log 重写指定日期的 daily_income_stat金额/笔数与日志完全一致)。
reset_day=True 时先删掉该日所有日统计行,再按日志重建(纠正重复记账)。
返回更新的俱乐部行数。
"""
if reset_day:
DailyIncomeStat.objects.filter(date=stat_date).delete()
try:
return record_wechat_income_once(amount, club_id, biz_ref)
except Exception as exc:
logger.error('补记收支失败 biz_ref=%s err=%s', biz_ref, exc, exc_info=True)
return False
log_qs = PlatformIncomeLog.objects.filter(CreateTime__date=stat_date)
except _PLATFORM_LOG_ERRORS as exc:
logger.error('platform_income_log 不可用: %s', exc)
return 0
club_ids = list(
log_qs.values_list('club_id', flat=True).distinct()
)
updated = 0
y, m, d = stat_date.year, stat_date.month, stat_date.day
with transaction.atomic():
for cid in club_ids:
agg = log_qs.filter(club_id=cid).aggregate(
total_amount=Sum('amount'),
total_count=Count('id'),
)
amt = agg['total_amount'] or _ZERO
cnt = int(agg['total_count'] or 0)
stat, _ = DailyIncomeStat.objects.select_for_update().get_or_create(
date=stat_date,
club_id=cid,
defaults={
'year': y,
'month': m,
'day': d,
'total_amount': amt,
'total_count': cnt,
},
)
stat.year = y
stat.month = m
stat.day = d
stat.total_amount = amt
stat.total_count = cnt
stat.save(update_fields=['year', 'month', 'day', 'total_amount', 'total_count'])
updated += 1
return updated
def sync_szjilu_income_from_logs():
"""按 platform_income_log 全量重算各俱乐部 szjilu 累计收入/流水(纠正重复累加)。"""
try:
rows = PlatformIncomeLog.objects.values('club_id').annotate(
total=Sum('amount'),
cnt=Count('id'),
)
except _PLATFORM_LOG_ERRORS as exc:
logger.error('platform_income_log 不可用: %s', exc)
return 0
today = date.today()
updated = 0
with transaction.atomic():
for row in rows:
cid = normalize_szjilu_club_id(row['club_id'])
total = row['total'] or _ZERO
today_agg = PlatformIncomeLog.objects.filter(
CreateTime__date=today, club_id=cid,
).aggregate(t=Sum('amount'))
today_flow = today_agg['t'] or _ZERO
szjilu, _ = Szjilu.objects.select_for_update().get_or_create(
club_id=cid,
defaults=dict(_SZJILU_DEFAULTS),
)
szjilu.TotalIncome = total
szjilu.TotalFlow = total
szjilu.DailyFlow = today_flow
szjilu.save(update_fields=['TotalIncome', 'TotalFlow', 'DailyFlow', 'UpdateTime'])
updated += 1
return updated
# 订单已付款且非退款repair 错误地把 Status=5 退款单也计入了收入)
_ORDER_INCOME_STATUSES = [1, 2, 3, 4, 6, 7, 8]
# 微信充值类 czjilu会员/押金/积分/商家/罚单/考核
_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), ...]
"""
from datetime import datetime, timedelta
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)
seen = {}
items = []
def _add(ref, club_id, amount, kind):
ref = (ref or '').strip()
if not ref or ref in seen:
return
jine = Decimal(str(amount or 0))
if jine <= 0:
return
cid = normalize_szjilu_club_id(club_id)
seen[ref] = True
items.append((ref, cid, jine, kind))
for od in Order.query.filter(
UpdateTime__gte=day_start,
UpdateTime__lt=day_end,
Status__in=_ORDER_INCOME_STATUSES,
):
_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,
zhuangtai=3,
leixing__in=_CZ_WECHAT_LEIXING,
):
_add(f'cz:{cz.dingdan_id}', getattr(cz, 'club_id', None), cz.jine, f'cz_{cz.leixing}')
return items
def recover_income_stats_for_date(stat_date, dry_run=True):
"""
撤销 repair_wechat_income 造成的重复记账:
1) 删除该日 platform_income_log
2) 按真实支付重建日志(每笔唯一)
3) 重写 daily_income_stat
4) 重算 szjilu 累计收入
返回 dict 供命令行展示。
"""
from datetime import datetime
from django.db.models import Sum
before_daily = DailyIncomeStat.objects.filter(date=stat_date).aggregate(
amt=Sum('total_amount'),
cnt=Sum('total_count'),
)
before_daily_by_club = {
row.club_id: {
'amt': Decimal(str(row.total_amount or 0)),
'cnt': int(row.total_count or 0),
}
for row in DailyIncomeStat.objects.filter(date=stat_date)
}
before_log_cnt = 0
before_log_amt = _ZERO
try:
before_log_cnt = PlatformIncomeLog.objects.filter(CreateTime__date=stat_date).count()
agg = PlatformIncomeLog.objects.filter(CreateTime__date=stat_date).aggregate(
t=Sum('amount'),
)
before_log_amt = agg['t'] or _ZERO
except _PLATFORM_LOG_ERRORS:
pass
items = collect_canonical_wechat_income_for_date(stat_date)
rebuild_amt = sum((x[2] for x in items), _ZERO)
rebuild_cnt = len(items)
result = {
'date': str(stat_date),
'before_daily_amount': float(before_daily['amt'] or 0),
'before_daily_count': int(before_daily['cnt'] or 0),
'before_log_amount': float(before_log_amt),
'before_log_count': before_log_cnt,
'after_amount': float(rebuild_amt),
'after_count': rebuild_cnt,
'items': items,
'dry_run': dry_run,
}
if dry_run:
return result
from datetime import time as dt_time
from django.utils import timezone
pay_ts = timezone.make_aware(datetime.combine(stat_date, dt_time(12, 0, 0)))
with transaction.atomic():
try:
PlatformIncomeLog.objects.filter(CreateTime__date=stat_date).delete()
except _PLATFORM_LOG_ERRORS as exc:
raise RuntimeError(f'platform_income_log 不可用: {exc}') from exc
DailyIncomeStat.objects.filter(date=stat_date).delete()
for ref, cid, jine, _kind in items:
log = PlatformIncomeLog.objects.create(
biz_ref=ref,
club_id=cid,
amount=jine,
)
PlatformIncomeLog.objects.filter(pk=log.pk).update(CreateTime=pay_ts)
y, m, d = stat_date.year, stat_date.month, stat_date.day
by_club = {}
for ref, cid, jine, _kind in items:
if cid not in by_club:
by_club[cid] = {'amt': _ZERO, 'cnt': 0}
by_club[cid]['amt'] += jine
by_club[cid]['cnt'] += 1
for cid, agg in by_club.items():
DailyIncomeStat.objects.create(
date=stat_date,
club_id=cid,
year=y,
month=m,
day=d,
total_amount=agg['amt'],
total_count=agg['cnt'],
)
# 只扣掉 repair 多记的部分,不用日志全量重算 szjilu避免抹掉无日志的历史收入
all_clubs = set(before_daily_by_club.keys()) | set(by_club.keys())
for cid in all_clubs:
old_amt = before_daily_by_club.get(cid, {}).get('amt', _ZERO)
new_amt = by_club.get(cid, {}).get('amt', _ZERO)
excess = old_amt - new_amt
if excess <= 0 and stat_date != date.today():
continue
szjilu = _get_szjilu_for_update(cid)
if excess > 0:
szjilu.TotalIncome -= excess
szjilu.TotalFlow -= excess
if stat_date == date.today():
szjilu.DailyFlow = new_amt
szjilu.save(update_fields=['TotalIncome', 'TotalFlow', 'DailyFlow', 'UpdateTime'])
return result
def apply_szjilu_expense(amount, club_id=None):