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(