# -*- coding: utf-8 -*- """库存变更任务""" import json from .base_task import BaseTask from loaders.facts.inventory_change import InventoryChangeLoader from models.parsers import TypeParser class InventoryChangeTask(BaseTask): """同步库存变化记录""" def get_task_code(self) -> str: return "INVENTORY_CHANGE" def execute(self) -> dict: self.logger.info("开始执行 INVENTORY_CHANGE 任务") window_start, window_end, _ = self._get_time_window() params = { "storeId": self.config.get("app.store_id"), "startTime": TypeParser.format_timestamp(window_start, self.tz), "endTime": TypeParser.format_timestamp(window_end, self.tz), } try: records, _ = self.api.get_paginated( endpoint="/Inventory/ChangeList", params=params, page_size=self.config.get("api.page_size", 200), data_path=("data", "queryDeliveryRecordsList"), ) parsed = [] for raw in records: mapped = self._parse_change(raw) if mapped: parsed.append(mapped) loader = InventoryChangeLoader(self.db) inserted, updated, skipped = loader.upsert_changes(parsed) self.db.commit() counts = { "fetched": len(records), "inserted": inserted, "updated": updated, "skipped": skipped, "errors": 0, } self.logger.info(f"INVENTORY_CHANGE 完成: {counts}") return self._build_result("SUCCESS", counts) except Exception: self.db.rollback() self.logger.error("INVENTORY_CHANGE 失败", exc_info=True) raise def _parse_change(self, raw: dict) -> dict | None: change_id = TypeParser.parse_int( raw.get("siteGoodsStockId") or raw.get("site_goods_stock_id") ) if not change_id: self.logger.warning("跳过缺少变动 id 的库存记录: %s", raw) return None store_id = self.config.get("app.store_id") return { "store_id": store_id, "change_id": change_id, "site_goods_id": TypeParser.parse_int( raw.get("siteGoodsId") or raw.get("site_goods_id") ), "stock_type": raw.get("stockType") or raw.get("stock_type"), "goods_name": raw.get("goodsName"), "change_time": TypeParser.parse_timestamp( raw.get("createTime") or raw.get("create_time"), self.tz ), "start_qty": TypeParser.parse_int(raw.get("startNum")), "end_qty": TypeParser.parse_int(raw.get("endNum")), "change_qty": TypeParser.parse_int(raw.get("changeNum")), "unit": raw.get("unit"), "price": TypeParser.parse_decimal(raw.get("price")), "operator_name": raw.get("operatorName"), "remark": raw.get("remark"), "goods_category_id": TypeParser.parse_int(raw.get("goodsCategoryId")), "goods_second_category_id": TypeParser.parse_int( raw.get("goodsSecondCategoryId") ), "raw_data": json.dumps(raw, ensure_ascii=False), }