From bb28718daf2544b83bec8adae9220e208878d900 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Mon, 9 Feb 2026 13:37:06 +0000 Subject: [PATCH] Don't emit graph embeddings if there aren't any. Don't store graph embeddings in a knowledge store if there's an empty list. Translate between Cassandra's 'null' representing an empty list and an empty list which is what the surrounding code wants (and stored in the first place). Avoid emitting empty embedding lists Avoid output empty triple lists --- .../embeddings/graph_embeddings/embeddings.py | 11 ++--- .../extract/kg/definitions/extract.py | 42 +++++++++--------- .../trustgraph/extract/kg/ontology/extract.py | 44 +++++-------------- .../extract/kg/relationships/extract.py | 21 ++++----- .../trustgraph/storage/knowledge/store.py | 6 ++- .../trustgraph/tables/knowledge.py | 36 ++++++++------- 6 files changed, 76 insertions(+), 84 deletions(-) diff --git a/trustgraph-flow/trustgraph/embeddings/graph_embeddings/embeddings.py b/trustgraph-flow/trustgraph/embeddings/graph_embeddings/embeddings.py index 4726be4d..7b2c779b 100755 --- a/trustgraph-flow/trustgraph/embeddings/graph_embeddings/embeddings.py +++ b/trustgraph-flow/trustgraph/embeddings/graph_embeddings/embeddings.py @@ -73,12 +73,13 @@ class Processor(FlowProcessor): ) ) - r = GraphEmbeddings( - metadata=v.metadata, - entities=entities, - ) + if entities: + r = GraphEmbeddings( + metadata=v.metadata, + entities=entities, + ) - await flow("output").send(r) + await flow("output").send(r) except Exception as e: logger.error("Exception occurred", exc_info=True) diff --git a/trustgraph-flow/trustgraph/extract/kg/definitions/extract.py b/trustgraph-flow/trustgraph/extract/kg/definitions/extract.py index 8c278c25..9cf87e9a 100755 --- a/trustgraph-flow/trustgraph/extract/kg/definitions/extract.py +++ b/trustgraph-flow/trustgraph/extract/kg/definitions/extract.py @@ -168,27 +168,29 @@ class Processor(FlowProcessor): entities.append(ec) - await self.emit_triples( - flow("triples"), - Metadata( - id=v.metadata.id, - metadata=[], - user=v.metadata.user, - collection=v.metadata.collection, - ), - triples - ) + if triples: + await self.emit_triples( + flow("triples"), + Metadata( + id=v.metadata.id, + metadata=[], + user=v.metadata.user, + collection=v.metadata.collection, + ), + triples + ) - await self.emit_ecs( - flow("entity-contexts"), - Metadata( - id=v.metadata.id, - metadata=[], - user=v.metadata.user, - collection=v.metadata.collection, - ), - entities - ) + if entities: + await self.emit_ecs( + flow("entity-contexts"), + Metadata( + id=v.metadata.id, + metadata=[], + user=v.metadata.user, + collection=v.metadata.collection, + ), + entities + ) except Exception as e: logger.error(f"Definitions extraction exception: {e}", exc_info=True) diff --git a/trustgraph-flow/trustgraph/extract/kg/ontology/extract.py b/trustgraph-flow/trustgraph/extract/kg/ontology/extract.py index adcbe0a8..1644f729 100644 --- a/trustgraph-flow/trustgraph/extract/kg/ontology/extract.py +++ b/trustgraph-flow/trustgraph/extract/kg/ontology/extract.py @@ -282,17 +282,6 @@ class Processor(FlowProcessor): if not ontology_subsets: logger.warning("No relevant ontology elements found for chunk") - # Emit empty outputs - await self.emit_triples( - flow("triples"), - v.metadata, - [] - ) - await self.emit_entity_contexts( - flow("entity-contexts"), - v.metadata, - [] - ) return # Merge subsets if multiple ontologies matched @@ -327,35 +316,26 @@ class Processor(FlowProcessor): entity_contexts = self.build_entity_contexts(all_triples) # Emit all triples (extracted + ontology definitions) - await self.emit_triples( - flow("triples"), - v.metadata, - all_triples - ) + if all_triples: + await self.emit_triples( + flow("triples"), + v.metadata, + all_triples + ) # Emit entity contexts - await self.emit_entity_contexts( - flow("entity-contexts"), - v.metadata, - entity_contexts - ) + if entity_contexts: + await self.emit_entity_contexts( + flow("entity-contexts"), + v.metadata, + entity_contexts + ) logger.info(f"Extracted {len(triples)} content triples + {len(ontology_triples)} ontology triples " f"= {len(all_triples)} total triples and {len(entity_contexts)} entity contexts") except Exception as e: logger.error(f"OntoRAG extraction exception: {e}", exc_info=True) - # Emit empty outputs on error - await self.emit_triples( - flow("triples"), - v.metadata, - [] - ) - await self.emit_entity_contexts( - flow("entity-contexts"), - v.metadata, - [] - ) async def extract_with_simplified_format( self, diff --git a/trustgraph-flow/trustgraph/extract/kg/relationships/extract.py b/trustgraph-flow/trustgraph/extract/kg/relationships/extract.py index be0bd705..c92a33e9 100755 --- a/trustgraph-flow/trustgraph/extract/kg/relationships/extract.py +++ b/trustgraph-flow/trustgraph/extract/kg/relationships/extract.py @@ -181,16 +181,17 @@ class Processor(FlowProcessor): o=Term(type=IRI, iri=v.metadata.id) )) - await self.emit_triples( - flow("triples"), - Metadata( - id=v.metadata.id, - metadata=[], - user=v.metadata.user, - collection=v.metadata.collection, - ), - triples - ) + if triples: + await self.emit_triples( + flow("triples"), + Metadata( + id=v.metadata.id, + metadata=[], + user=v.metadata.user, + collection=v.metadata.collection, + ), + triples + ) except Exception as e: logger.error(f"Relationship extraction exception: {e}", exc_info=True) diff --git a/trustgraph-flow/trustgraph/storage/knowledge/store.py b/trustgraph-flow/trustgraph/storage/knowledge/store.py index a79b7b83..475604b6 100644 --- a/trustgraph-flow/trustgraph/storage/knowledge/store.py +++ b/trustgraph-flow/trustgraph/storage/knowledge/store.py @@ -64,12 +64,14 @@ class Processor(FlowProcessor): async def on_triples(self, msg, consumer, flow): v = msg.value() - await self.table_store.add_triples(v) + if v.triples: + await self.table_store.add_triples(v) async def on_graph_embeddings(self, msg, consumer, flow): v = msg.value() - await self.table_store.add_graph_embeddings(v) + if v.entities: + await self.table_store.add_graph_embeddings(v) @staticmethod def add_args(parser): diff --git a/trustgraph-flow/trustgraph/tables/knowledge.py b/trustgraph-flow/trustgraph/tables/knowledge.py index 0e5d4e1d..6ea16499 100644 --- a/trustgraph-flow/trustgraph/tables/knowledge.py +++ b/trustgraph-flow/trustgraph/tables/knowledge.py @@ -435,14 +435,17 @@ class KnowledgeTableStore: else: metadata = [] - triples = [ - Triple( - s = tuple_to_term(elt[0], elt[1]), - p = tuple_to_term(elt[2], elt[3]), - o = tuple_to_term(elt[4], elt[5]), - ) - for elt in row[3] - ] + if row[3]: + triples = [ + Triple( + s = tuple_to_term(elt[0], elt[1]), + p = tuple_to_term(elt[2], elt[3]), + o = tuple_to_term(elt[4], elt[5]), + ) + for elt in row[3] + ] + else: + triples = [] await receiver( Triples( @@ -491,13 +494,16 @@ class KnowledgeTableStore: else: metadata = [] - entities = [ - EntityEmbeddings( - entity = tuple_to_term(ent[0][0], ent[0][1]), - vectors = ent[1] - ) - for ent in row[3] - ] + if row[3]: + entities = [ + EntityEmbeddings( + entity = tuple_to_term(ent[0][0], ent[0][1]), + vectors = ent[1] + ) + for ent in row[3] + ] + else: + entities = [] await receiver( GraphEmbeddings(