# 日结服务设计方案 ## 一、概述 日结服务是一个基于 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 日结配置模型 ```python 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 = '日结配置' ``` **说明**: - 新商户不创建配置,则不参与日结统计 - 配置中不包含 `run_hour` 和 `run_minute`,统一在 settings.py 中配置 ## 五、信号设计 ### 5.1 信号定义 ```python 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 信号处理器 ```python 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 开版订单统计 ```python 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 生产订单统计 ```python 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 主调度任务 ```python @shared_task(bind=True) def run_daily_settlement(self): """ 执行所有商户的日结任务(固定时间触发) 在 settings.py 中配置的固定时间触发, 拉取所有配置了统计模块的商户,逐个执行日结。 """ from basic_info.models import Merchant # 获取所有配置了统计模块的商户 configs = DailySettlementConfig.objects.filter( settlement_modules__len__gt=0 ).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 商户日结任务 ```python @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` 中添加: ```python # 日结任务 - 每天凌晨5点执行 'daily_settlement': { 'task': 'settlement.tasks.run_daily_settlement', 'schedule': crontab(hour=5, minute=0), }, ``` **说明**: - 运行时间统一在 settings.py 中配置(`hour=5, minute=0`) - 不需要每小时检查,直接在固定时间触发 ## 九、settings.py 配置 ```python # 日结服务配置 # ------------------------------------------------------------------------------ # 日结运行时间(固定时间,所有商户统一) SETTLEMENT_RUN_HOUR = env.int('SETTLEMENT_RUN_HOUR', default=5) SETTLEMENT_RUN_MINUTE = env.int('SETTLEMENT_RUN_MINUTE', default=0) # ------------------------------------------------------------------------------ ``` 然后在 CELERY_BEAT_SCHEDULE 中使用: ```python 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 ), }, ``` ## 十、应用配置 ```python 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 本身支持任务队列和并发控制,后续可以轻松扩展 ## 十三、实施步骤 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 中只做日志输出,便于调试