1
0
forked from erp-dev/erp
Files
erpnew/docs/settlement/DESIGN.md

15 KiB
Raw Blame History

日结服务设计方案

一、概述

日结服务是一个基于 Celery 的定时任务系统用于在每天凌晨5点自动执行各商户的日结统计。服务采用独立模块设计支持商户级别的配置包括统计模块选择和通知渠道配置。

二、设计原则

  1. 独立模块settlement 模块独立,职责清晰
  2. 配置驱动:每个商户可独立配置统计模块和通知渠道
  3. 统一调度:运行时间在 settings.py 中统一配置,便于管理
  4. 商户隔离:每个商户的日结任务独立执行,互不干扰
  5. 信号通知:通过 signal 通知日结完成,便于后续扩展
  6. 模块化统计:统计函数集中在 services.py,便于后续填充具体逻辑
  7. 容错性强:单个模块失败不影响其他模块
  8. 可扩展性:未来可轻松添加新的统计模块和通知渠道

三、模块结构

settlement/
├── __init__.py
├── apps.py              # 应用配置,注册信号处理器
├── models.py            # 日结配置模型
├── signals.py           # 日结完成信号定义
├── handlers.py          # 信号处理器(日志输出)
├── services.py          # 统计服务函数(空壳)
├── tasks.py            # Celery定时任务
└── migrations/         # 数据库迁移文件

四、模型设计

4.1 日结配置模型

class NotificationChannelEnum(models.TextChoices):
    """通知渠道枚举"""
    WECOM = 'wecom', '企业微信'
    NONE = 'none', '不通知'


class SettlementModuleEnum(models.TextChoices):
    """日结模块枚举"""
    PLATE_ORDER = 'plate_order', '开版订单'
    PRINTING_ORDER = 'printing_order', '生产订单'


class DailySettlementConfig(ModelBase):
    """日结配置 - 每个商户一条记录(可选)"""
    id = models.BigAutoField(primary_key=True)
    merchant = models.OneToOneField(
        'basic_info.Merchant',
        on_delete=models.CASCADE,
        related_name='settlement_config',
        verbose_name='所属商户'
    )
    
    # 统计模块配置JSON数组存储枚举值
    settlement_modules = models.JSONField(
        default=list,
        verbose_name='统计模块',
        help_text='例如: ["plate_order", "printing_order"]'
    )
    
    # 通知配置
    notification_enabled = models.BooleanField(default=False, verbose_name='启用通知')
    notification_channel = models.CharField(
        max_length=20,
        choices=NotificationChannelEnum.choices,
        default=NotificationChannelEnum.NONE,
        verbose_name='通知渠道'
    )
    
    # 其他配置
    description = models.TextField(blank=True, null=True, verbose_name='备注')
    
    class Meta:
        verbose_name = '日结配置'
        verbose_name_plural = '日结配置'
        db_table = 'daily_settlement_config'

说明

  • 新商户不创建配置,则不参与日结统计
  • 配置中不包含 run_hourrun_minute,统一在 settings.py 中配置

五、信号设计

5.1 信号定义

from django.dispatch import Signal

# 日结完成信号
# 参数:
#   merchant: Merchant 实例
#   settlement_date: datetime.date
#   status: 'success' / 'failed' / 'partial'
#   modules: list[str] - 参与统计的模块列表
#   errors: dict - 模块错误详情 {'module_name': 'error message'}
#   task_id: str - Celery 任务ID
daily_settlement_completed = Signal()

5.2 信号处理器

def on_daily_settlement_completed(sender, **kwargs):
    """日结完成时的处理(日志输出)"""
    merchant = kwargs.get('merchant')
    settlement_date = kwargs.get('settlement_date')
    status = kwargs.get('status')
    modules = kwargs.get('modules', [])
    errors = kwargs.get('errors', {})
    task_id = kwargs.get('task_id')
    
    # 检查配置是否启用通知
    config = merchant.settlement_config
    if not config or not config.notification_enabled:
        logger.info(
            f'[settlement.handlers] 商户 {merchant.id} 日结完成,通知未启用'
        )
        return
    
    # 根据通知渠道处理(目前只有企业微信)
    if config.notification_channel == NotificationChannelEnum.WECOM:
        logger.info(
            f'[settlement.handlers] 商户 {merchant.id} 日结完成,'
            f'状态={status}, 模块={modules}, 错误={errors}, '
            f'通知渠道=企业微信(暂不发送具体内容)'
        )
        # 未来可以在这里添加具体的企业微信通知逻辑

六、统计服务函数(空壳)

6.1 开版订单统计

