Files
KCGL/inventory-backend/app/services/bom_service.py
yueli 3a07f1a272 fix(bom): BOM 详情查询改为 LEFT JOIN,子件物料被删后仍返回父件结构
- get_bom_detail 使用 outerjoin 替代 join
- 子件物料被删除后 child_name 显示 [已删除物料] 而非误报 BOM 不存在
- 修复按 BOM 出库时因子件删除导致 404 BOM 不存在
2026-08-31 10:28:22 +08:00

588 lines
22 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

from app.extensions import db
from app.models.bom import BomTable
from app.models.base import MaterialBase
from app.models.inbound.buy import StockBuy
from app.models.inbound.semi import StockSemi
from app.models.inbound.product import StockProduct
from app.utils.decorators import get_current_company_filter
from sqlalchemy import func, distinct, or_, case
from collections import defaultdict
import uuid
import json
import logging
from datetime import datetime
logger = logging.getLogger(__name__)
# Redis 缓存键前缀 + TTL
BOM_CACHE_PREFIX = 'bom:tree'
BOM_CACHE_TTL = 43200 # 12小时
def _get_redis():
"""
获取 Redis 客户端实例,带容错保护。
若 extensions 中没有 redis_client 或连接失败,返回 None。
绝不抛出异常,确保业务不因此中断。
"""
try:
from app.extensions import redis_client
return redis_client
except Exception:
return None
def _cache_get(bom_no, version=None):
"""
从 Redis 读取 BOM 缓存。
键 = bom:tree:{bom_no} 或 bom:tree:{bom_no}:{version}
返回:反序列化后的 dict 或 None
"""
client = _get_redis()
if not client:
return None
key = f"{BOM_CACHE_PREFIX}:{bom_no}" + (f":{version}" if version else "")
try:
raw = client.get(key)
if raw:
logger.debug(f"[BOM Cache] HIT key={key}")
return json.loads(raw)
logger.debug(f"[BOM Cache] MISS key={key}")
return None
except Exception as e:
logger.warning(f"[BOM Cache] GET 失败,降级查库. key={key}, err={e}")
return None
def _cache_set(bom_no, version, data):
"""
将 BOM 数据写入 Redis设置 12 小时 TTL。
即使写入失败也只是打日志,不阻断业务流程。
"""
client = _get_redis()
if not client:
return
key = f"{BOM_CACHE_PREFIX}:{bom_no}" + (f":{version}" if version else "")
try:
client.setex(key, BOM_CACHE_TTL, json.dumps(data, ensure_ascii=False))
logger.debug(f"[BOM Cache] SET key={key} ttl={BOM_CACHE_TTL}s")
except Exception as e:
logger.warning(f"[BOM Cache] SET 失败,已忽略. key={key}, err={e}")
def _cache_delete(bom_no, version=None):
"""
删除 Redis 中指定 BOM 的缓存条目。
在写操作(增/改/删)成功后调用,确保后续读请求拿到最新数据。
"""
client = _get_redis()
if not client:
return
# 删除版本级缓存
if version:
key = f"{BOM_CACHE_PREFIX}:{bom_no}:{version}"
try:
client.delete(key)
logger.debug(f"[BOM Cache] DEL key={key}")
except Exception as e:
logger.warning(f"[BOM Cache] DEL 失败,已忽略. key={key}, err={e}")
# 同时删除"最新版"缓存(不带版本后缀),避免缓存不一致
key_latest = f"{BOM_CACHE_PREFIX}:{bom_no}"
try:
client.delete(key_latest)
logger.debug(f"[BOM Cache] DEL key={key_latest}")
except Exception as e:
logger.warning(f"[BOM Cache] DEL 失败,已忽略. key={key_latest}, err={e}")
class BomService:
# ====================== 新版 BOM 逻辑(基于 bom_no ======================
@staticmethod
def generate_bom_no():
"""生成唯一的 BOM 编号 (作为默认备选)"""
timestamp = datetime.now().strftime('%Y%m%d%H%M%S')
unique = str(uuid.uuid4())[:8]
return f'BOM-{timestamp}-{unique}'
@staticmethod
def get_bom_list(keyword=None, active_only=False, category=None, page=1, limit=15):
"""
获取所有 BOM 配方(按 bom_no + version 分组,单条 SQL 聚合 + 分页)
性能优化v2
- 消除 N+1单条 GROUP BY + string_agg 查询替代循环内逐条查询
- 消除全量加载DB 层 .paginate() 替代 .all() + Python 内存分页
- 消除 Python 二次过滤:关键词过滤完全下沉到 SQL子查询 + EXISTS 语义)
Args:
category: 可选,按 parent_category 精确过滤(用于懒加载分组展开)
"""
child_alias = db.aliased(MaterialBase)
# ===== 主聚合查询(单条 SQLGROUP BY + string_agg =====
query = db.session.query(
BomTable.bom_no,
BomTable.version,
BomTable.parent_id,
MaterialBase.name.label('parent_name'),
MaterialBase.spec_model.label('parent_spec'),
MaterialBase.category.label('parent_category'),
BomTable.is_enabled,
func.count(BomTable.child_id).label('child_count'),
func.string_agg(child_alias.name, ', ').label('child_names'),
func.string_agg(child_alias.spec_model, ', ').label('child_specs')
).join(
MaterialBase, BomTable.parent_id == MaterialBase.id
).outerjoin(
child_alias, BomTable.child_id == child_alias.id
).group_by(
BomTable.bom_no, BomTable.version, BomTable.parent_id,
MaterialBase.name, MaterialBase.spec_model, MaterialBase.category,
BomTable.is_enabled
)
# 过滤禁用状态
if active_only:
query = query.filter(BomTable.is_enabled == True)
# 【行级数据隔离】基于 JWT 多租户公司过滤
company_limit = get_current_company_filter()
if company_limit is not None:
query = query.filter(MaterialBase.company_name == company_limit)
# 按类别过滤(用于懒加载分组展开)
if category:
query = query.filter(MaterialBase.category == category)
# ===== 关键词过滤(完全下沉到 SQL消除 Python 二次过滤) =====
if keyword:
kw = f'%{keyword.strip()}%'
# 子查询:找到所有匹配的 (bom_no, version) 对
kw_child_alias = db.aliased(MaterialBase)
match_subq = db.session.query(
BomTable.bom_no,
BomTable.version
).join(
MaterialBase, BomTable.parent_id == MaterialBase.id
).outerjoin(
kw_child_alias, BomTable.child_id == kw_child_alias.id
).filter(
or_(
BomTable.bom_no.ilike(kw),
MaterialBase.name.ilike(kw),
MaterialBase.spec_model.ilike(kw),
kw_child_alias.name.ilike(kw),
kw_child_alias.spec_model.ilike(kw)
)
).distinct().subquery()
# 用子查询结果过滤主查询
query = query.join(
match_subq,
db.and_(
BomTable.bom_no == match_subq.c.bom_no,
BomTable.version == match_subq.c.version
)
)
# 排序(最新在前)
query = query.order_by(BomTable.bom_no.desc(), BomTable.version.desc())
# ===== 数据库层分页(不再 .all() 到内存) =====
pagination = query.paginate(page=page, per_page=limit, error_out=False)
# 组装结果
items = []
for row in pagination.items:
items.append({
'bom_no': row.bom_no,
'version': row.version,
'parent_id': row.parent_id,
'parent_name': row.parent_name,
'parent_spec': row.parent_spec or '',
'parent_category': row.parent_category or '',
'is_enabled': row.is_enabled,
'child_count': row.child_count,
'child_names': row.child_names or '',
'child_specs': row.child_specs or ''
})
return {
'items': items,
'total': pagination.total,
'pages': pagination.pages,
'current_page': page
}
@staticmethod
def get_bom_summary(keyword=None):
"""
BOM 分组摘要 API极轻量单条 GROUP BY + COUNT
SQL:
SELECT m.category, COUNT(DISTINCT (b.bom_no, b.version)) AS count
FROM bom_table b
JOIN material_base m ON b.parent_id = m.id
WHERE ...
GROUP BY m.category
ORDER BY m.category
Returns:
[{"category": "IRIS/半成品/无人机U", "count": 15}, ...]
"""
query = db.session.query(
MaterialBase.category,
func.count(func.distinct(
func.concat(BomTable.bom_no, '|', BomTable.version)
)).label('count')
).join(
MaterialBase, BomTable.parent_id == MaterialBase.id
)
# 行级数据隔离
company_limit = get_current_company_filter()
if company_limit is not None:
query = query.filter(MaterialBase.company_name == company_limit)
# 关键词搜索
if keyword:
kw = f'%{keyword.strip()}%'
query = query.filter(or_(
BomTable.bom_no.ilike(kw),
MaterialBase.name.ilike(kw),
MaterialBase.spec_model.ilike(kw)
))
query = query.group_by(MaterialBase.category) \
.order_by(MaterialBase.category)
rows = query.all()
return [
{"category": r.category or "未分类", "count": r.count}
for r in rows
]
@staticmethod
def get_bom_detail(bom_no, version=None):
"""
根据 bom_no (和 version) 获取配方详情。
Cache-Aside 模式(三步走):
1. 先查 Redis有值直接返回Cache Hit
2. 无值或 Redis 报错查数据库Cache Miss → Fallback
3. 数据库查好后写入 RedisTTL=12h供下次命中
注意:查询"最新版"version=None缓存键不带版本后缀
写入时也写入不带版本的键,这样无需指定 version 就能命中。
"""
# ===== 第一步:尝试从 Redis 读取缓存 =====
cached = _cache_get(bom_no, version if version else None)
if cached is not None:
# Cache Hit直接返回缓存数据不再查库
logger.debug(f"[BOM] get_bom_detail bom_no={bom_no} version={version} → 命中缓存")
return cached
# ===== 第二步Cache Miss → 查数据库 =====
logger.debug(f"[BOM] get_bom_detail bom_no={bom_no} version={version} → 查询数据库")
query = db.session.query(
BomTable,
MaterialBase.name.label('child_name'),
MaterialBase.spec_model.label('child_spec')
).outerjoin(
# ★ LEFT JOIN子件物料被删除后仍返回父件 BOM 结构,避免误报 "BOM 不存在"
MaterialBase, BomTable.child_id == MaterialBase.id
).filter(
BomTable.bom_no == bom_no
)
if version:
query = query.filter(BomTable.version == version)
else:
latest_ver = db.session.query(BomTable.version).filter_by(bom_no=bom_no) \
.order_by(BomTable.version.desc()).limit(1).scalar()
if not latest_ver:
return None
query = query.filter(BomTable.version == latest_ver)
# 记录本次实际查的版本,用于缓存键
version = latest_ver
rows = query.all()
if not rows:
return None
first = rows[0]
parent_id = first.BomTable.parent_id
parent_material = MaterialBase.query.get(parent_id)
children = []
for bom, child_name, child_spec in rows:
children.append({
'child_id': bom.child_id,
'child_name': child_name or '[已删除物料]',
'child_spec': child_spec or '',
'dosage': float(bom.dosage) if bom.dosage else 0.0,
'remark': bom.remark or ''
})
result = {
'bom_no': bom_no,
'version': first.BomTable.version,
'parent_id': parent_id,
'parent_name': parent_material.name if parent_material else '',
'parent_spec': parent_material.spec_model if parent_material else '',
'is_enabled': first.BomTable.is_enabled,
'children': children
}
# ===== 第三步:写入 Redis 缓存TTL=12h失败只打日志不阻断 =====
_cache_set(bom_no, version, result)
return result
@staticmethod
def save_bom(data):
"""保存 BOM (支持多版本),新增跨版本内容查重"""
bom_no = data.get('bom_no')
version = data.get('version', 'V1.0')
parent_id = data['parent_id']
children = data['children']
is_enabled = data.get('is_enabled', True)
if not bom_no:
raise ValueError('BOM编号不能为空')
for child in children:
if child['child_id'] == parent_id:
raise ValueError('父件与子件不能是同一物料')
# ===== 跨版本内容查重 =====
# 将当前提交的 children 转换为可比较的集合 (child_id, dosage)
current_children_set = set()
for child in children:
# 用 (child_id, dosage) 元组表示dosage 转为整数比较
dosage_val = int(child.get('dosage', 0)) if child.get('dosage') else 0
current_children_set.add((child['child_id'], dosage_val))
# 查询该 bom_no 下所有其他版本的子件配置
existing_versions = db.session.query(
BomTable.version,
BomTable.child_id,
BomTable.dosage
).filter(
BomTable.bom_no == bom_no,
BomTable.version != version # 排除当前正在保存的版本
).all()
# 按版本分组,构建每个版本的子件集合
version_children = {}
for ver, child_id, dosage in existing_versions:
if ver not in version_children:
version_children[ver] = set()
dosage_val = int(dosage) if dosage else 0
version_children[ver].add((child_id, dosage_val))
# 比对每个版本
for ver, existing_set in version_children.items():
if current_children_set == existing_set:
raise ValueError(f'保存失败!当前子件配置与已有版本 {ver} 完全一致,请勿重复保存')
# ===== 执行保存 =====
# 仅删除当前版本的旧记录(改为对象级删除以触发审计事件)
old_records = BomTable.query.filter_by(bom_no=bom_no, version=version).all()
for rec in old_records:
db.session.delete(rec)
# 【核心修复】:强制立即执行 DELETE 语句,为后续的 INSERT 腾出唯一键空间
db.session.flush()
for child in children:
bom = BomTable(
bom_no=bom_no,
version=version,
parent_id=parent_id,
child_id=child['child_id'],
dosage=child.get('dosage', 0),
remark=child.get('remark', ''),
is_enabled=is_enabled
)
db.session.add(bom)
db.session.commit()
# ===== 写入后立刻清除缓存Cache Invalidation =====
# 确保后续 get_bom_detail 读取到最新数据,而不是 stale cache
_cache_delete(bom_no, version)
logger.info(f"[BOM Cache] save_bom → 缓存已失效 bom_no={bom_no} version={version}")
return bom_no
@staticmethod
def get_bom_with_stock_by_bom_no(bom_no):
"""
根据 bom_no 获取配方详情,并计算(已修复 N+1 性能问题)
"""
detail = BomService.get_bom_detail(bom_no)
if not detail or not detail.get('children'):
return detail
# 1. 提取所有子件的 ID 列表
child_ids = [child['child_id'] for child in detail['children']]
# 2. 分别查三张库存表Python 侧合并聚合(避免 UNION ALL 子查询列名问题)
from collections import defaultdict
stock_agg = defaultdict(lambda: {'qty': 0.0, 'locs': set()})
for model in (StockBuy, StockSemi, StockProduct):
rows = db.session.query(
model.base_id,
model.available_quantity,
model.warehouse_location
).filter(
model.base_id.in_(child_ids),
model.available_quantity > 0
).all()
for base_id, qty, loc in rows:
stock_agg[base_id]['qty'] += float(qty or 0)
if loc:
stock_agg[base_id]['locs'].add(loc)
# 3. 将聚合结果转换为字典 (Map),方便后续 O(1) 极速匹配
stock_map = {
base_id: {
'qty': agg['qty'],
'loc': ', '.join(sorted(agg['locs'])) if agg['locs'] else ''
}
for base_id, agg in stock_agg.items()
}
# 4. 遍历组装数据(纯内存操作,极快)
for child in detail['children']:
base_id = child['child_id']
stat = stock_map.get(base_id, {'qty': 0, 'loc': ''})
stock_qty = float(stat['qty'])
dosage = float(child['dosage']) if child.get('dosage') else 0
child['current_stock'] = stock_qty
child['warehouse_location'] = stat['loc']
child['max_producible'] = int(stock_qty // dosage) if dosage > 0 else 0
return detail
# ====================== 兼容旧接口 ======================
@staticmethod
def get_bom_no_by_parent(parent_id):
row = BomTable.query.filter_by(parent_id=parent_id).order_by(BomTable.version.desc()).first()
return row.bom_no if row else None
@staticmethod
def create_or_update_bom(parent_id, child_list, bom_no=None, version='V1.0'):
try:
if not bom_no:
existing = BomTable.query.filter_by(parent_id=parent_id).first()
bom_no = existing.bom_no if existing else BomService.generate_bom_no()
# 改为对象级删除以触发审计事件
old_records = BomTable.query.filter_by(bom_no=bom_no, version=version).all()
for rec in old_records:
db.session.delete(rec)
for item in child_list:
bom = BomTable(
bom_no=bom_no, version=version, parent_id=parent_id,
child_id=item['child_id'], dosage=item.get('dosage', 0), remark=item.get('remark', '')
)
db.session.add(bom)
db.session.commit()
# ===== 写入后立刻清除缓存Cache Invalidation =====
_cache_delete(bom_no, version)
logger.info(f"[BOM Cache] create_or_update_bom → 缓存已失效 bom_no={bom_no} version={version}")
except Exception as e:
db.session.rollback()
logger.error(f"[BOM] create_or_update_bom 失败 bom_no={bom_no}: {e}")
raise
return True
@staticmethod
def get_bom_with_stock(parent_id):
bom_no = BomService.get_bom_no_by_parent(parent_id)
if not bom_no: return []
detail = BomService.get_bom_with_stock_by_bom_no(bom_no)
return detail['children'] if detail else []
@staticmethod
def calculate_cascade_inventory(bom_no, order_qty):
"""
根据 bom_no 和订单数量,计算所有子件的级联库存缺口。
返回结构供 AI 消费每个子件包含parent_name / spec / name / level_type /
need_qty / available_stock / suggested_qty / gap
若 BOM 不存在返回 None。
"""
# 1. 获取 BOM 明细
detail = BomService.get_bom_detail(bom_no)
if not detail or not detail.get('children'):
return None
parent_name = detail.get('parent_name', '')
# 2. 提取所有子件 ID分别查三张库存表Python 侧合并聚合
child_ids = [child['child_id'] for child in detail['children']]
from collections import defaultdict
stock_agg = defaultdict(float)
for model in (StockBuy, StockSemi, StockProduct):
rows = db.session.query(
model.base_id,
model.available_quantity
).filter(
model.base_id.in_(child_ids)
).all()
for base_id, qty in rows:
stock_agg[base_id] += float(qty or 0)
stock_map = dict(stock_agg)
# 3. 提取所有子件的基础物料信息(名称/规格/类型)
materials = db.session.query(
MaterialBase.id,
MaterialBase.name,
MaterialBase.spec_model
).filter(MaterialBase.id.in_(child_ids)).all()
mat_map = {
m.id: {'name': m.name or '', 'spec': m.spec_model or ''}
for m in materials
}
# 4. 遍历子件,计算每个子件的缺口数据
results = []
for child in detail['children']:
child_id = child['child_id']
dosage = float(child.get('dosage') or 0)
need_qty = dosage * order_qty
available_stock = stock_map.get(child_id, 0)
suggested_qty = max(0.0, min(need_qty, available_stock))
gap = available_stock - need_qty
mat_info = mat_map.get(child_id, {'name': '', 'spec': ''})
results.append({
'parent_name': parent_name,
'spec': mat_info['spec'],
'name': mat_info['name'],
'level_type': 'child',
'need_qty': round(need_qty, 4),
'available_stock': round(available_stock, 4),
'suggested_qty': round(suggested_qty, 4),
'gap': round(gap, 4),
})
return results