import threading import requests from django.conf import settings from peizhi.models import Club from shangpin.models import ShangpinLeixing from dingdan.models import Dingdan # 根据实际路径导入 import logging from datetime import date from django.db import transaction from django.db.models import F from peizhi.models import DailyDispatchStat # 假设模型在 peizhi 应用中 logger = logging.getLogger(__name__) # 可选:线程池限制(避免大量线程同时创建,建议使用 Celery,这里简单用线程) THREAD_POOL_SIZE = 10 _semaphore = threading.Semaphore(THREAD_POOL_SIZE) def sync_order_to_partners(dingdan_id): """ 异步向所有开启互通的合作平台同步订单。 使用信号量限制并发线程数。 """ def _task(): _semaphore.acquire() try: _sync_order_task(dingdan_id) finally: _semaphore.release() thread = threading.Thread(target=_task) thread.daemon = True thread.start() def _sync_order_task(dingdan_id): """实际同步任务(在线程中执行)""" try: # 获取订单对象(只取必要字段,提升性能) dingdan = Dingdan.objects.only( 'dingdan_id', 'zhuangtai', 'jine', 'dashou_fencheng', 'jieshao', 'beizhu', 'nicheng', 'leixing_id', 'yaoqiuleixing', 'huiyuan_id', 'yongjin', 'zhiding_id', 'create_time' ).get(dingdan_id=dingdan_id) # 基本条件判断 if dingdan.zhuangtai != 1: logger.info(f"订单 {dingdan_id} 状态 {dingdan.zhuangtai} 不是1,不同步") return if dingdan.zhiding_id: logger.info(f"订单 {dingdan_id} 是指定单,不同步") return if dingdan.yaoqiuleixing != 1: logger.info(f"订单 {dingdan_id} 抢单要求类型不是会员抢单,不同步") return if not dingdan.leixing_id: logger.info(f"订单 {dingdan_id} 无商品类型,不同步") return try: leixing = ShangpinLeixing.objects.get(id=dingdan.leixing_id) if not leixing.is_cross_enabled: logger.info(f"订单 {dingdan_id} 的商品类型未开启互通,不同步") return except ShangpinLeixing.DoesNotExist: logger.info(f"订单 {dingdan_id} 的商品类型不存在,不同步") return # 查询所有开启互通的俱乐部 clubs = Club.objects.filter(is_interop_enabled=True).select_related() if not clubs.exists(): return # 为每个俱乐部发送订单 for club in clubs: _send_to_single_partner(dingdan, club) except Exception as e: logger.error(f"同步订单任务异常: {e}", exc_info=True) def _send_to_single_partner(dingdan, club): """向单个俱乐部发送订单(带重试,幂等处理)""" # 计算实际金额:订单打手分成 × 对方费率 actual_amount = dingdan.dashou_fencheng * club.dashou_rate # 构建请求数据 payload = { 'token': settings.CROSS_PLATFORM_TOKEN, 'our_club_id': club.club_id, # 我方在对方系统中的标识(存储于对方表的 partner_club_id) 'partner_club_id': club.partner_club_id, # 对方自己的俱乐部ID 'partner_order_id': dingdan.dingdan_id, # 我方订单ID作为对方的外部订单号 'jine': str(actual_amount), # 对方平台应展示的订单金额 'dashou_fencheng': str(actual_amount), # 对方平台打手分成(与金额一致) 'jieshao': dingdan.jieshao, 'beizhu': dingdan.beizhu, 'create_time': dingdan.create_time.isoformat(), 'nicheng': dingdan.nicheng, 'leixing_id': dingdan.leixing_id, # 我方商品类型ID,对方需要映射 'user_id': get_order_user_id(dingdan), # 对方平台下单用户ID(不带前缀) } # 重试机制(最多3次,指数退避) for attempt in range(1, 4): try: response = requests.post( club.partner_domain + '/dingdan/partner_sync_order', json=payload, timeout=5 ) if response.status_code == 200: data = response.json() # 对方返回 code=0 表示成功(包括新建或已存在) if data.get('code') == 0: # 判断本次响应是新创建还是已存在 status = data.get('data', {}).get('status', '') if status == 'existed': # 订单已经存在于对方系统,不重复统计 logger.info(f"订单 {dingdan.dingdan_id} 在俱乐部 {club.club_id} 已存在,跳过统计") else: # 新创建成功,更新统计 try: update_daily_dispatch_stat( direction=1, # 我方派单 partner_club_id=club.partner_club_id, dispatch_amount=actual_amount, ) logger.info(f"成功同步订单 {dingdan.dingdan_id} 到俱乐部 {club.club_id}") except Exception as e: logger.error(f"更新派单统计失败: {e}", exc_info=True) # 标记订单为跨平台(仅当第一次成功时标记,避免重复标记) if not dingdan.is_cross: Dingdan.objects.filter(dingdan_id=dingdan.dingdan_id).update( is_cross=1, dispatch_type=1 # 我方派单 ) logger.info(f"订单 {dingdan.dingdan_id} 已标记为跨平台订单") break else: # 对方返回业务错误(如映射不存在等),不重试 logger.error(f"同步订单失败(业务): {data}") break else: logger.error(f"同步订单HTTP失败: {response.status_code}") except Exception as e: logger.error(f"同步订单异常(尝试 {attempt}/3): {e}") if attempt < 3: import time time.sleep(2 ** attempt) # 指数退避 else: logger.error(f"同步订单 {dingdan.dingdan_id} 到俱乐部 {club.club_id} 最终失败") # 获取下单方ID(根据订单类型) def get_order_user_id(order): from dingdan.models import DingdanPingtai, DingdanShangjia if order.fadan_pingtai == 1: try: ext = DingdanPingtai.objects.get(dingdan=order) return ext.laoban_id # 老板ID except DingdanPingtai.DoesNotExist: return None elif order.fadan_pingtai == 2: try: ext = DingdanShangjia.objects.get(dingdan=order) return ext.shangjia_id # 商家ID except DingdanShangjia.DoesNotExist: return None return None # dingdan/utils/dispatch_stat.py import logging import logging from datetime import date from decimal import Decimal from django.db import transaction from django.db.models import F from peizhi.models import DailyDispatchStat # 请根据实际模型位置修改导入路径 logger = logging.getLogger(__name__) def update_daily_dispatch_stat(direction, partner_club_id, dispatch_amount=None, dispatch_count=1, claimed_amount=None, claimed_count=1, success_amount=None, success_count=1): """ 原子更新每日派单统计(支持派单、抢单、成交三种类型)。 只更新传递了非None参数的字段,未传递的字段保持不变。 参数: direction: 1=我方派单,2=对方派单 partner_club_id: 对方俱乐部ID dispatch_amount: 派单金额(Decimal),如果提供则累加到 dispatch_amount 和 dispatch_count dispatch_count: 派单数量,默认为1,与 dispatch_amount 同时提供时生效 claimed_amount: 抢单金额(Decimal),如果提供则累加到 claimed_amount 和 claimed_count claimed_count: 抢单数量,默认为1,与 claimed_amount 同时提供时生效 success_amount: 成交金额(Decimal),如果提供则累加到 success_amount 和 success_count success_count: 成交数量,默认为1,与 success_amount 同时提供时生效 """ today = date.today() year, month, day = today.year, today.month, today.day with transaction.atomic(): # 使用 select_for_update 锁定行,避免并发累加冲突 stat, created = DailyDispatchStat.objects.select_for_update().get_or_create( date=today, year=year, month=month, day=day, direction=direction, partner_club_id=partner_club_id, defaults={ 'dispatch_count': dispatch_count if dispatch_amount is not None else 0, 'dispatch_amount': dispatch_amount or Decimal('0.00'), 'claimed_count': claimed_count if claimed_amount is not None else 0, 'claimed_amount': claimed_amount or Decimal('0.00'), 'success_count': success_count if success_amount is not None else 0, 'success_amount': success_amount or Decimal('0.00'), } ) if not created: update_fields = [] if dispatch_amount is not None: stat.dispatch_count = F('dispatch_count') + dispatch_count stat.dispatch_amount = F('dispatch_amount') + dispatch_amount update_fields.extend(['dispatch_count', 'dispatch_amount']) if claimed_amount is not None: stat.claimed_count = F('claimed_count') + claimed_count stat.claimed_amount = F('claimed_amount') + claimed_amount update_fields.extend(['claimed_count', 'claimed_amount']) if success_amount is not None: stat.success_count = F('success_count') + success_count stat.success_amount = F('success_amount') + success_amount update_fields.extend(['success_count', 'success_amount']) if update_fields: stat.save(update_fields=update_fields) logger.info(f"更新派单统计: 日期={today}, 方向={direction}, 俱乐部={partner_club_id}, " f"派单+{dispatch_count if dispatch_amount else 0}单/{dispatch_amount or 0}, " f"抢单+{claimed_count if claimed_amount else 0}单/{claimed_amount or 0}, " f"成交+{success_count if success_amount else 0}单/{success_amount or 0}") from peizhi.models import DailyIncomeStat logger = logging.getLogger(__name__) def update_daily_income(amount): """ 原子更新当日收入统计(金额累加,笔数加1)。 用于微信支付成功回调。 参数: amount: Decimal 本次收入金额 """ today = date.today() year, month, day = today.year, today.month, today.day with transaction.atomic(): stat, created = DailyIncomeStat.objects.select_for_update().get_or_create( date=today, defaults={ 'year': year, 'month': month, 'day': day, 'total_amount': amount, 'total_count': 1 } ) if not created: # 原子累加 stat.total_amount = F('total_amount') + amount stat.total_count = F('total_count') + 1 stat.save(update_fields=['total_amount', 'total_count']) logger.info( f"每日收入更新: {today}, +{amount}元, 总金额={stat.total_amount if created else stat.total_amount + amount}, 总笔数={stat.total_count if created else stat.total_count + 1}") from peizhi.models import DailyPayoutStat logger = logging.getLogger(__name__) def update_daily_payout(amount): """ 原子更新当日出款统计(金额累加,笔数加1)。 用于提现成功、结算打款等场景。 参数: amount: Decimal 本次出款金额 """ today = date.today() year, month, day = today.year, today.month, today.day with transaction.atomic(): stat, created = DailyPayoutStat.objects.select_for_update().get_or_create( date=today, defaults={ 'year': year, 'month': month, 'day': day, 'total_amount': amount, 'total_count': 1 } ) if not created: # 原子累加 stat.total_amount = F('total_amount') + amount stat.total_count = F('total_count') + 1 stat.save(update_fields=['total_amount', 'total_count']) logger.info(f"每日出款更新: {today}, +{amount}元, 总金额={stat.total_amount if created else stat.total_amount + amount}, 总笔数={stat.total_count if created else stat.total_count + 1}") from decimal import Decimal from django.db import IntegrityError from dingdan.models import Lilubiao def _normalize_lilu_rate(raw): """Lilubiao 存小数比例;若后台误填 15 表示 15%,自动除以 100""" if raw is None: return Decimal('0') rate = Decimal(str(raw)) if rate <= 0: return Decimal('0') if rate > 1: rate = (rate / Decimal('100')).quantize(Decimal('0.0001')) return rate def calc_shangjia_order_fencheng(jine): """ 商家发单时计算打手/管事分成。 - 打手分成(Lilubiao=3) - 管事打手接单分红(Lilubiao=13) - 管事商家派单分红(Lilubiao=15) 三项合计超过订单金额时,优先削减打手分成,再削减商家派单管事分红,最后削减打手接单管事分红。 """ jine = Decimal(str(jine)) try: lilu_obj = Lilubiao.objects.get(fadanpingtai='3') rate_dashou = _normalize_lilu_rate(lilu_obj.lilu) if lilu_obj.lilu else Decimal('1') if rate_dashou <= 0: rate_dashou = Decimal('1') except Lilubiao.DoesNotExist: rate_dashou = Decimal('1') lilu_guanshi_dashou_obj = Lilubiao.objects.filter(fadanpingtai='13').first() rate_guanshi_dashou = _normalize_lilu_rate( lilu_guanshi_dashou_obj.lilu if lilu_guanshi_dashou_obj else None ) lilu_guanshi_shangjia_obj = Lilubiao.objects.filter(fadanpingtai='15').first() rate_guanshi_shangjia = _normalize_lilu_rate( lilu_guanshi_shangjia_obj.lilu if lilu_guanshi_shangjia_obj else None ) dashou_fencheng = min( (jine * rate_dashou).quantize(Decimal('0.01')), jine, ) guanshi_fencheng = (jine * rate_guanshi_dashou).quantize(Decimal('0.01')) guanshi_shangjia_fencheng = (jine * rate_guanshi_shangjia).quantize(Decimal('0.01')) total = dashou_fencheng + guanshi_fencheng + guanshi_shangjia_fencheng if total > jine: overflow = total - jine cut = min(overflow, dashou_fencheng) dashou_fencheng -= cut overflow -= cut if overflow > 0: cut = min(overflow, guanshi_shangjia_fencheng) guanshi_shangjia_fencheng -= cut overflow -= cut if overflow > 0: cut = min(overflow, guanshi_fencheng) guanshi_fencheng -= cut return dashou_fencheng, guanshi_fencheng, guanshi_shangjia_fencheng def _pay_guanshi_order_fenhong(order, guanshi_id, amount, related_user_id, nicheng, avatar, fenhong_leixing, shuoming): """向管事发放单笔商家订单分红(幂等)""" from shangpin.models import Gsfenhong from yonghu.models import UserGuanshi amount = amount or Decimal('0') if amount <= 0 or not guanshi_id: return False if Gsfenhong.objects.filter(dingdan_id=order.dingdan_id, fenhong_leixing=fenhong_leixing).exists(): logger.info(f"管事分红已处理: 订单{order.dingdan_id}, 类型{fenhong_leixing}") return False try: with transaction.atomic(): guanshi = UserGuanshi.objects.select_for_update().get(user__yonghuid=guanshi_id) UserGuanshi.objects.filter(id=guanshi.id).update( yue=F('yue') + amount, chongzhifenrun=F('chongzhifenrun') + amount, ) Gsfenhong.objects.create( dingdan_id=order.dingdan_id, guanshi=guanshi_id, dashouid=related_user_id, shuoming=shuoming, fenhong=amount, avatar=avatar, nicheng=nicheng or '未知用户', fenhong_leixing=fenhong_leixing, ) logger.info( f"管事分红成功: 订单{order.dingdan_id}, 管事{guanshi_id}, " f"关联用户{related_user_id}, 类型{fenhong_leixing}, 金额{amount}元" ) except UserGuanshi.DoesNotExist: logger.warning(f"管事分红跳过: 管事{guanshi_id}不存在") return False except IntegrityError as e: logger.error( f"管事分红写入失败(数据库约束): 订单{order.dingdan_id}, 类型{fenhong_leixing}, " f"关联用户{related_user_id}, 错误={e}. " f"请确认已执行 migrate shangpin 0010 与 migrate dingdan 0017", exc_info=True, ) return False except Exception as e: logger.error(f"管事分红失败: {e}", exc_info=True) return False try: from houtai.utils import update_guanshi_daily_by_action update_guanshi_daily_by_action( yonghuid=guanshi_id, action=3, amount=amount, ) except Exception as e: logger.error(f"管事每日统计更新失败: {e}") return True def _resolve_dashou_inviter_guanshi(dashou_profile): guanshi_id = (dashou_profile.yaoqingren or '').strip() return guanshi_id or None def settle_shangjia_order_guanshi_fenhong(order, dashou_id): """ 商家订单结单时结算管事分红(幂等,失败不抛异常阻断主流程)。 - 管事打手接单分红:接单打手的邀请管事(guanshi_fencheng) - 管事商家派单分红:派单商家打手身份的邀请管事(guanshi_shangjia_fencheng) 同一管事可同时获得两笔分红。 """ from dingdan.models import Dingdan, DingdanShangjia from yonghu.models import UserDashou if getattr(order, 'fadan_pingtai', None) != 2: return # 重新读订单,确保拿到发单时写入的分红金额 try: order = Dingdan.objects.get(dingdan_id=order.dingdan_id) except Dingdan.DoesNotExist: logger.warning(f"管事分红跳过: 订单{order.dingdan_id}不存在") return guanshi_fencheng = order.guanshi_fencheng or Decimal('0') guanshi_shangjia_fencheng = order.guanshi_shangjia_fencheng or Decimal('0') logger.info( f"商家订单管事分红结算: 订单{order.dingdan_id}, 打手{dashou_id}, " f"打手接单分红={guanshi_fencheng}, 商家派单分红={guanshi_shangjia_fencheng}" ) if guanshi_fencheng <= 0 and guanshi_shangjia_fencheng <= 0: logger.info( f"管事分红跳过: 订单{order.dingdan_id} 两项分红均为0 " f"(请检查 Lilubiao 13/15 费率及打手分成是否占满订单金额)" ) return # 1. 接单打手邀请管事分红 if guanshi_fencheng > 0 and dashou_id: try: dashou = UserDashou.objects.select_related('user').get(user__yonghuid=dashou_id) except UserDashou.DoesNotExist: logger.info(f"管事打手分红跳过: 打手{dashou_id}不存在") dashou = None guanshi_id = _resolve_dashou_inviter_guanshi(dashou) if dashou else None if dashou and guanshi_id: user_main = dashou.user _pay_guanshi_order_fenhong( order=order, guanshi_id=guanshi_id, amount=guanshi_fencheng, related_user_id=dashou_id, nicheng=dashou.nicheng, avatar=user_main.avatar if user_main else None, fenhong_leixing=3, shuoming='商家订单打手分红', ) else: logger.info(f"管事打手分红跳过: 打手{dashou_id}无邀请管事(yaoqingren为空)") # 2. 派单商家邀请管事分红 if guanshi_shangjia_fencheng > 0: try: shangjia_ext = DingdanShangjia.objects.select_related('dingdan').get( dingdan__dingdan_id=order.dingdan_id ) shangjia_id = shangjia_ext.shangjia_id except DingdanShangjia.DoesNotExist: logger.info(f"管事商家分红跳过: 订单{order.dingdan_id}无商家扩展") shangjia_id = None if shangjia_id: try: merchant_dashou = UserDashou.objects.select_related('user').get(user__yonghuid=shangjia_id) except UserDashou.DoesNotExist: logger.info(f"管事商家分红跳过: 商家{shangjia_id}无打手扩展") merchant_dashou = None guanshi_id = _resolve_dashou_inviter_guanshi(merchant_dashou) if merchant_dashou else None if merchant_dashou and guanshi_id: user_main = merchant_dashou.user _pay_guanshi_order_fenhong( order=order, guanshi_id=guanshi_id, amount=guanshi_shangjia_fencheng, related_user_id=shangjia_id, nicheng=shangjia_ext.sjnicheng or merchant_dashou.nicheng, avatar=user_main.avatar if user_main else None, fenhong_leixing=4, shuoming='商家派单管事分红', ) else: logger.info( f"管事商家分红跳过: 商家{shangjia_id}打手身份无邀请管事(yaoqingren为空)" ) def pay_deposit_guanshi_fenhong(dashou, order, recharge_amount): """ 打手充值押金后向邀请管事发放分红(幂等:dingdan_id + fenhong_leixing=2)。 网络重试/微信重复回调时不会重复入账。 """ from shangpin.models import Gsfenhong from yonghu.models import UserGuanshi guanshi_id = (dashou.yaoqingren or '').strip() if not guanshi_id: logger.info(f"押金分红跳过: 打手{dashou.user.yonghuid}无邀请管事") return False if Gsfenhong.objects.filter(dingdan_id=order.dingdan_id, fenhong_leixing=2).exists(): logger.info(f"押金分红已处理: 订单{order.dingdan_id}") return True lilu_obj = Lilubiao.objects.filter(fadanpingtai='12').first() rate = lilu_obj.lilu if lilu_obj and lilu_obj.lilu is not None else Decimal('0') rate = _normalize_lilu_rate(rate) if rate else Decimal('0') if rate <= 0: logger.info("押金分红跳过: 费率未配置或为0") return False recharge_amount = Decimal(str(recharge_amount)) fenhong_jine = (recharge_amount * rate).quantize(Decimal('0.01')) if fenhong_jine <= 0: return False try: with transaction.atomic(): guanshi = UserGuanshi.objects.select_for_update().get(user__yonghuid=guanshi_id) UserGuanshi.objects.filter(pk=guanshi.pk).update( yue=F('yue') + fenhong_jine, chongzhifenrun=F('chongzhifenrun') + fenhong_jine, ) user_main = dashou.user Gsfenhong.objects.create( dingdan_id=order.dingdan_id, guanshi=guanshi_id, dashouid=dashou.user.yonghuid, shuoming='押金分红', fenhong=fenhong_jine, avatar=user_main.avatar if user_main else None, nicheng=dashou.nicheng or '未知打手', fenhong_leixing=2, ) logger.info( f"押金分红成功: 订单{order.dingdan_id}, 管事{guanshi_id}, " f"打手{dashou.user.yonghuid}, 金额{fenhong_jine}元" ) try: from houtai.utils import update_guanshi_daily_by_action update_guanshi_daily_by_action( yonghuid=guanshi_id, action=4, amount=fenhong_jine, ) except Exception as e: logger.error(f"押金分红管事日统计更新失败: {e}") return True except UserGuanshi.DoesNotExist: logger.warning(f"押金分红跳过: 管事{guanshi_id}不存在") return False except IntegrityError: logger.info(f"押金分红已处理(并发): 订单{order.dingdan_id}") return True except Exception as e: logger.error(f"押金分红处理失败: {e}", exc_info=True) return False