src.dackar.RCA.equipment_similarity.kg_equipment_poller¶
kg_equipment_poller — KGEquipmentPoller.
One-time population script: queries all element_usage → element_definition nodes from Neo4j (with linked failure modes), builds spec text via EquipmentSpecBuilder, and upserts into EquipmentSpecStore.
Run at KG setup time and re-run whenever the KG is updated.
Usage:
from py2neo import Graph
from storage.chroma_store import ChromaRecordStore
from equipment_similarity.equipment_spec_store import EquipmentSpecStore
from equipment_similarity.kg_equipment_poller import KGEquipmentPoller
neo4j = Graph("bolt://localhost:7687", auth=("neo4j", "password"))
chroma = ChromaRecordStore(persist_directory="./chroma_store")
spec_store = EquipmentSpecStore(chroma)
poller = KGEquipmentPoller(client=neo4j)
n = poller.poll_and_upsert(spec_store, batch_size=100)
print(f"Upserted {n} component specs")
Attributes¶
Classes¶
Polls Neo4j for all element_usage / element_definition pairs and |
Module Contents¶
- src.dackar.RCA.equipment_similarity.kg_equipment_poller._SPEC_QUERY = Multiline-String[source]¶
Show Value
""" MATCH (c:element_usage)-[:instance_of]->(def:element_definition) OPTIONAL MATCH (c)-[:subject_to|applies_to]-(fm:failure_mode) RETURN c.id AS component_id, c.name AS component_name, def.domain_category AS domain_category, def.structural_kind AS structural_kind, def.nominal_size AS nominal_size, def.design_pressure AS design_pressure, def.design_temperature AS design_temperature, def.material_spec AS material_spec, def.manufacturer AS manufacturer, def.model_number AS model_number, collect(DISTINCT fm.name) AS failure_mode_names, collect(DISTINCT fm.failure_mechanism) AS failure_mechanisms ORDER BY component_id """
- class src.dackar.RCA.equipment_similarity.kg_equipment_poller.KGEquipmentPoller(client, database=None)[source]¶
Polls Neo4j for all element_usage / element_definition pairs and upserts their spec embeddings into EquipmentSpecStore.
- Parameters:
client (Any) – A
py2neo.Graph(or any object with arun(query, **params)method returning an iterable of record-like objects).database (Optional[str]) – Optional Neo4j database name. If
None, the default database for the connection is used.
- poll_and_upsert(spec_store, batch_size=100)[source]¶
Query all components from the KG and upsert their specs.
- Parameters:
spec_store (src.dackar.RCA.equipment_similarity.equipment_spec_store.EquipmentSpecStore) –
EquipmentSpecStoreto upsert into.batch_size (int) – Number of records per Chroma upsert call. Larger batches are faster but use more memory.
- Returns:
Total number of component specs upserted.
- Return type:
int
Notes
This is upsert-only. Each component’s spec is written by a stable
record_id(equip_spec::{component_id}), so re-running refreshes existing components and adds new ones, but specs for components that were removed from the KG are not deleted from the Chroma collection — they linger as stale vectors. A full refresh that drops removed components requires rebuilding theequipment_specscollection from scratch (delete the collection, then re-poll).