第 10 章
存储与出口
场景还原
数据平台组刚把知识图谱跑通,老板要求在三个环境里各接一个向量库:开发用 SQLite,测试用 pgvector,生产用 Qdrant。工程师照文档写调用,三处都传 top_k=10。结果生产环境直接抛 TypeError: got multiple values for keyword argument 'top_k',测试环境返回的 score 是 0.9,开发环境返回 0.78,同一个查询、同一批向量,排序结果却对不上。
更麻烦的在后面。下游数据治理组要一份 JSON-LD 给数据目录系统,还要一份 Turtle 给语义网工具。导出跑了两次,第二次合并进目录系统时,同一批实体全部重复了一份,因为每次导出都给图节点现造了一个新 IRI,时间戳差几微秒,IRI 就变一次。查了三天才发现:根子在「标识符用墙钟时间铸造」这条暗线上,和数据库选型无关。
这一章拆两个问题:存储层怎么把六种向量后端、四种图后端收敛到一套调用面而不丢掉各自特性;出口层怎么让十几种格式共享同一套图谱数据而保持互操作。两条暗线贯穿:一是「统一结果 schema」,二是「内容寻址标识符」。
逐行精读
VectorStore 门面:先定返回结构,再谈后端
存储层的收敛策略和大多数框架相反:先固定「搜索返回什么」,再让后端适配。SearchResult 这个 TypedDict 就是那份契约:
81class SearchResult(TypedDict):82 """Canonical schema returned by VectorStore.search_vectors().8384 Required fields (always present):85 id – string or integer identifier of the stored vector.86 score – float similarity score, higher is better (normalised to 0.0–1.0 across all backends).87 metadata – dict of associated metadata; empty dict when none is stored.88 vector – np.ndarray when the backend returns the raw vector, otherwise None.89 distance – raw native distance value preserved for backends that expose it90 (FAISS L2, Weaviate cosine), otherwise None.91 """9293 id: Union[str, int]94 score: float95 metadata: Dict[str, Any]96 vector: Optional[Any] # np.ndarray | None97 distance: Optional[float]这张契约解决的是「六种后端各自返回不同键名」的问题。score 被规定为越大越相似、归一化到 0 到 1,这是给上层排序用的统一口径;distance 保留后端原生的原始距离值,谁需要精确的 L2 或 cosine 距离就拿它。注释里写明了归一化是「across all backends」的承诺,这意味着每个后端的适配器都要负责把自家度量转成这个尺度。
门面本身靠一张白名单挡住拼写错误:
100class VectorStore:101 """102 Vector store interface and management.103104 • Stores and manages vector embeddings105 • Provides similarity search capabilities106 • Handles vector indexing and retrieval107 • Manages vector metadata and provenance108 • Supports multiple vector store backends109 • Provides vector store operations110 """111112 SUPPORTED_BACKENDS = {"faiss", "weaviate", "qdrant", "milvus", "pinecone", "pgvector", "inmemory", "sqlite"}SUPPORTED_BACKENDS 是个集合字面量,八个值。__init__ 第一件事就是把传入的 backend 做 lower() 再查这个集合,不在就直接 ValueError。集合用 sorted() 展开成提示语,报错时用户能看到全部合法选项。注意 inmemory 也在集合里,它是「不接任何外部服务」的本地兜底,后面所有后端归一化的逻辑都以它为参照。
六种后端的差异被压缩在 _init_backend_store 的一个 if/elif 链里,每个分支做两件事:从配置文件读参数、实例化对应适配器。pgvector 分支额外检查 connection_string 是否缺失,缺失就抛带示例的 ValueError;sqlite 分支检查 db_path。这套「参数校验前置」让错误发生在构造时,把失败提前到任何查询之前。
真正考验设计的是写和查两个方法。store_vectors 的委托逻辑:
507 # Delegate to backend store if available508 if self._backend_store:509 # Handle different method names across backend stores510 if hasattr(self._backend_store, 'add'):511 return self._backend_store.add(vectors, metadata, **options)512 elif hasattr(self._backend_store, 'add_vectors'):513 # FAISSStore and others use add_vectors with different signature514 if hasattr(self._backend_store, 'store_vectors'):515 # Some stores have store_vectors method516 return self._backend_store.store_vectors(vectors, metadata=metadata, **options)517 else:518 try:519 add_vectors_params = inspect.signature(self._backend_store.add_vectors).parameters520 supports_metadata = 'metadata' in add_vectors_params or any(521 p.kind == inspect.Parameter.VAR_KEYWORD for p in add_vectors_params.values()522 )523 except (ValueError, TypeError):524 supports_metadata = True525 if supports_metadata:526 return self._backend_store.add_vectors(vectors, metadata=metadata, **options)527 return self._backend_store.add_vectors(vectors, **options)528 else:529 raise NotImplementedError(f"Backend store {type(self._backend_store).__name__} does not have add or add_vectors method")这里没有强制后端实现统一接口,而是用 hasattr 运行时探测方法名。有的后端叫 add,有的叫 add_vectors,有的叫 store_vectors,逐个试。最精细的一处是 inspect.signature:当只有 add_vectors 可用时,它反射这个方法的参数表,判断是否接受 metadata 关键字,接受就传,不接受就省掉,避免给 FAISS 这类老接口传多余参数炸掉。这条链的终点是 NotImplementedError,让「这个后端不支持写入」以异常形式暴露,绝不静默吞掉。
搜索侧同款逻辑,但多了一个针对性注释:
696 # Delegate to backend store if available697 if self._backend_store:698 # Handle different method names across backend stores699 if hasattr(self._backend_store, 'search'):700 return self._backend_store.search(query_vector, top_k=k, **options)701 elif hasattr(self._backend_store, 'search_similar'):702 return self._backend_store.search_similar(query_vector, k=k, **options)703 elif hasattr(self._backend_store, 'search_vectors'):704 # Some stores (QdrantStore, MilvusStore, PineconeStore) name705 # their count parameter differently (limit vs k), so bind it706 # positionally rather than guessing the keyword.707 return self._backend_store.search_vectors(query_vector, k, **options)708 else:709 raise NotImplementedError(f"Backend store {type(self._backend_store).__name__} does not have search, search_similar, or search_vectors method")top_k、k、limit 三种命名在不同后端里打架。注释给了解法:对有 search_vectors 方法的后端,数量参数按位置传,不猜关键字。这就是场景还原里 top_k 冲突的反面教材的来源,门面这一层用「位置绑定」绕开了命名差异。
count() 则是「能力缺失要显式」的样本:
827 def count(self) -> int:828 """Return the number of vectors in the store, backend-agnostic.829830 The inmemory backend counts its local dict; persistent backends831 delegate to a ``count()`` on the wrapped store when available.832 Following the get_vector()/get_metadata() precedent (#843) and the833 NotImplementedError-on-unsupported-capability precedent of834 _filter_by_metadata() (#848), a persistent backend that cannot835 report a count raises NotImplementedError so callers can tell836 "no vectors" apart from "counting not supported" — including when837 the wrapped backend store is missing entirely (never silently838 report an uninitialized store as empty).839 """840 if self.backend == "inmemory":841 return len(self.vectors)842 elif self._backend_store is not None:843 count_attr = getattr(self._backend_store, "count", None)844 if callable(count_attr):845 return cast(int, count_attr())846 raise NotImplementedError(847 f"Backend store {type(self._backend_store).__name__} does not "848 "implement a count() method. Add a count() method to the "849 "backend store adapter to enable vector counting for this backend."850 )851 raise NotImplementedError(852 f"Backend store is not initialized; cannot count vectors for "853 f"backend {self.backend!r}."854 )docstring 把设计意图写得比代码还长:count() 有四种结局,本地字典直接数、后端有 count 就委托、后端没实现就 NotImplementedError、连后端实例都没有也 NotImplementedError。区分「没有向量」和「不支持计数」这两件事,靠的就是不返回 0 当哑值。注释里 (#843)、(#848) 是 issue 编号,说明这套语义是从前两个方法的教训里沿用下来的。
PgVectorStore:一个具体后端要扛多少事
门面省掉的脏活,全落到适配器里。pgvector 是最有代表性的一个,因为它要处理三种距离度量、SQL 注入、扩展缺失、连接池四类问题。
距离度量和 SQL 运算符的映射在 search 里:
374 # Determine operator based on distance metric375 operator_map = {376 "cosine": "<=>", # cosine distance377 "l2": "<->", # L2 distance378 "inner_product": "<#>", # negative inner product (for ordering)379 }380 op = operator_map[self.distance_metric]pgvector 的三种度量对应三个运算符,且 <=> 和 <-> 返回的是距离,越小越好;<#> 返回的是负内积,排序方向相反。这套差异在下游被折成统一分数:
437 # Convert distance to similarity score438 if self.distance_metric in ("cosine", "l2"):439 # Lower distance = higher similarity440 similarity = 1.0 / (1.0 + float(score))441 else: # inner_product442 # Negative because we used <#> operator443 similarity = -float(score)444445 results.append({446 "id": vec_id,447 "score": similarity,448 "metadata": meta if isinstance(meta, dict) else json.loads(meta),449 "vector": None,450 "distance": None,451 })cosine 和 l2 用 1/(1+距离) 把「距离越小越好」翻成「分数越大越好」,顺便把范围压进 0 到 1;内积因为运算符自带负号,取反就是真内积。结果字典里 vector 和 distance 都填 None,因为 SQL 查询默认不把原始向量捞回来,原生距离值在换算后也不再单独暴露。这就是 SearchResult 契约里那两个 Optional 字段的实际落点。
元数据过滤里的注入防线:
459 def _is_safe_identifier(self, key: str) -> bool:460 """461 Validate that a string is safe to use as a SQL/jsonb identifier.462463 Only allows alphanumeric characters, underscores, and hyphens.464 Rejects any string that could be used for SQL injection.465 """466 if not isinstance(key, str):467 return False468 if not key:469 return False470 # Only allow: alphanumeric, underscore, hyphen471 # Must start with letter or underscore472 return bool(re.match(r'^[a-zA-Z_][a-zA-Z0-9_-]*$', key))因为过滤器键会被拼进 metadata->>{} 的 SQL 片段,任何没被正则放行的键都会在 search 和 filter_by_metadata 里提前抛 ValidationError。这条正则允许字母数字下划线连字符,且首字符只能是字母或下划线,堵住了「键名里塞一段 SQL」的口子。表名和维度则走 psycopg_sql.SQL 的占位组合,双保险。
混合检索:RRF 落在哪一层
RRF 不在 VectorStore 门面里,而在 hybrid_search.py 的 SearchRanker。这是「融合」和「检索」分层的证据:门面只管单一向量源的读写,多路结果的合并是独立组件。
140class SearchRanker:141 """Search result ranker."""142143 def __init__(self, strategy: str = "reciprocal_rank_fusion"):144 """Initialize search ranker."""145 self.strategy = strategy146 self.logger = get_logger("search_ranker")147148 def reciprocal_rank_fusion(149 self, results: List[List[Dict[str, Any]]], k: int = 60150 ) -> List[Dict[str, Any]]:151 """152 Reciprocal Rank Fusion (RRF) algorithm.153154 Args:155 results: List of result lists from different sources156 k: RRF constant157158 Returns:159 Fused and ranked results160 """161 scores: Dict[str, float] = {}162163 for result_list in results:164 for rank, result in enumerate(result_list, start=1):165 result_id = result.get("id", str(id(result)))166 score = 1.0 / (k + rank)167 scores[result_id] = scores.get(result_id, 0.0) + score168169 # Sort by score170 ranked = sorted(scores.items(), key=lambda x: x[1], reverse=True)171172 # Reconstruct results173 fused_results = []174 result_map = {}175 for result_list in results:176 for result in result_list:177 result_id = result.get("id", str(id(result)))178 result_map[result_id] = result179180 for result_id, score in ranked:181 if result_id in result_map:182 result = result_map[result_id].copy()183 result["score"] = score184 fused_results.append(result)185186 return fused_results187188 def weighted_average(189 self, results: List[List[Dict[str, Any]]], weights: List[float]190 ) -> List[Dict[str, Any]]:191 """192 Weighted average fusion.193194 Args:195 results: List of result lists196 weights: Weights for each result list197198 Returns:199 Fused results200 """201 if len(weights) != len(results):202 weights = [1.0 / len(results)] * len(results)203204 scores: Dict[str, Tuple[float, Dict[str, Any]]] = {}205206 for weight, result_list in zip(weights, results):207 for result in result_list:208 result_id = result.get("id", str(id(result)))209 score = result.get("score", 0.0) * weight210211 if result_id not in scores:212 scores[result_id] = (0.0, result)213214 scores[result_id] = (scores[result_id][0] + score, scores[result_id][1])215216 # Sort by score217 ranked = sorted(scores.values(), key=lambda x: x[0], reverse=True)218219 fused_results = []220 for score, result in ranked:221 result_copy = result.copy()222 result_copy["score"] = score223 fused_results.append(result_copy)224225 return fused_results226227 def rank(228 self, results: List[List[Dict[str, Any]]], **options229 ) -> List[Dict[str, Any]]:230 """231 Rank and fuse results.232233 Args:234 results: List of result lists235 **options: Ranking options236237 Returns:238 Fused and ranked results239 """240 if self.strategy == "reciprocal_rank_fusion":241 k = options.get("k", 60)242 return self.reciprocal_rank_fusion(results, k)243 elif self.strategy == "weighted_average":244 weights = options.get("weights", [1.0 / len(results)] * len(results))245 return self.weighted_average(results, weights)246 else:247 return self.reciprocal_rank_fusion(results)RRF 的核心是 score = 1.0 / (k + rank):每个来源按名次给分,同名次累加,跨来源的结果就按「在多少条链里排多靠前」重新排序。result_map 的二次遍历是为了把分数写回原始结果字典,同时保留原始内容。weighted_average 是另一种融合策略,按权重乘分数再求和,适合来源本身分数可比、只是重要性不同的场景。rank 是统一入口,按构造时选的 strategy 分派。
GraphStore:图后端的同款门面,但收敛方式不同
图后端数量比向量少,但门面结构几乎照搬。差异在于:向量门面靠 hasattr 鸭子类型探测,图门面靠 if/elif 显式分派,因为图数据库的操作面更宽,靠反射猜方法名会失控。
525class GraphStore:526 """527 Main graph store interface.528529 Provides a unified interface for working with property graph databases,530 supporting Neo4j and FalkorDB backends.531 """532533 def __init__(534 self,535 backend: Optional[str] = None,536 **config,537 ):538 """539 Initialize graph store.540541 Args:542 backend: Backend type ("neo4j", "falkordb")543 **config: Backend-specific configuration544 """545 self.logger = get_logger("graph_store")546 self.progress_tracker = get_progress_tracker()547 # Ensure progress tracker is enabled548 if not self.progress_tracker.enabled:549 self.progress_tracker.enabled = True550551 # Determine backend552 self.backend = (553 backend554 or config.get("backend")555 or graph_store_config.get("default_backend", "neo4j")556 )557 self.config = config558559 # Initialize store backend560 self._store_backend = None561 self._manager = None562 self._initialize_store_backend()后端选择有三层回落:显式 backend 参数、config 字典里的 backend 键、全局配置的 default_backend,最后兜底 neo4j。这和向量门面「不传就默认 faiss」是同一套哲学,但多了一层 config 透传,因为图库的连接参数更多。
564 def _initialize_store_backend(self) -> None:565 """Initialize the appropriate store backend based on backend type."""566 if self.backend == "neo4j":567 from .neo4j_store import Neo4jStore568569 neo4j_config = graph_store_config.get_neo4j_config()570 neo4j_config.update(self.config)571 self._store_backend = Neo4jStore(**neo4j_config)572573 elif self.backend == "falkordb":574 from .falkordb_store import FalkorDBStore575576 falkordb_config = graph_store_config.get_falkordb_config()577 falkordb_config.update(self.config)578 self._store_backend = FalkorDBStore(**falkordb_config)579580 elif self.backend == "neptune" or self.backend == "amazon_neptune":581 from .amazon_neptune import AmazonNeptuneStore582583 neptune_config = graph_store_config.get_neptune_config()584 neptune_config.update(self.config)585 self._store_backend = AmazonNeptuneStore(**neptune_config)586587 elif self.backend == "age" or self.backend == "apache_age":588 from .age_store import ApacheAgeStore589590 age_config = graph_store_config.get_age_config()591 age_config.update(self.config)592 self._store_backend = ApacheAgeStore(**age_config)593594 else:595 raise ValidationError(f"Unknown backend: {self.backend}")596597 self._manager = GraphManager(self._store_backend)四种后端,每个分支都是「读模块级配置、用 self.config 覆盖、实例化、交给 GraphManager」四步。self.config 覆盖全局配置,意味着调用方在构造时传的参数优先级最高。neptune 和 amazon_neptune、age 和 apache_age 的别名合并在一个分支,用 or 连接,说明后端命名允许两种叫法。未知后端走 ValidationError,报错时没有「支持列表」的提示,这一点比向量门面弱。
QueryEngine 是图门面里最有价值的一个独立类,完整定义:
240class QueryEngine:241 """Engine for query execution and optimization."""242243 def __init__(self, backend: Any):244 """245 Initialize query engine.246247 Args:248 backend: Graph database backend instance249 """250 self.backend = backend251 self.logger = get_logger("query_engine")252 self._cache: Dict[str, Any] = {}253 self._cache_enabled = True254255 def execute(256 self,257 query: str,258 parameters: Optional[Dict[str, Any]] = None,259 use_cache: bool = False,260 **options,261 ) -> Dict[str, Any]:262 """263 Execute a Cypher/OpenCypher query.264265 Args:266 query: Query string267 parameters: Query parameters268 use_cache: Whether to use query caching269 **options: Additional options270271 Returns:272 Query results273 """274 # Check cache275 if use_cache and self._cache_enabled:276 cache_key = self._generate_cache_key(query, parameters)277 if cache_key in self._cache:278 return self._cache[cache_key]279280 # Execute query281 result = self.backend.execute_query(query, parameters, **options)282283 # Cache result284 if use_cache and self._cache_enabled:285 self._cache[cache_key] = result286287 return result288289 def _generate_cache_key(290 self,291 query: str,292 parameters: Optional[Dict[str, Any]],293 ) -> str:294 """Generate cache key for query."""295 import hashlib296297 key_str = f"{query}:{str(parameters)}"298 return hashlib.md5(key_str.encode()).hexdigest() # nosec B324 - cache key, not security-sensitive299300 def clear_cache(self) -> None:301 """Clear query cache."""302 self._cache.clear()303304 def enable_cache(self) -> None:305 """Enable query caching."""306 self._cache_enabled = True307308 def disable_cache(self) -> None:309 """Disable query caching."""310 self._cache_enabled = False缓存键是「查询文本 + 参数字符串」的 md5,注释标了 nosec B324,明确说这不是安全敏感场景,只是缓存键。缓存默认开启但查询默认不启用缓存(use_cache=False),需要调用方显式传参,因为 Cypher 查询经常带写操作,无脑缓存会返回过期数据。这个类的存在说明:图门面不只是一个转发器,它在转发之上叠加了缓存这个横切能力。
出口:三种代表性格式的互操作手法
导出层有 19 个文件,但真正决定「互操作」质量的是三个手法:IRI 铸造、@context 管理、Cypher 字符串转义。
RDF 侧先看置信度数据类型这一处注释,它解释了为什么一个数值要费这么大劲:
53#: The one datatype every serializer writes confidence in.54#:55#: The four paths used to disagree: Turtle wrote the value bare, which the56#: Turtle grammar reads as xsd:decimal, N-Triples typed it xsd:float, RDF/XML57#: emitted a plain literal with no datatype, and JSON-LD emitted a native58#: number, which becomes xsd:double. Those are four distinct RDF terms for one59#: value (issue #1100).60#:61#: xsd:decimal is the choice because it is what the Turtle path already62#: produced, so the most used output is unchanged, and because it is exact:63#: xsd:float is 32 bit binary, and cannot represent 0.9 or 0.95 at all.64CONFIDENCE_DATATYPE = "http://www.w3.org/2001/XMLSchema#decimal"同一个 confidence 值,四种序列化器曾经写出四个不同的 RDF 术语,合进同一个图里就是四个不同的三元组。统一成 xsd:decimal 的理由写在注释里:和 Turtle 既有输出一致,且十进制能精确表示 0.9、0.95 这类值,而 xsd:float 的 32 位二进制表示不了。这是「跨格式互操作」最具体的代价:一个数值的合法表示方式不止一种,选错就静默制造重复。
实体没有 id 时的 IRI 铸造:
123def mint_entity_iri(text: str) -> str:124 """Mint a stable IRI for an entity that arrived without an id.125126 Python's builtin ``hash()`` is randomised per process (PYTHONHASHSEED), so127 minting from it gave the same entity a different IRI on every run: exports128 could not be diffed, deduplicated against an earlier load, or joined to a129 provenance record written by an earlier process. SHA-256 is stable across130 runs and machines, which is what an identifier has to be.131 """132 digest = hash_data(str(text))[:16]133 return f"{SEMANTICA_NS}entity_{digest}"用 Python 内建 hash() 会因为 PYTHONHASHSEED 每个进程不同而产出不同 IRI,导致两次导出不可 diff、不可去重、接不上之前的溯源记录。改用 SHA-256 后,同样的文本在任何机器任何进程都得到同样的 IRI,这就是「内容寻址标识符」:标识符由内容决定,不由时间或随机数决定。
Turtle 序列化把实体转成三元组子句:
833 entity_type = entity.get("type", "semantica:Entity")834 lines.append(835 f"{subject} <http://www.w3.org/1999/02/22-rdf-syntax-ns#type> {expand_uri(entity_type)} ."836 )837838 # Text property839 text = entity.get("text") or entity.get("label", "")840 if text:841 safe_text = text.replace('"', '\\"').replace("\n", "\\n")842 lines.append(843 f'{subject} {expand_uri("semantica:text")} "{safe_text}" .'844 )845846 # Confidence property. The default matches the other serializers,847 # which have always written one; omitting it here was half of why848 # Turtle and N-Triples of one KG were different graphs (#1100).849 raw_confidence = entity.get("confidence", 1.0)850 confidence = normalize_confidence(raw_confidence)851 if confidence is None:852 self.logger.warning(853 f"Entity {entity.get('id')} has a confidence that is not a "854 f"number ({raw_confidence!r}), so no confidence is written"855 )856 else:857 lines.append(858 f'{subject} {expand_uri("semantica:confidence")} '859 f'"{confidence}"^^<{CONFIDENCE_DATATYPE}> .'860 )子句先攒成列表,a <类型> 是 rdf:type 的 Turtle 简写,semantica:text 是文本字面量,confidence 用 ^^<...> 显式标上数据类型,metadata 通过 _metadata_statements 展开成额外的谓词-宾语对。最后按 Turtle 的 ; 分句、. 收尾规则拼接。clauses[1:-1] 的切片处理了「只有一个子句」的边界,首句和末句分别用 ; 和 . 收尾。
JSON-LD 侧的图节点 IRI 也走了内容寻址:
37def _content_iri(prefix: str, payload: Any) -> str:38 """Mint a document IRI from what was exported, not when.3940 Minting from ``utc_now_iso()`` gave every export of the same graph a new41 identity a few microseconds apart, so re-exporting an unchanged graph was42 never idempotent and merging exports duplicated every node (#1147). This43 mirrors ``mint_entity_iri`` (#1109): identical content hashes to the same44 IRI, and any change to the content changes it too. ``default=str`` keeps45 the hash defined for values ``json.dumps`` would otherwise reject, such as46 ``datetime`` objects a caller may have left in the graph.4748 Args:49 prefix: IRI prefix the digest is appended to50 payload: JSON-serializable value whose content determines the digest5152 Returns:53 A stable IRI of the form ``{prefix}{16-hex-char digest}``54 """55 canonical = json.dumps(payload, sort_keys=True, default=str)56 digest = hash_data(canonical)[:16]57 return f"{prefix}{digest}"docstring 第一句就点破:IRI 从「导出了什么」铸造,从「何时导出」铸造。sort_keys=True 保证键顺序不影响哈希,default=str 兜住 datetime 这类 json.dumps 会拒绝的值。这就是场景还原里「合并重复实体」的正解:同一张图两次导出的 IRI 一致,合并时自然去重。
图节点用这个函数命名:
609 # Minted from the graph's own content rather than the wall clock610 # (#1147): re-exporting an unchanged graph must produce the same611 # subject, or merging repeated exports duplicates every node.612 "@id": options.get("graph_uri")613 or _content_iri("https://semantica.dev/graph/", kg),614 "@type": "semantica:KnowledgeGraph",调用方可以传 graph_uri 覆盖,传了就用自己的,没传就用内容哈希铸一个,两者取或。默认路径保证幂等,覆盖路径保留灵活性。
Cypher 侧看转义,这是 LPG 导出最容易翻车的地方:
229 def _create_nodes_batch(self, nodes: List[Dict[str, Any]]) -> str:230 """Create a batch of nodes in a single Cypher query."""231 if not nodes:232 return ""233234 lines = []235 for idx, node in enumerate(nodes):236 node_id = node.get("id") or node.get("entity_id", f"node_{idx}")237 node_type = node.get("type") or node.get("entity_type", "Entity")238 label = node.get("label") or node.get("name") or node.get("text", "")239240 # Escape special characters241 node_id_escaped = self._escape_cypher_string(str(node_id))242 label_escaped = self._escape_cypher_string(str(label))243244 # Build properties245 properties = {"id": node_id_escaped, "name": label_escaped}246247 # Add other properties248 for key, value in node.items():249 if key not in [250 "id",251 "entity_id",252 "type",253 "entity_type",254 "label",255 "name",256 "text",257 ]:258 if isinstance(value, (str, int, float, bool)):259 properties[key] = self._format_cypher_value(value)260261 # Format properties262 props_str = ", ".join([f"{k}: {v}" for k, v in properties.items()])263264 # Create node265 lines.append(f"CREATE (n{idx}:{node_type} {{{props_str}}});")266267 return "\n".join(lines)id、name 两个固定属性先行转义,其余标量属性过滤后逐个 _format_cypher_value,非标量(dict、list)直接丢弃,因为 Cypher 属性值只收标量。类型直接当节点标签用,不做映射。这里有个隐患会在边界条件里展开:node_type 没转义,如果实体类型名里带空格或特殊字符,生成的 Cypher 就是非法语句。
关系侧:
281 def _create_relationships_batch(self, edges: List[Dict[str, Any]]) -> str:282 """Create a batch of relationships in a single Cypher query."""283 if not edges:284 return ""285286 lines = []287 for idx, edge in enumerate(edges):288 source_id = edge.get("source") or edge.get("source_id")289 target_id = edge.get("target") or edge.get("target_id")290 rel_type = edge.get("type") or edge.get("relationship_type", "RELATED_TO")291292 if not source_id or not target_id:293 continue294295 # Escape special characters296 source_escaped = self._escape_cypher_string(str(source_id))297 target_escaped = self._escape_cypher_string(str(target_id))298 rel_type_escaped = rel_type.upper().replace(" ", "_")299300 # Build properties301 properties = {}302 for key, value in edge.items():303 if key not in [304 "source",305 "source_id",306 "target",307 "target_id",308 "type",309 "relationship_type",310 ]:311 if isinstance(value, (str, int, float, bool)):312 properties[key] = self._format_cypher_value(value)313314 # Format properties315 if properties:316 props_str = ", ".join([f"{k}: {v}" for k, v in properties.items()])317 props_str = f" {{{props_str}}}"318 else:319 props_str = ""320321 # Create relationship322 lines.append(323 f"MATCH (a {{id: '{source_escaped}'}}), (b {{id: '{target_escaped}'}}) "324 f"CREATE (a)-[r:{rel_type_escaped}{props_str}]->(b);"325 )326327 return "\n".join(lines)关系用 MATCH (a {id: ...}), (b {id: ...}) 定位两端节点,再用 CREATE (a)-[r:类型 {属性}]->(b) 连边。源和目标 id 转义了,但关系类型只做了 upper() 加空格换下划线,和节点标签有同样的未转义隐患。if not source_id or not target_id: continue 是静默跳过,坏边不报错,只在结果里少一条。
转义函数本身:
329 def _escape_cypher_string(self, value: str) -> str:330 """Escape special characters in Cypher strings."""331 return value.replace("'", "\\'").replace("\\", "\\\\")332333 def _format_cypher_value(self, value: Any) -> str:334 """Format a value for Cypher query."""335 if isinstance(value, str):336 escaped = self._escape_cypher_string(value)337 return f"'{escaped}'"338 elif isinstance(value, bool):339 return str(value).lower()340 else:341 return str(value)单引号和反斜杠被转义,字符串包单引号,布尔转小写,其余直接 str()。这套转义覆盖了「值里带引号」的常见翻车,但覆盖不了标签和关系类型的注入。
最后看出口层的总调度:
924 # Auto-detect format from file extension if not specified925 if not format:926 file_path_obj = Path(file_path)927 ext = file_path_obj.suffix.lower()928 format_map = {929 ".json": "json",930 ".jsonld": "json-ld",931 ".csv": "csv",932 ".ttl": "turtle",933 ".rdf": "rdfxml",934 ".graphml": "graphml",935 ".gexf": "gexf",936 ".dot": "dot",937 ".yaml": "yaml",938 ".yml": "yaml",939 ".owl": "owl-xml",940 ".cypher": "cypher",941 ".aql": "aql",942 ".parquet": "parquet",943 }944 format = format_map.get(ext, "json")945946 # Route to appropriate exporter947 if format in ["json", "json-ld"]:948 export_json(knowledge_graph, file_path, format=format, method=method, **kwargs)949 elif format == "csv":950 export_csv(knowledge_graph, file_path, method=method, **kwargs)951 elif format in ["turtle", "rdfxml", "jsonld", "ntriples", "n3"]:952 export_rdf(knowledge_graph, file_path, format=format, method=method, **kwargs)953 elif format in ["graphml", "gexf", "dot"]:15 个扩展名映射到 11 个导出族,每个族再派给对应的 export_* 函数。.ttl 走 RDF 的 turtle,.cypher 走 LPG,.jsonld 走 JSON-LD。没传 format 时从扩展名猜,猜不到就默认 json。这就是「十几种导出格式」的统一入口:一个函数、一张表、一条 if/elif 链。
设计决策分析
存储层最核心的决策是「鸭子类型收口,不强制统一接口」。Semantica 没有给后端定义一个 AbstractVectorStore 抽象基类,而是让 VectorStore 用 hasattr 逐个探测方法名。这样做的代价是每次委托都要写一条 if/elif 链,好处是新后端只要实现了「长得像」的方法就能挂进来,不需要继承任何东西。docs/architecture.md 第 152 行把这条哲学写成了原则:
152Every component works standalone. `NERExtractor` runs without a graph store. `VectorStore` runs without decision tracking. The framework never forces a full stack instantiation: you pay only for what you import.「pay only for what you import」是这句话的落点。换后端时调用方一行代码不改,这句也写进了性能特性表:
184| **Backend flexibility** | Swap in-memory NetworkX for Neo4j / FalkorDB with no API changes |「with no API changes」是设计目标,也是衡量门面成败的标准。向量门面和图门面各自用不同手段达成:向量用鸭子类型探测,图用显式分派。
出口层的关键决策是「内容寻址标识符」。mint_entity_iri 和 _content_iri 都用 SHA-256 哈希内容的规范序列化,前者哈希实体文本,后者哈希整个图结构的 json.dumps(sort_keys=True)。这个决策的直接收益是幂等:同一张图导出两次,IRI 相同,下游合并不会重复。代价是任何内容变化都会改变 IRI,包括 metadata 里一个无关字段的改动,这会让「基于 IRI 的增量更新」误判为全新节点。注释里没有掩盖这个取舍,_content_iri 的 docstring 明确写了「any change to the content changes it too」。
第三个决策是「元数据键到 RDF 术语的显式映射,宁可丢不可猜」。DEFAULT_METADATA_TERMS 表只认九个 Semantica 自己产出的键,调用方自定义的键会被警告并跳过,绝不瞎猜命名空间。这是对「互操作」的保守解:写一个错命名空间的三元组比不写更糟,因为下游解析器会把它当成合法但语义错误的数据。
边界条件剖析
如果后端名拼错了会怎样。 VectorStore.__init__ 第一行就把 backend.lower() 拿去查 SUPPORTED_BACKENDS 集合,不在集合里直接抛 ValueError,消息里带 sorted() 展开的八个合法值。这是第 116 到 121 行的逻辑,发生在任何后端实例化之前,所以拼错 qdrnt 连一行网络请求都不会发出。图门面则相反,_initialize_store_backend 的 else 分支只抛 ValidationError(f"Unknown backend: {self.backend}"),不列支持列表,落在第 592 到 593 行,报错信息比向量侧弱。
如果 pgvector 表没建、扩展没装会怎样。 PgVectorStore.__init__ 依次调 _init_pool、_verify_pgvector_extension、_ensure_table_exists。扩展没装时,_verify_pgvector_extension 查 pg_extension 表发现没有 vector 行,抛 ProcessingError 并附安装链接,这是第 210 到 227 行。扩展在但表不存在时,_ensure_table_exists 用 CREATE TABLE IF NOT EXISTS 建表,幂等,不会报错。连接串缺失则更早,在 VectorStore._init_backend_store 的 pgvector 分支抛带示例的 ValueError。
如果查询时传了危险过滤键会怎样。 search 遍历 filter 时对每个键调 _is_safe_identifier,键里只要有一个字符不在 [a-zA-Z0-9_-] 范围内,就抛 ValidationError,这是第 390 到 395 行的校验。a; DROP TABLE vectors;-- 这样的键过不了正则,到不了 SQL 拼接那一步。
如果后端不支持 count 会怎样。 count() 先看 self.backend == "inmemory",是就数本地字典;再看 self._backend_store 是否有可调用的 count 属性,有就委托,没有就抛 NotImplementedError,这是第 840 到 850 行。关键在「连后端实例都没有」的分支,第 851 到 854 行也抛 NotImplementedError,把「未初始化」和「不支持」区分开,绝不返回 0 冒充空库。
横向对比
对象是 GraphRAG 的存储层。两边都要回答「多个存储后端怎么统一」,但答案形态完全不同。
Semantica 的 VectorStore 是一个具体类,用 hasattr 探测方法名收口(第 508 到 529 行);GraphRAG 的 VectorStore 是一个抽象基类,把接口钉死在抽象方法上:
56class VectorStore(ABC):57 """The base class for vector storage data-access classes."""5859 def __init__(60 self,61 index_name: str = "vector_index",62 id_field: str = "id",63 vector_field: str = "vector",64 create_date_field: str = "create_date",65 update_date_field: str = "update_date",66 vector_size: int = 3072,GraphRAG 的每个后端必须继承这个 ABC,实现 connect、create_index、load_documents、similarity_search_by_vector 等抽象方法。工厂用 match 语句按枚举类型注册:
73 # Lazy load built-in implementations74 if strategy not in vector_store_factory:75 match strategy:76 case VectorStoreType.LanceDB:77 from graphrag_vectors.lancedb import LanceDBVectorStore7879 register_vector_store(VectorStoreType.LanceDB, LanceDBVectorStore)80 case VectorStoreType.AzureAISearch:81 from graphrag_vectors.azure_ai_search import AzureAISearchVectorStore8283 register_vector_store(84 VectorStoreType.AzureAISearch, AzureAISearchVectorStore85 )86 case VectorStoreType.CosmosDB:87 from graphrag_vectors.cosmosdb import CosmosDBVectorStore8889 register_vector_store(VectorStoreType.CosmosDB, CosmosDBVectorStore)90 case _:91 msg = f"Vector store type '{strategy}' is not registered in the VectorStoreFactory. Registered types: {', '.join(vector_store_factory.keys())}."92 raise ValueError(msg)两种收口方式的取舍很清晰:Semantica 的鸭子类型让新后端零继承成本挂进来,代价是每处委托都要手写探测链,且拼错方法名在运行时才暴露;GraphRAG 的抽象基类让拼错在类定义时就被 abstractmethod 抓住,代价是每个后端都要完整实现整个接口面,哪怕只用其中一半。
分数语义的差异更直观。Semantica 在 pgvector 适配器里把距离换算成 0 到 1 的相似度(第 437 到 443 行),门面承诺 score 跨后端可比;GraphRAG 的 LanceDB 适配器把分数算成 1 - abs(distance),注释没有,语义直接丢给调用方:
216 VectorStoreSearchResult(217 document=VectorStoreDocument(218 id=doc[self.id_field],219 vector=doc[self.vector_field] if include_vectors else None,220 data=self._extract_data(doc, select),221 create_date=doc.get(self.create_date_field),222 update_date=doc.get(self.update_date_field),223 ),224 score=1 - abs(float(doc["_distance"])),225 )1 - abs(distance) 和 Semantica 的 1/(1+distance) 数值不等,但两边都没有在文档层面统一「score 必须落在哪个区间」。Semantica 至少把承诺写进了 SearchResult 的 docstring,GraphRAG 的 VectorStoreSearchResult 只写了 score: float,注释「Similarity score between -1 and 1」是唯一约束,而 LanceDB 这条路径产出的值其实可以落在别的区间,承诺与实现之间存在裂缝。
存储格式的对比更根本。GraphRAG 的 Storage 抽象处理的是「键值文件存储」:
13class Storage(ABC):14 """Provide a storage interface."""1516 @abstractmethod17 def __init__(self, **kwargs: Any) -> None:18 """Create a storage instance."""1920 @abstractmethod21 def find(22 self,23 file_pattern: re.Pattern[str],24 ) -> Iterator[str]:25 """Find files in the storage using a file pattern.它管的是 parquet、csv、lance 这类文件的后端(File、Memory、AzureBlob、AzureCosmos);图数据库与向量数据库的连接不归这层管。GraphRAG 的整条索引产物天然落在 parquet 表和 LanceDB 向量库上,所以它不需要一个「把图导出成 RDF/JSON-LD/Cypher」的出口层:存储格式本身就是出口,下游直接用 parquet/lance 读。
检索证据:在 graphrag 仓库的 packages/ 下用 rdf、jsonld、json-ld、turtle、owl 作关键词搜 Python 文件,除 knowledge、graphrag 这类误命中外无结果;find packages -iname '*export*' 无文件。GraphRAG 没有出口层,它用 graphrag-storage 的 ParquetTableProvider、CsvTableProvider 和 graphrag-vectors 的 LanceDB 顶替了这个职责。Semantica 为什么需要出口层:它把知识图谱当成独立资产,要交付给数据目录、语义网工具、图数据库三套不同的下游,三者读取格式互不相通,于是 19 个文件的 export/ 模块存在。
互动演示设计
一句话结论: 存储层的门面靠「统一结果结构 + 运行时探测方法名」收敛后端差异,出口层靠「内容寻址 IRI + 显式术语映射」保证幂等与互操作,两个层都在用「宁可显式失败,不静默猜」兜底。
舞台元素与比喻: 把 VectorStore 想象成一个酒店前台。客人(调用方)只说「存」「查」,前台(门面)手里有一本《各分店服务名对照表》(hasattr 探测链),不管客人住的是 FAISS 店还是 pgvector 店,前台都回一张同样格式的回执(SearchResult)。出口层是酒店的行李托运处,同一件行李(kg 字典)要贴三张不同的标签(Turtle、JSON-LD、Cypher),每张贴法不同,但行李编号(IRI)都由行李内容本身生成。
分步动画:
第一步,客人递来「存三份向量」,前台先翻白名单确认分店名合法。字幕:「backend 不在 SUPPORTED_BACKENDS,第 116 行直接拒绝,连门都不出。」
第二步,前台翻对照表找分店的服务名,逐个 hasattr 试。字幕:「是 add 还是 add_vectors?第 509 到 528 行逐个探测,最后一个都不中就抛 NotImplementedError。」
第三步,分店返回结果,前台按统一格式整理后递回。字幕:「score 归一化到 0 到 1,原生 distance 原样保留,第 437 到 451 行在 pgvector 侧完成换算。」
第四步,客人要三份不同格式的导出,行李托运处按扩展名分派。字幕:「.ttl 走 RDF,.jsonld 走 JSON-LD,.cypher 走 LPG,第 928 到 948 行的 format_map 一张表定去向。」
第五步,行李编号由内容哈希生成,两次托运同一件行李编号相同。字幕:「_content_iri 用 json.dumps(sort_keys=True) 的 SHA-256 铸 IRI,第 55 到 56 行,墙钟时间不参与。」
读者可操作项: 把 VectorStore(backend="qdrnt") 跑一遍,观察 ValueError 里的支持列表;把同一份 kg 用 export_json(kg, "a.jsonld", format="json-ld") 跑两遍,diff 两份文件的 @id,确认一致;给 LPG 导出的实体 type 字段填 "Person Name",观察生成的 Cypher 里标签带空格后语法是否还能过。
逻辑轨迹面板伪代码(右侧标真实行号):
if backend not in SUPPORTED_BACKENDS: raise # vector_store.py:116-121
if self._backend_store: # vector_store.py:509
if hasattr(store, 'add'): return add(...) # vector_store.py:511-512
elif hasattr(store, 'add_vectors'): ... # vector_store.py:513-528
similarity = 1.0 / (1.0 + distance) # pgvector_store.py:440
score = 1.0 / (k + rank) # hybrid_search.py:166
digest = hash_data(canonical)[:16] # json_exporter.py:56
format = format_map.get(ext, "json") # methods.py:943可迁移结论
值得抄的: 门面收口时先定返回结构,再谈后端。SearchResult 那五个字段是这个设计的灵魂:score 统一、distance 保留原值、vector 可空,等于给所有后端一个共同的目标形态,新后端照着填就行。这套不依赖 Python,任何语言的存储门面都能抄。
最小成本形态: 如果只有两个后端,不需要 SUPPORTED_BACKENDS 集合加 hasattr 探测链的完整配置,一个 if backend == "a": return AStore(...) 的 if/elif 就够。但「统一结果结构」和「能力缺失抛异常」这两条要留,否则换后端时调用方会拿到不同键名的结果。出口层若只做幂等,_content_iri 一个函数就是全部:规范序列化加哈希,二十行内能落地。
哪些是过度设计: VectorStore 里 store_decision、search_decisions、build_decision_context、explain_decision 这整块决策追踪能力,对纯存储场景是负担,它把向量库和一个「决策记忆」概念绑在一起,而门面的核心职责只是读写向量。GraphStore 的 build_from_conversations、add_nodes、add_edges 这类「兼容方法」里塞满了 # Let's check...、# Assuming... 的注释,是典型的未完成代码,抄之前要确认自己真需要从对话直接建图。RDF 侧的 export_shacl、RDFValidator 的 namespace 校验,属于规范完整性,单团队内部分享图谱时用不上。
不依赖 Python 的一条结论: 内容寻址标识符可以用任何有哈希函数的语言实现,核心是「先规范化再哈希」。规范化要满足三件事:键排序、类型稳定、可序列化值兜底。这三件事分别对应 json.dumps(sort_keys=True)、default=str、以及只哈希内容不哈希时间。这条从 Semantica 的两个 IRI 铸造函数里都能独立抽出来。
思考题
-
VectorStore.search_vectors委托给有search_vectors方法的后端时,数量参数按位置传而不按关键字传。如果某个后端把这个位置参数实现成了「返回向量本身的最大维度」,会出什么问题?结合第 704 到 708 行的注释,说说位置绑定牺牲了什么。 -
_content_iri用json.dumps(payload, sort_keys=True)铸 IRI,其中payload是整个 kg 字典。如果 kg 里有一个键的值是 dict,且该 dict 的两个键顺序在一次代码升级后互换,IRI 会不会变?为什么sort_keys=True能兜住这一层、却兜不住「键本身被重命名」? -
动手验证:在 semantica 仓库根目录起 Python 交互环境,运行:
from semantica.export.methods import export_json
kg = {"entities": [{"id": "e1", "text": "Alice", "type": "Person"}],
"relationships": [{"source_id": "e1", "target_id": "e2", "type": "knows"}]}
export_json(kg, "/tmp/a.jsonld", format="json-ld")
export_json(kg, "/tmp/b.jsonld", format="json-ld")
print(open("/tmp/a.jsonld").read() == open("/tmp/b.jsonld").read())观察输出是 True 还是 False。把 kg 里 "text": "Alice" 改成 "text": "Alic" 再导出第三份,对比三份文件的 @id 字段各是什么,说明内容寻址的粒度。再跑 VectorStore(backend="qdrnt"),读它抛出的 ValueError 消息里出现了哪八个后端名。
LPGExporter._create_nodes_batch对node_type不做转义,对node_id和label做了转义。如果实体type是"Person Name",生成的 Cypher 里标签Person Name含空格,导入 Neo4j 时会发生什么?为什么作者只转义了值、没转义标签和关系类型,这是遗漏还是「标签由受控本体生成」的隐含假设?