def calculate_plate_order_daily_summary(
    merchant_id: int,
    settlement_date: datetime.date
) -> dict:
    """
    计算开版订单日结汇总(空壳实现)
    
    Args:
        merchant_id: 商户ID
        settlement_date: 结算日期
    
    Returns:
        dict: 统计结果(空壳)
    """
    logger.info(
        f'[settlement.services] 计算开版订单日结汇总: '
        f'merchant_id={merchant_id}, date={settlement_date}'
    )
    
    return {
        'merchant_id': merchant_id,
        'settlement_date': str(settlement_date),
        'total_orders': 0,
        'completed_orders': 0,
        'in_progress_orders': 0,
    }

6.2 生产订单统计

def calculate_printing_order_daily_summary(
    merchant_id: int,
    settlement_date: datetime.date
) -> dict:
    """
    计算生产订单日结汇总(空壳实现)
    
    Args:
        merchant_id: 商户ID
        settlement_date: 结算日期
    
    Returns:
        dict: 统计结果(空壳)
    """
    logger.info(
        f'[settlement.services] 计算生产订单日结汇总: '
        f'merchant_id={merchant_id}, date={settlement_date}'
    )
    
    return {
        'merchant_id': merchant_id,
        'settlement_date': str(settlement_date),
        'total_orders': 0,
        'completed_orders': 0,
        'in_progress_orders': 0,
    }

七、Celery 任务设计

7.1 主调度任务

@shared_task(bind=True)
def run_daily_settlement(self):
    """
    执行所有商户的日结任务(固定时间触发)
    
    在 settings.py 中配置的固定时间触发,
    拉取所有配置了统计模块的商户,逐个执行日结。
    """
    from basic_info.models import Merchant
    
    # 获取所有配置了统计模块的商户
    configs = DailySettlementConfig.objects.exclude(
        settlement_modules=[]
    ).select_related('merchant')
    
    logger.info(
        f'[settlement.tasks] 开始执行日结,共 {configs.count()} 个商户需要统计'
    )
    
    # 为每个商户触发日结任务
    for config in configs:
        settlement_date = timezone.localdate() - datetime.timedelta(days=1)
        run_merchant_daily_settlement.delay(
            merchant_id=config.merchant.id,
            settlement_date=str(settlement_date)
        )
    
    return {
        'task_id': self.request.id,
        'total_merchants': configs.count(),
        'settlement_date': str(timezone.localdate() - datetime.timedelta(days=1)),
    }

7.2 商户日结任务

@shared_task(bind=True)
def run_merchant_daily_settlement(
    self,
    merchant_id: int,
    settlement_date: str | None = None
):
    """
    执行单个商户的日结
    
    Args:
        merchant_id: 商户ID
        settlement_date: 结算日期YYYY-MM-DD默认为昨天
    
    Returns:
        dict: 包含 task_id 和执行结果的 payload
    """
    from basic_info.models import Merchant
    
    # 解析日期
    if settlement_date:
        date_obj = datetime.datetime.strptime(settlement_date, '%Y-%m-%d').date()
    else:
        date_obj = timezone.localdate() - datetime.timedelta(days=1)
    
    merchant = Merchant.objects.get(id=merchant_id)
    config = merchant.settlement_config
    
    modules = config.settlement_modules
    summaries = {}
    errors = {}
    overall_status = 'success'
    
    logger.info(
        f'[settlement.tasks] 开始执行商户 {merchant_id} 的日结,'
        f'日期={date_obj}, 模块={modules}'
    )
    
    # 执行各模块的统计
    for module in modules:
        try:
            if module == SettlementModuleEnum.PLATE_ORDER:
                summary = calculate_plate_order_daily_summary(
                    merchant_id=merchant_id,
                    settlement_date=date_obj
                )
                summaries['plate_order'] = summary
                
            elif module == SettlementModuleEnum.PRINTING_ORDER:
                summary = calculate_printing_order_daily_summary(
                    merchant_id=merchant_id,
                    settlement_date=date_obj
                )
                summaries['printing_order'] = summary
                
        except Exception as e:
            logger.exception(
                f'[settlement.tasks] 商户 {merchant_id}{module} 统计失败: {e}'
            )
            errors[module] = str(e)
            overall_status = 'partial'
    
    # 如果所有模块都失败
    if len(errors) == len(modules):
        overall_status = 'failed'
    
    # 发送信号(不包含统计结果详情)
    daily_settlement_completed.send(
        sender=run_merchant_daily_settlement,
        merchant=merchant,
        settlement_date=date_obj,
        status=overall_status,
        modules=modules,
        errors=errors,
        task_id=self.request.id
    )
    
    payload = {
        'task_id': self.request.id,
        'merchant_id': merchant_id,
        'settlement_date': str(date_obj),
        'status': overall_status,
        'modules': modules,
        'errors': errors,
    }
    
    logger.info(
        f'[settlement.tasks] 商户 {merchant_id} 日结完成: {payload}'
    )
    return payload

八、Celery Beat 配置

settings.pyCELERY_BEAT_SCHEDULE 中添加:

# 日结任务 - 每天凌晨5点执行
'daily_settlement': {
    'task': 'settlement.tasks.run_daily_settlement',
    'schedule': crontab(hour=5, minute=0),
},

