diff --git a/netx_api/biz_state/compare_service.py b/netx_api/biz_state/compare_service.py index 5c30368..983f7fb 100644 --- a/netx_api/biz_state/compare_service.py +++ b/netx_api/biz_state/compare_service.py @@ -1348,69 +1348,88 @@ def _port_map_dict(db: Session, mapping_id: str) -> dict[str, str]: def _load_metric_rows(db: Session, *, batch_id: str, metric_id: str) -> list[dict[str, Any]]: - """Load metric rows with streaming fetch to avoid holding the full ORM set.""" + """Load metric rows in LIMIT/OFFSET chunks (avoid ORM yield_per + unique() clash).""" if metric_id == "lldp_neighbor": - q = ( - db.query(BizStateLldpNeighbor) - .filter(BizStateLldpNeighbor.batch_id == batch_id) - .execution_options(stream_results=True, yield_per=_LOAD_YIELD_PER) - ) out: list[dict[str, Any]] = [] - for n in q: - out.append( - { - "local_if": n.local_if, - "remote_sys": n.remote_sys, - "remote_if": n.remote_if, - "remote_ip": n.remote_ip, - "protocol": n.protocol, - "_netx": { - "batch_id": batch_id, - "batch_command_id": n.batch_command_id or "", - "task_id": n.task_id or "", - "ne_id": n.ne_id or "", - "collected_at": n.collected_at.isoformat() + "Z" - if n.collected_at - else None, - "row_id": n.id, - }, - } + offset = 0 + while True: + chunk = ( + db.query(BizStateLldpNeighbor) + .filter(BizStateLldpNeighbor.batch_id == batch_id) + .order_by(BizStateLldpNeighbor.id.asc()) + .offset(offset) + .limit(_LOAD_YIELD_PER) + .all() ) - db.expunge(n) + if not chunk: + break + for n in chunk: + out.append( + { + "local_if": n.local_if, + "remote_sys": n.remote_sys, + "remote_if": n.remote_if, + "remote_ip": n.remote_ip, + "protocol": n.protocol, + "_netx": { + "batch_id": batch_id, + "batch_command_id": n.batch_command_id or "", + "task_id": n.task_id or "", + "ne_id": n.ne_id or "", + "collected_at": n.collected_at.isoformat() + "Z" + if n.collected_at + else None, + "row_id": n.id, + }, + } + ) + db.expunge(n) + offset += len(chunk) + if len(chunk) < _LOAD_YIELD_PER: + break return out # Generic tabular metrics (ISIS / interface / ARP / ND6 / BGP …) from ..models import BizStateMetricRow - q = ( - db.query(BizStateMetricRow) - .filter( - BizStateMetricRow.batch_id == batch_id, - BizStateMetricRow.metric_id == metric_id, - ) - .order_by(BizStateMetricRow.seq.asc(), BizStateMetricRow.id.asc()) - .execution_options(stream_results=True, yield_per=_LOAD_YIELD_PER) - ) out = [] - for r in q: - # Raw rows only — filtering belongs to the compare sheet template - # (``row_filters``), not metric-specific branches here. - # ``_netx`` is collector provenance (stripped before field compare). - out.append( - { - **dict(r.data_json or {}), - "_netx": { - "batch_id": batch_id, - "batch_command_id": r.batch_command_id or "", - "task_id": r.task_id or "", - "ne_id": r.ne_id or "", - "collected_at": r.collected_at.isoformat() + "Z" - if r.collected_at - else None, - "row_id": r.id, - }, - } + offset = 0 + while True: + chunk = ( + db.query(BizStateMetricRow) + .filter( + BizStateMetricRow.batch_id == batch_id, + BizStateMetricRow.metric_id == metric_id, + ) + .order_by(BizStateMetricRow.seq.asc(), BizStateMetricRow.id.asc()) + .offset(offset) + .limit(_LOAD_YIELD_PER) + .all() ) - db.expunge(r) + if not chunk: + break + for r in chunk: + # Raw rows only — filtering belongs to the compare sheet template + # (``row_filters``), not metric-specific branches here. + # ``_netx`` is collector provenance (stripped before field compare). + out.append( + { + **dict(r.data_json or {}), + "_netx": { + "batch_id": batch_id, + "batch_command_id": r.batch_command_id or "", + "task_id": r.task_id or "", + "ne_id": r.ne_id or "", + "collected_at": r.collected_at.isoformat() + "Z" + if r.collected_at + else None, + "row_id": r.id, + }, + } + ) + db.expunge(r) + offset += len(chunk) + if len(chunk) < _LOAD_YIELD_PER: + break if out: return out # Known metric with zero rows is OK; unknown metric still errors