# 文件路径: app/services/inbound/base_service.py from flask_jwt_extended import get_jwt from app.extensions import db from app.models.base import MaterialBase, MaterialWarningSetting from app.models.inbound.buy import StockBuy from app.models.inbound.semi import StockSemi from app.models.inbound.product import StockProduct # from app.models.inbound.service import StockService from app.models.bom import BomTable from sqlalchemy import or_, and_, func, case, desc, asc, cast, Numeric, text import traceback import json import io import datetime import numpy as np from app.utils.ai_vision import extract_and_embed from app.services.image_embedding_service import ImageEmbeddingService # 需要 pip install openpyxl from openpyxl import Workbook from openpyxl.styles import Font, Alignment, Border, Side, PatternFill from collections import defaultdict class MaterialBaseService: """ 基础物料服务层 负责处理 MaterialBase 的增删改查及搜索逻辑 """ @staticmethod def search_material(keyword): """ 根据关键字搜索已启用的基础物料 (供 /api/v1/inbound/base/search 接口调用) """ try: if not keyword: return [] query = MaterialBase.query.filter( MaterialBase.is_enabled == True, or_( MaterialBase.name.ilike(f'%{keyword}%'), MaterialBase.common_name.ilike(f'%{keyword}%'), MaterialBase.spec_model.ilike(f'%{keyword}%'), # 支持搜索公司名 MaterialBase.company_name.ilike(f'%{keyword}%') ) ) # 【行级数据隔离】基于 JWT 多租户公司过滤 from app.utils.decorators import get_current_company_filter company_limit = get_current_company_filter() if company_limit is not None: query = query.filter(MaterialBase.company_name == company_limit) # [修改1] 增加返回数量限制 # 原为 limit(20),现改为 1000,确保前端能获取所有(或足够多)的数据 query = query.limit(1000) # 获取查询结果对象列表 db_items = query.all() # [修改2] 规格型号排序逻辑 # 要求:只考虑 '/' 前面的内容进行排序 # 使用 Python 的 sort 方法,提取 spec_model 中 '/' 前的部分 def get_sort_key(item): if not item.spec_model: return "" # 如果包含 '/',取前半部分;否则取整个字符串 parts = item.spec_model.split('/') return parts[0] if len(parts) > 0 else item.spec_model # 执行排序 db_items.sort(key=get_sort_key) results = [] for item in db_items: results.append({ 'id': item.id, # 必须保留ID供前端逻辑使用,视觉上的隐藏请在前端处理 'companyName': item.company_name, 'name': item.name, 'commonName': item.common_name, 'spec': item.spec_model, 'category': item.category, 'unit': item.unit, 'type': item.material_type, 'status': '启用', # 强制质检标记 'isInspectionRequired': bool(item.is_inspection_required) }) return results except Exception as e: traceback.print_exc() return [] @staticmethod def _get_stock_counts(stock_query): """ 辅助函数:安全计算库存列表的总数量 """ total_inv = 0 total_avail = 0 try: items = list(stock_query) # 触发查询 except: items = [] for x in items: # 1. 获取库存数 (兼容不同字段名) q = getattr(x, 'stock_quantity', getattr(x, 'in_quantity', 0)) # 优先取库存,其次入库 # 2. 获取可用数 a = getattr(x, 'available_quantity', q) try: total_inv += float(q if q is not None else 0) total_avail += float(a if a is not None else 0) except: pass return total_inv, total_avail @staticmethod def get_list(page, limit, filters=None, user_permissions=None): """ 获取基础信息列表 (带分页、高级筛选和全字段排序) 支持库存预警功能(如果用户有 view_warning 权限) 修复说明: 1. 使用子查询方案确保聚合列(total_inv)能正确参与预警排序计算 2. 高级筛选中针对聚合字段使用 HAVING 而非 WHERE """ try: # 检查用户是否有查看预警的权限 has_warning_permission = False if user_permissions: # 超级管理员有所有权限 if '*' in user_permissions or 'material_list:*' in user_permissions: has_warning_permission = True elif 'material_list:view_warning' in user_permissions: has_warning_permission = True # 构建聚合子查询 buy_sub = db.session.query( StockBuy.base_id, func.sum(StockBuy.stock_quantity).label('buy_inv'), func.sum(StockBuy.available_quantity).label('buy_avail') ).group_by(StockBuy.base_id).subquery() semi_sub = db.session.query( StockSemi.base_id, func.sum(StockSemi.stock_quantity).label('semi_inv'), func.sum(StockSemi.available_quantity).label('semi_avail') ).group_by(StockSemi.base_id).subquery() prod_sub = db.session.query( StockProduct.base_id, func.sum(StockProduct.stock_quantity).label('prod_inv'), func.sum(StockProduct.available_quantity).label('prod_avail') ).group_by(StockProduct.base_id).subquery() # 总库存和可用数的 SQL 表达式(用于后续计算) total_inv = func.coalesce(buy_sub.c.buy_inv, 0) + \ func.coalesce(semi_sub.c.semi_inv, 0) + \ func.coalesce(prod_sub.c.prod_inv, 0) total_avail = func.coalesce(buy_sub.c.buy_avail, 0) + \ func.coalesce(semi_sub.c.semi_avail, 0) + \ func.coalesce(prod_sub.c.prod_avail, 0) # 导入预警设置模型 from app.models.base import MaterialWarningSetting # ============================================================ # 【核心修复】使用子查询方案:将聚合结果作为子查询 # 这样 total_inv 在外层查询中是一个普通列,可安全参与计算 # ============================================================ # 内层子查询:基础数据 + 聚合库存 inner_sub = db.session.query( MaterialBase.id.label('base_id'), MaterialBase, total_inv.label('total_inv'), total_avail.label('total_avail') ).outerjoin(buy_sub, MaterialBase.id == buy_sub.c.base_id) \ .outerjoin(semi_sub, MaterialBase.id == semi_sub.c.base_id) \ .outerjoin(prod_sub, MaterialBase.id == prod_sub.c.base_id) # 将内层子查询具体化 inner_sub = inner_sub.subquery() # 外层查询:关联预警设置,并在外层处理排序和高级筛选 query = db.session.query( MaterialBase, inner_sub.c.total_inv, inner_sub.c.total_avail, MaterialWarningSetting.is_enabled.label('warning_enabled'), MaterialWarningSetting.yellow_threshold.label('warning_yellow'), MaterialWarningSetting.red_threshold.label('warning_red'), MaterialWarningSetting.red_emails.label('warning_red_emails'), MaterialWarningSetting.yellow_emails.label('warning_yellow_emails') ).outerjoin(inner_sub, MaterialBase.id == inner_sub.c.base_id) \ .outerjoin(MaterialWarningSetting, MaterialBase.id == MaterialWarningSetting.base_id) # 【修复】关键词搜索必须在主查询上过滤,不能在 inner_sub 上(因为使用了 Outer Join) if filters: search_field = filters.get('searchField', 'all') keyword = filters.get('keyword') if keyword: kw = f"%{keyword}%" if search_field == 'name': query = query.filter(MaterialBase.name.ilike(kw)) elif search_field == 'common_name': query = query.filter(MaterialBase.common_name.ilike(kw)) elif search_field == 'spec': query = query.filter(MaterialBase.spec_model.ilike(kw)) else: query = query.filter(or_( MaterialBase.name.ilike(kw), MaterialBase.common_name.ilike(kw), MaterialBase.spec_model.ilike(kw) )) # ============================================================ # 【行级数据隔离】基于 JWT 多租户公司过滤 # ============================================================ from app.utils.decorators import get_current_company_filter company_limit = get_current_company_filter() if company_limit is not None: query = query.filter(MaterialBase.company_name == company_limit) category = filters.get('category') if category is not None and category != '': # 在末尾拼接 '%' 实现前缀模糊匹配 query = query.filter(MaterialBase.category.ilike(f"{category.strip()}%")) type_val = filters.get('type') if type_val is not None and type_val != '': query = query.filter(MaterialBase.material_type.ilike(type_val.strip())) if filters.get('isEnabled') is not None: val_str = str(filters['isEnabled']).lower() is_active = val_str in ['1', 'true', 'yes', 't'] query = query.filter(MaterialBase.is_enabled == is_active) # 库存状态筛选 (has_stock) - 使用子查询列 has_stock = filters.get('has_stock') if has_stock and str(has_stock).lower() in ['true', '1', 'yes']: query = query.filter(inner_sub.c.total_inv > 0) # ============================================================ # 【修复3】高级筛选:统一使用 WHERE (filter) # 聚合字段同样是子查询列,直接用 filter 过滤即可 # ============================================================ advanced_filters = filters.get('advancedFilters', []) if filters else [] filter_conditions = [] if advanced_filters: allowed_fields = { 'companyName': 'company_name', 'name': 'name', 'commonName': 'common_name', 'category': 'category', 'type': 'material_type', 'spec': 'spec_model', 'unit': 'unit', 'referencePrice': 'reference_price', 'inventoryCount': 'total_inv', 'availableCount': 'total_avail' } # 聚合字段列表 aggregate_fields = {'inventoryCount', 'availableCount'} field_permission_map = { 'companyName': 'material_list:companyName', 'name': 'material_list:name', 'commonName': 'material_list:commonName', 'category': 'material_list:category', 'type': 'material_list:type', 'spec': 'material_list:spec', 'unit': 'material_list:unit', 'referencePrice': 'material_list:referencePrice', 'inventoryCount': 'material_list:inventoryCount', 'availableCount': 'material_list:availableCount' } for condition in advanced_filters: field = condition.get('field') operator = condition.get('operator') value = condition.get('value') if not field or not operator or value is None: continue db_field = allowed_fields.get(field) if not db_field: continue # 权限校验 if user_permissions is not None: perm_code = field_permission_map.get(field) if 'material_list:*' in user_permissions: pass elif perm_code and perm_code not in user_permissions: continue # 判断是否为聚合字段 is_aggregate = field in aggregate_fields if is_aggregate: # 聚合字段:直接使用子查询列,通过 WHERE (filter) 过滤 col = inner_sub.c.total_inv if field == 'inventoryCount' else inner_sub.c.total_avail try: num_val = float(value) except ValueError: continue if operator == 'eq': filter_conditions.append(col == num_val) elif operator == 'ne': filter_conditions.append(col != num_val) elif operator == 'ge': filter_conditions.append(col >= num_val) elif operator == 'le': filter_conditions.append(col <= num_val) else: # 非聚合字段:使用 WHERE column = getattr(MaterialBase, db_field, None) if column is None: continue if operator == 'eq': filter_conditions.append(column == value) elif operator == 'ne': filter_conditions.append(column != value) elif operator == 'contains': filter_conditions.append(column.ilike(f'%{value}%')) elif operator == 'not_contains': filter_conditions.append(~column.ilike(f'%{value}%')) elif operator == 'ge': try: num_val = float(value) filter_conditions.append(column >= num_val) except ValueError: continue elif operator == 'le': try: num_val = float(value) filter_conditions.append(column <= num_val) except ValueError: continue # 应用 WHERE 条件(统一使用 filter) if filter_conditions: query = query.filter(and_(*filter_conditions)) # 排序处理(支持全字段) order_by_column = filters.get('orderByColumn', '') is_asc = filters.get('isAsc', None) # 信任前端传递的预警排序开关(前端已基于权限注入) enable_warning_sort = filters.get('enableWarningSort', False) if enable_warning_sort: print("====== [DEBUG] 成功进入预警强排逻辑 ======") inv_val = inner_sub.c.total_inv red_val = cast(MaterialWarningSetting.red_threshold, Numeric) yellow_val = cast(MaterialWarningSetting.yellow_threshold, Numeric) # 预警等级计算:红=2, 黄=1, 正常=0 warning_level = case( (and_(MaterialWarningSetting.is_enabled.is_(True), red_val.isnot(None), inv_val <= red_val), 2), (and_(MaterialWarningSetting.is_enabled.is_(True), yellow_val.isnot(None), inv_val <= yellow_val), 1), else_=0 ) # 统一计算缺口 (Shortage) = 目标阈值 - 当前库存 # 红色算红色的缺口,黄色算黄色的缺口,越大说明缺的越多 shortage = case( (and_(MaterialWarningSetting.is_enabled.is_(True), red_val.isnot(None), inv_val <= red_val), red_val - inv_val), (and_(MaterialWarningSetting.is_enabled.is_(True), yellow_val.isnot(None), inv_val <= yellow_val), yellow_val - inv_val), else_=0 ) query = query.order_by( desc(warning_level), # 1. 先按红、黄、正常排 desc(shortage), # 2. 同级别内,缺口越大的排越上面 desc(inv_val), # 3. 缺口一样,库存多的排上面 desc(MaterialBase.id) ) elif order_by_column: # 字段映射 - 使用子查询列进行排序 sort_field_map = { 'companyName': MaterialBase.company_name, 'name': MaterialBase.name, 'commonName': MaterialBase.common_name, 'category': MaterialBase.category, 'type': MaterialBase.material_type, 'spec': MaterialBase.spec_model, 'unit': MaterialBase.unit, 'referencePrice': MaterialBase.reference_price, 'inventoryCount': inner_sub.c.total_inv, 'availableCount': inner_sub.c.total_avail } sort_column = sort_field_map.get(order_by_column) if sort_column is not None: if is_asc == 'asc': query = query.order_by(sort_column.asc()) elif is_asc == 'desc': query = query.order_by(sort_column.desc()) else: # 默认排序:优先按总库存数降序,当库存相同时,再按规格型号升序 query = query.order_by(inner_sub.c.total_inv.desc(), MaterialBase.spec_model.asc()) # 分页 pagination = query.paginate(page=page, per_page=limit, error_out=False) items_list = [] for row in pagination.items: # 防弹解包逻辑:直接判断自身是否有 to_dict 方法 if hasattr(row, 'to_dict'): # 说明查询只返回了 MaterialBase 单一对象 item = row inv = avail = warning_enabled = warning_yellow = warning_red = None else: # 说明返回了 Row 对象 (包含多个字段) item = row[0] inv = row[1] if len(row) > 1 else None avail = row[2] if len(row) > 2 else None warning_enabled = row[3] if len(row) > 3 else False warning_yellow = row[4] if len(row) > 4 else 0 warning_red = row[5] if len(row) > 5 else 0 warning_red_emails = row[6] if len(row) > 6 else None warning_yellow_emails = row[7] if len(row) > 7 else None # 安全兜底 if not hasattr(item, 'to_dict'): continue item_dict = item.to_dict() item_dict['inventoryCount'] = float(inv) if inv is not None else 0 item_dict['availableCount'] = float(avail) if avail is not None else 0 # 处理预警信息(仅当用户有权限时) if has_warning_permission: item_dict['warningEnabled'] = bool(warning_enabled) if warning_enabled is not None else False item_dict['warningYellow'] = float(warning_yellow) if warning_yellow is not None else None item_dict['warningRed'] = float(warning_red) if warning_red is not None else None item_dict['warningRedEmails'] = warning_red_emails or '' item_dict['warningYellowEmails'] = warning_yellow_emails or '' # 计算预警状态 if warning_enabled: invQty = item_dict['inventoryCount'] # 优先判断红色预警(如果设置了红阈值,且库存 <= 红阈值) if warning_red is not None and invQty <= warning_red: item_dict['warningStatus'] = 2 # 红色 # 其次判断黄色预警(如果设置了黄阈值,且库存 <= 黄阈值) elif warning_yellow is not None and invQty <= warning_yellow: item_dict['warningStatus'] = 1 # 黄色 # 都不满足则正常 else: item_dict['warningStatus'] = 0 # 正常 else: item_dict['warningStatus'] = 0 # 未开启或未设置 items_list.append(item_dict) return {"total": pagination.total, "items": items_list} except Exception as e: traceback.print_exc() print(f"查询基础信息列表失败: {e}") return {"total": 0, "items": []} @staticmethod def get_odoo_summary(keyword=None, is_enabled=None): """ Odoo 分组摘要 API(极轻量,不 JOIN 任何库存表) 执行: SELECT category, COUNT(id) AS count FROM material_base WHERE ... GROUP BY category ORDER BY category 返回: [{"category": "IRIS/半成品/高塔监测T", "count": 54}, ...] """ try: query = db.session.query( MaterialBase.category, func.count(MaterialBase.id).label('count') ) # 状态过滤 if is_enabled is not None: query = query.filter(MaterialBase.is_enabled == is_enabled) # 关键词搜索(与 get_list 行为一致) if keyword: kw = f'%{keyword.strip()}%' query = query.filter(or_( MaterialBase.name.ilike(kw), MaterialBase.common_name.ilike(kw), MaterialBase.spec_model.ilike(kw) )) # 行级数据隔离 from app.utils.decorators import get_current_company_filter company_limit = get_current_company_filter() if company_limit is not None: query = query.filter(MaterialBase.company_name == company_limit) # GROUP BY + ORDER query = query.group_by(MaterialBase.category) \ .order_by(MaterialBase.category) rows = query.all() return [ {"category": row.category or "未分类", "count": row.count} for row in rows ] except Exception as e: traceback.print_exc() print(f"查询 Odoo 摘要失败: {e}") return [] @staticmethod def get_distinct_options(): """ 获取所有已存在的类别、类型、公司 (去重且排序) """ try: # 1. 类别 (获取后在内存或前端做层级处理,这里先按字母序返回扁平列表) categories = db.session.query(MaterialBase.category) \ .filter(MaterialBase.category != None, MaterialBase.category != '') \ .distinct().all() # 对类别进行排序 sorted_categories = sorted([c[0] for c in categories]) # 2. 类型 types = db.session.query(MaterialBase.material_type) \ .filter(MaterialBase.material_type != None, MaterialBase.material_type != '') \ .distinct().all() sorted_types = sorted([t[0] for t in types]) # 3. 公司 companies = db.session.query(MaterialBase.company_name) \ .filter(MaterialBase.company_name != None, MaterialBase.company_name != '') \ .distinct().all() sorted_companies = sorted([c[0] for c in companies]) return { "categories": sorted_categories, "types": sorted_types, "companies": sorted_companies } except Exception as e: traceback.print_exc() return {"categories": [], "types": [], "companies": []} @staticmethod def get_distinct_units(): """ 获取所有已存在且非空的计量单位(去重 + 排序)。 用于前端"基础信息"新增/编辑弹窗的"计量单位"下拉历史记录。 SQL 语义: SELECT DISTINCT unit FROM material_base WHERE unit IS NOT NULL AND unit != '' ORDER BY unit ASC """ try: rows = db.session.query(MaterialBase.unit) \ .filter(MaterialBase.unit.isnot(None), MaterialBase.unit != '') \ .distinct() \ .all() sorted_units = sorted([u[0] for u in rows if u[0]]) return sorted_units except Exception as e: traceback.print_exc() print(f"查询计量单位字典失败: {e}") return [] @staticmethod def create_material(data): """新增基础信息""" try: if not data.get('name') or not data.get('spec'): raise ValueError("名称和规格型号不能为空") exist = MaterialBase.query.filter_by( name=data['name'], spec_model=data['spec'] ).first() if exist: raise ValueError(f"已存在相同名称和规格的数据 (ID: {exist.id})") # 【核心修改】:兼容前端传来的布尔值 raw_enabled = data.get('isEnabled', True) is_enabled_val = str(raw_enabled).lower() in ['1', 'true', 'yes', 't'] if raw_enabled is not None else True new_material = MaterialBase( company_name=data.get('companyName'), name=data['name'], common_name=data.get('commonName'), spec_model=data['spec'], category=data.get('category'), material_type=data.get('type'), unit=data.get('unit'), visibility_level=data.get('visibilityLevel'), manual_link=json.dumps(data.get('generalManual', [])), product_image=json.dumps(data.get('generalImage', [])), product_image_remark=data.get('productImageRemark', ''), manual_link_remark=data.get('manualLinkRemark', ''), reference_price=data.get('referencePrice'), is_enabled=is_enabled_val ) db.session.add(new_material) db.session.flush() # 获取 new_material.id # 先提交主事务,图片向量异步后台提取 db.session.commit() image_list = data.get('generalImage', []) if isinstance(image_list, list) and image_list: from flask import current_app from app.utils.executor import run_embedding_task run_embedding_task( ImageEmbeddingService.save_embeddings_background, current_app._get_current_object(), ImageEmbeddingService.MODULE_MATERIAL_BASE, new_material.id, image_list ) return new_material except Exception as e: db.session.rollback() raise e @staticmethod def update_material(m_id, data): """修改基础信息""" try: material = MaterialBase.query.get(m_id) if not material: raise ValueError("数据不存在") # 更新字段 if 'companyName' in data: material.company_name = data['companyName'] if 'name' in data: material.name = data['name'] if 'commonName' in data: material.common_name = data['commonName'] if 'spec' in data: material.spec_model = data['spec'] if 'category' in data: material.category = data['category'] if 'type' in data: material.material_type = data['type'] if 'unit' in data: material.unit = data['unit'] if 'visibilityLevel' in data: material.visibility_level = data['visibilityLevel'] if 'productImageRemark' in data: material.product_image_remark = data['productImageRemark'] if 'manualLinkRemark' in data: material.manual_link_remark = data['manualLinkRemark'] if 'generalManual' in data: material.manual_link = json.dumps(data['generalManual']) # ★ 仅当 payload 明确携带 generalImage 时才更新图片(含传 [] 表示清空); # 修改其他字段时不动图片,避免误清空 product_image if 'generalImage' in data: new_photo_list = data['generalImage'] material.product_image = json.dumps(new_photo_list) # 立即触发异步向量提取,不阻塞主事务提交 if isinstance(new_photo_list, list) and new_photo_list: from flask import current_app from app.utils.executor import run_embedding_task run_embedding_task( ImageEmbeddingService.save_embeddings_background, current_app._get_current_object(), ImageEmbeddingService.MODULE_MATERIAL_BASE, material.id, new_photo_list ) else: # 图片被清空([])时,同步清理该物料的图片向量 ImageEmbeddingService.delete_embeddings( ImageEmbeddingService.MODULE_MATERIAL_BASE, material.id ) # 【核心修改】:兼容前端传来的布尔值 if 'referencePrice' in data: material.reference_price = data['referencePrice'] if 'isEnabled' in data: raw_enabled = data['isEnabled'] material.is_enabled = str(raw_enabled).lower() in ['1', 'true', 'yes', 't'] db.session.commit() # ★★★ 级联缓存失效:物料信息变更后,清除所有涉及该物料的 BOM 树缓存 try: from app.models.bom import BomTable from app.extensions import redis_client # 查出所有以该物料为父件或子件的 bom_no(去重) affected = db.session.query(BomTable.bom_no).filter( or_(BomTable.parent_id == m_id, BomTable.child_id == m_id) ).distinct().all() affected_bom_nos = [r[0] for r in affected] if affected_bom_nos: for bom_no in affected_bom_nos: # 清除最新版缓存 redis_client.delete(f'bom:tree:{bom_no}') # 清除所有版本缓存(通配符) for key in redis_client.scan_iter(f'bom:tree:{bom_no}:*'): redis_client.delete(key) print(f"🔁 物料 {m_id} 变更,已级联清除 {len(affected_bom_nos)} 个 BOM 缓存: {affected_bom_nos}") except Exception as cache_err: # Redis 报错不阻断业务返回 print(f"⚠️ BOM 缓存级联清除失败(不阻断业务): {cache_err}") return material except Exception as e: db.session.rollback() raise e @staticmethod def delete_material(m_id): """ 删除基础信息 (带依赖检查) """ try: material = MaterialBase.query.get(m_id) if not material: raise ValueError("数据不存在") # 提前获取物料名称用于审计日志 material_name = material.name buy_usage_count = StockBuy.query.filter_by(base_id=m_id).count() semi_usage_count = StockSemi.query.filter_by(base_id=m_id).count() prod_usage_count = StockProduct.query.filter_by(base_id=m_id).count() bom_usage_count = BomTable.query.filter( or_(BomTable.parent_id == m_id, BomTable.child_id == m_id) ).count() total_usage = buy_usage_count + semi_usage_count + prod_usage_count + bom_usage_count if total_usage > 0: raise ValueError( f"无法删除:该基础物料正被使用中。\n" f"- 采购库存记录: {buy_usage_count} 条\n" f"- 半成品库存记录: {semi_usage_count} 条\n" f"- 成品库存记录: {prod_usage_count} 条\n" f"- BOM配方引用: {bom_usage_count} 条\n" f"请先清理相关库存/BOM绑定或仅「禁用」此条目。" ) # 删除时同步清理向量记录 ImageEmbeddingService.delete_embeddings( ImageEmbeddingService.MODULE_MATERIAL_BASE, material.id ) db.session.delete(material) db.session.commit() return material_name except Exception as e: db.session.rollback() print(f"删除基础信息失败: {e}") raise e # ============================================================================== # [核心修改] 统一资产统计导出(增加最高单价计算逻辑) # ============================================================================== @staticmethod def export_excel(filters=None, user_permissions=None): """ 全口径资产统计报表: 根据基础信息列表和库存表,计算出每个物料的最高历史单价并进行导出 """ try: # 1. 构造基础信息的筛选条件 (用于过滤库存) filter_conditions = [] # ============================================================ # 【行级数据隔离】基于 JWT 多租户公司过滤 # ★ 必须无条件执行:此前这段嵌套在 `if filters:` 内部,只要调用方 # 不传或传空筛选条件,公司隔离就会被整段跳过且不报错(静默失效)。 # company_limit 的取值规则见 get_current_company_filter(): # - 普通用户 → JWT 中记录的本公司,强制隔离 # - 超管 / 跨域角色 → 仅当请求显式带公司标识时才限制,否则返回 None # ============================================================ from app.utils.decorators import get_current_company_filter company_limit = get_current_company_filter() if company_limit is not None: filter_conditions.append(MaterialBase.company_name == company_limit) if filters: if filters.get('keyword'): kw = f"%{filters['keyword']}%" filter_conditions.append(or_( MaterialBase.name.ilike(kw), MaterialBase.common_name.ilike(kw), MaterialBase.spec_model.ilike(kw), MaterialBase.company_name.ilike(kw) )) category = filters.get('category') if category is not None and category != '': # 同样在末尾拼接 '%' filter_conditions.append(MaterialBase.category.ilike(f"{category.strip()}%")) type_val = filters.get('type') if type_val is not None and type_val != '': filter_conditions.append(MaterialBase.material_type.ilike(type_val.strip())) # 【核心修改】:增强布尔值解析 if filters.get('isEnabled') is not None: val_str = str(filters['isEnabled']).lower() is_active = val_str in ['1', 'true', 'yes', 't'] filter_conditions.append(MaterialBase.is_enabled == is_active) # 2. 查询三表(加 ORDER BY 替代内存排序,yield_per 分块流式读取) query_buy = db.session.query(StockBuy, MaterialBase).join( MaterialBase, StockBuy.base_id == MaterialBase.id ).filter(StockBuy.stock_quantity > 0) for cond in filter_conditions: query_buy = query_buy.filter(cond) query_buy = query_buy.order_by( MaterialBase.company_name, MaterialBase.spec_model, MaterialBase.id, StockBuy.batch_number ) query_semi = db.session.query(StockSemi, MaterialBase).join( MaterialBase, StockSemi.base_id == MaterialBase.id ).filter(StockSemi.stock_quantity > 0) for cond in filter_conditions: query_semi = query_semi.filter(cond) query_semi = query_semi.order_by( MaterialBase.company_name, MaterialBase.spec_model, MaterialBase.id ) query_product = db.session.query(StockProduct, MaterialBase).join( MaterialBase, StockProduct.base_id == MaterialBase.id ).filter(StockProduct.stock_quantity > 0) for cond in filter_conditions: query_product = query_product.filter(cond) query_product = query_product.order_by( MaterialBase.company_name, MaterialBase.spec_model, MaterialBase.id ) # 预先计算最高单价(分块扫描,O(unique_base_ids) 内存) buy_max_prices, semi_max_prices, product_max_prices = {}, {}, {} for stock, base in query_buy.yield_per(2000): price = float(stock.pre_tax_unit_price or 0) if price > buy_max_prices.get(base.id, 0): buy_max_prices[base.id] = price for stock, base in query_semi.yield_per(2000): price = float(stock.manual_cost or 0) if price > semi_max_prices.get(base.id, 0): semi_max_prices[base.id] = price for stock, base in query_product.yield_per(2000): price = float(stock.manual_cost or 0) if price > product_max_prices.get(base.id, 0): product_max_prices[base.id] = price def get_highest_price(base_id): if base_id in buy_max_prices and buy_max_prices[base_id] > 0: return buy_max_prices[base_id] if base_id in semi_max_prices and semi_max_prices[base_id] > 0: return semi_max_prices[base_id] if base_id in product_max_prices and product_max_prices[base_id] > 0: return product_max_prices[base_id] return 0.0 # 3. 流式写入 Excel(write_only=True 不占内存) wb = Workbook(write_only=True) ws = wb.create_sheet("库存统计") # ★ 表头必须与下方 _write_stock_rows 里 ws.append([...]) 的 22 列严格一一对应, # 顺序错位会让整张报表串列。该变量此前漏定义,导致导出接口必抛 NameError。 headers = [ '所属公司', '物料名称', '规格型号', '物料类型', '一级分类', '二级分类', '三级分类', '四级分类', '五级分类', '单位', '库存类型', '批次/序列号', '库位', '供应商/负责人', '入库/生产日期', '库存数量', '可用数量', '不含税单价', '不含税金额', '税率(%)', '含税单价', '含税金额' ] ws.append(headers) def _write_stock_rows(iterable, type_name, price_getter, tax_getter=None): for stock, base in iterable: qty = float(stock.stock_quantity or 0) price = price_getter(stock, base) tax = tax_getter(stock) if tax_getter else 0.0 price_incl = price * (1 + tax / 100.0) if tax else price ident = (getattr(stock, 'batch_number', None) or getattr(stock, 'serial_number', None) or getattr(stock, 'barcode', None) or getattr(stock, 'sku', '') or '') date_val = (getattr(stock, 'in_date', None) or getattr(stock, 'production_date', None)) date_str = date_val.strftime('%Y-%m-%d') if isinstance(date_val, datetime.date) else '' source = (getattr(stock, 'supplier_name', None) or getattr(stock, 'production_manager', '') or '') cat_parts = (base.category or "").split('/') while len(cat_parts) < 5: cat_parts.append("") ws.append([ base.company_name, base.name, base.spec_model, base.material_type, cat_parts[0], cat_parts[1], cat_parts[2], cat_parts[3], cat_parts[4], base.unit, type_name, ident, getattr(stock, 'warehouse_location', '') or '', source, date_str, qty, float(stock.available_quantity or 0), price, qty * price, tax, price_incl, qty * price_incl ]) _write_stock_rows( query_buy.yield_per(2000), "采购件", lambda s, b: get_highest_price(b.id), lambda s: float(s.tax_rate or 0) ) _write_stock_rows( query_semi.yield_per(2000), "半成品", lambda s, b: float(s.manual_cost or 0) ) _write_stock_rows( query_product.yield_per(2000), "成品", lambda s, b: float(s.manual_cost or 0) ) output = io.BytesIO() wb.save(output) output.seek(0) return output except Exception as e: traceback.print_exc() raise e # 支持二级分类的前缀集合(如 OPT1, OPT2, LICA1 等子系列分组) SUB_CATEGORY_PREFIXES = {'OPT', 'LICA', 'M', 'UAV', 'CF', 'GPS'} @staticmethod def get_latest_specs(): """ 规格连号助手 — 智能分组统计(v2: 流式读取 + Redis缓存 + 宽松regex) 匹配模式: PREFIX[-_]?NUMBERS[SUFFIX], 如: OPT12046 → OPT, 1, 2046 LICA-3000 → LICA, 3000, '' M3x12 → M, 3, 'x12' 分组规则: 前缀在 SUB_CATEGORY_PREFIXES 中 → 前缀+首位数字作为key 其他 → 只用前缀作为key """ import re import json as json_module from collections import defaultdict CACHE_KEY = 'inventory:specs:grouped' CACHE_TTL = 3600 # ── Redis 缓存 ── try: from app.extensions import redis_client if redis_client: cached = redis_client.get(CACHE_KEY) if cached: return json_module.loads(cached) except Exception: pass # Redis 不可用时降级 # ── 流式查询(yield_per 分批 + limit 防 OOM) ── pattern = re.compile(r'^([A-Za-z]+)[-_]?(\d+)(.*)$') groups = defaultdict(list) rows = MaterialBase.query.with_entities( MaterialBase.id, MaterialBase.spec_model ).filter( MaterialBase.spec_model.isnot(None), MaterialBase.spec_model != '' ).limit(10000).all() for row in rows: spec = row.spec_model if not spec: continue base_spec = spec.split('/')[0] match = pattern.match(base_spec) if not match: continue prefix = match.group(1).upper() num_str = match.group(2) suffix = match.group(3) if not num_str: continue num = int(num_str) sub_cat = num_str[0] # 首位数字作为子分类 # 分组 key if prefix in MaterialBaseService.SUB_CATEGORY_PREFIXES and sub_cat: key = f"{prefix}_{sub_cat}" else: key = prefix groups[key].append((num, spec)) # ── 生成结果 ── result = [] for key, items in groups.items(): items.sort(key=lambda x: x[0]) max_num, max_spec = items[-1] result.append({ 'group': key, 'count': len(items), 'latest': max_spec, 'max_num': max_num }) result.sort(key=lambda x: (-x['count'], x['group'])) # ── 写入 Redis 缓存 ── try: from app.extensions import redis_client if redis_client: redis_client.setex(CACHE_KEY, CACHE_TTL, json_module.dumps(result)) except Exception: pass return result