说明

  • 运行时间统一在 settings.py 中配置(hour=5, minute=0
  • 不需要每小时检查,直接在固定时间触发

九、settings.py 配置

# 日结服务配置
# ------------------------------------------------------------------------------
# 日结运行时间(固定时间,所有商户统一)
SETTLEMENT_RUN_HOUR = env.int('SETTLEMENT_RUN_HOUR', default=5)
SETTLEMENT_RUN_MINUTE = env.int('SETTLEMENT_RUN_MINUTE', default=0)
# ------------------------------------------------------------------------------

然后在 CELERY_BEAT_SCHEDULE 中使用:

from celery.schedules import crontab

'daily_settlement': {
    'task': 'settlement.tasks.run_daily_settlement',
    'schedule': crontab(
        hour=settings.SETTLEMENT_RUN_HOUR,
        minute=settings.SETTLEMENT_RUN_MINUTE
    ),
},

十、应用配置

class SettlementConfig(AppConfig):
    default_auto_field = 'django.db.models.BigAutoField'
    name = 'settlement'
    verbose_name = '日结管理'

    def ready(self):
        """注册信号处理函数"""
        from .signals import daily_settlement_completed
        from . import handlers
        
        # 注册日结完成信号处理器
        daily_settlement_completed.connect(
            handlers.on_daily_settlement_completed,
            dispatch_uid='settlement.on_daily_settlement_completed'
        )
        
        logger.info(
            f'[settlement.apps] 已注册 daily_settlement_completed 信号处理器'
        )

十一、错误处理策略

  1. 模块级隔离:每个模块的统计失败不影响其他模块
  2. 错误记录:记录每个模块的错误信息到 errors 字典
  3. 状态判断
    • success:所有模块都成功
    • partial:部分模块成功
    • failed:所有模块都失败
  4. 日志记录:使用 logger 记录详细日志(开始、完成、异常)

十二、并发控制说明

暂不考虑并发控制,原因:

  1. 日结任务在凌晨时段触发,业务量较小
  2. 商户数量预期不多,可以逐个串行执行
  3. 未来如需提升性能,可以通过增加 Celery worker 数量来解决
  4. Celery 本身支持任务队列和并发控制,后续可以轻松扩展

十三、数据库优化评估PostgreSQL

13.1 当前结论(暂不实施)

已评估从数据库层面优化开版订单统计(索引、视图、物化视图),当前阶段暂不实施结构性优化,维持现有实现。

13.2 暂缓原因

  1. 当前测试与功能已稳定,优先保证行为一致性
  2. 生产环境要求“不可接受阻塞风险”,索引变更需专项窗口与监控保障
  3. 现阶段数据规模下,统计查询尚可接受

13.3 后续可选优化路线(按优先级)

  1. 索引优先:为 plate_order 统计路径增加复合/部分索引
  2. 表达式索引:针对 date(plate_date) 的筛选场景
  3. 物化视图:当数据规模明显增大时,将日粒度聚合前置

13.4 生产安全约束

若后续执行索引优化,必须遵循:

  1. 使用 PostgreSQL CREATE INDEX CONCURRENTLY
  2. 迁移使用 atomic = False
  3. 低峰分批执行,一次一个索引
  4. 全程监控 CPU/IO/WAL、慢查询与复制延迟

13.5 验证要求

任何数据库优化上线前,需要在预发(接近生产数据量)完成:

  1. EXPLAIN ANALYZE 对比
  2. 回归测试通过(settlement + api_v1.views.settlement
  3. 回滚脚本预演

十四、实施步骤

  1. 创建 settlement 模块目录结构
  2. 实现 models.py(配置模型)
  3. 实现 signals.py(信号定义)
  4. 实现 services.py(空壳统计函数)
  5. 实现 handlers.py(信号处理器)
  6. 实现 tasks.pyCelery 任务)
  7. 实现 apps.py(注册信号处理器)
  8. settings.py 中添加 settlementINSTALLED_APPS
  9. settings.py 中配置 Celery Beat 和日结运行时间
  10. 创建数据库迁移
  11. 编写测试用例
  12. 手动触发测试,验证结果

十五、方案优势

  1. 独立模块settlement 模块独立,职责清晰
  2. 配置驱动:每个商户可独立配置统计模块和通知渠道
  3. 统一调度:运行时间在 settings.py 中统一配置,便于管理
  4. 商户隔离:每个商户的日结任务独立执行,互不干扰
  5. 信号通知:通过 signal 通知日结完成,便于后续扩展
  6. 模块化统计:统计函数集中在 services.py,便于后续填充具体逻辑
  7. 容错性强:单个模块失败不影响其他模块
  8. 可扩展性:未来可轻松添加新的统计模块和通知渠道
  9. 简化信号:信号参数只包含基本信息,不包含统计结果详情
  10. 日志友好handler 中只做日志输出,便于调试