Files
along_django/dingdan/utils.py

511 lines
20 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 dingdan.models import Lilubiao
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 = Decimal(str(lilu_obj.lilu)) if lilu_obj.lilu and lilu_obj.lilu > 0 else Decimal('1')
except Lilubiao.DoesNotExist:
rate_dashou = Decimal('1')
lilu_guanshi_dashou_obj = Lilubiao.objects.filter(fadanpingtai='13').first()
rate_guanshi_dashou = (
Decimal(str(lilu_guanshi_dashou_obj.lilu))
if lilu_guanshi_dashou_obj and lilu_guanshi_dashou_obj.lilu is not None
else Decimal('0')
)
lilu_guanshi_shangjia_obj = Lilubiao.objects.filter(fadanpingtai='15').first()
rate_guanshi_shangjia = (
Decimal(str(lilu_guanshi_shangjia_obj.lilu))
if lilu_guanshi_shangjia_obj and lilu_guanshi_shangjia_obj.lilu is not None
else Decimal('0')
)
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'))
if dashou_fencheng + guanshi_fencheng > jine:
guanshi_fencheng = Decimal('0.00')
if dashou_fencheng + guanshi_fencheng + guanshi_shangjia_fencheng > jine:
guanshi_shangjia_fencheng = Decimal('0.00')
if dashou_fencheng + guanshi_fencheng + guanshi_shangjia_fencheng > jine:
guanshi_fencheng = Decimal('0.00')
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 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 settle_shangjia_order_guanshi_fenhong(order, dashou_id):
"""
商家订单结单时结算管事分红(幂等,失败不抛异常阻断主流程)。
- 管事打手接单分红接单打手的邀请管事guanshi_fencheng
- 管事商家派单分红派单商家打手身份的邀请管事guanshi_shangjia_fencheng
同一管事可同时获得两笔分红。
"""
from dingdan.models import DingdanShangjia
from yonghu.models import UserDashou
if getattr(order, 'fadan_pingtai', None) != 2:
return
# 1. 接单打手邀请管事分红
guanshi_fencheng = order.guanshi_fencheng or Decimal('0')
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
if dashou and dashou.yaoqingren:
user_main = dashou.user
_pay_guanshi_order_fenhong(
order=order,
guanshi_id=dashou.yaoqingren,
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}无邀请管事")
# 2. 派单商家邀请管事分红
guanshi_shangjia_fencheng = order.guanshi_shangjia_fencheng or Decimal('0')
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
if merchant_dashou and merchant_dashou.yaoqingren:
user_main = merchant_dashou.user
_pay_guanshi_order_fenhong(
order=order,
guanshi_id=merchant_dashou.yaoqingren,
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}打手身份无邀请管事")