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_hour 和 run_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.py 的 CELERY_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.py(Celery 任务)
  7. 实现 apps.py(注册信号处理器)
  8. 在 settings.py 中添加 settlement 到 INSTALLED_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 中只做日志输出,便于调试