杂乱数据如何变成一张可解释、可审计的知识图谱 10 / 10

第 10 章

存储与出口

源码核对基于 semantica-agi/semantica commit `1ee2ae88`,tag `course-anchor-20260824`

场景还原

数据平台组刚把知识图谱跑通,老板要求在三个环境里各接一个向量库:开发用 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 就是那份契约:

semantica/vector_store/vector_store.py81:97
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」的承诺,这意味着每个后端的适配器都要负责把自家度量转成这个尺度。

门面本身靠一张白名单挡住拼写错误:

semantica/vector_store/vector_store.py100:112
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__ 第一件事就是把传入的 backendlower() 再查这个集合,不在就直接 ValueError。集合用 sorted() 展开成提示语,报错时用户能看到全部合法选项。注意 inmemory 也在集合里,它是「不接任何外部服务」的本地兜底,后面所有后端归一化的逻辑都以它为参照。

六种后端的差异被压缩在 _init_backend_store 的一个 if/elif 链里,每个分支做两件事:从配置文件读参数、实例化对应适配器。pgvector 分支额外检查 connection_string 是否缺失,缺失就抛带示例的 ValueError;sqlite 分支检查 db_path。这套「参数校验前置」让错误发生在构造时,把失败提前到任何查询之前。

真正考验设计的是写和查两个方法。store_vectors 的委托逻辑:

semantica/vector_store/vector_store.py507:529
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,让「这个后端不支持写入」以异常形式暴露,绝不静默吞掉。

搜索侧同款逻辑,但多了一个针对性注释:

semantica/vector_store/vector_store.py696:709
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_kklimit 三种命名在不同后端里打架。注释给了解法:对有 search_vectors 方法的后端,数量参数按位置传,不猜关键字。这就是场景还原里 top_k 冲突的反面教材的来源,门面这一层用「位置绑定」绕开了命名差异。

count() 则是「能力缺失要显式」的样本:

semantica/vector_store/vector_store.py827:854
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 里:

semantica/vector_store/pgvector_store.py374:380
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 的三种度量对应三个运算符,且 <=><-> 返回的是距离,越小越好;<#> 返回的是负内积,排序方向相反。这套差异在下游被折成统一分数:

semantica/vector_store/pgvector_store.py437:451
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;内积因为运算符自带负号,取反就是真内积。结果字典里 vectordistance 都填 None,因为 SQL 查询默认不把原始向量捞回来,原生距离值在换算后也不再单独暴露。这就是 SearchResult 契约里那两个 Optional 字段的实际落点。

元数据过滤里的注入防线:

semantica/vector_store/pgvector_store.py459:472
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 片段,任何没被正则放行的键都会在 searchfilter_by_metadata 里提前抛 ValidationError。这条正则允许字母数字下划线连字符,且首字符只能是字母或下划线,堵住了「键名里塞一段 SQL」的口子。表名和维度则走 psycopg_sql.SQL 的占位组合,双保险。

混合检索:RRF 落在哪一层

RRF 不在 VectorStore 门面里,而在 hybrid_search.pySearchRanker。这是「融合」和「检索」分层的证据:门面只管单一向量源的读写,多路结果的合并是独立组件。

semantica/vector_store/hybrid_search.py140:247
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 显式分派,因为图数据库的操作面更宽,靠反射猜方法名会失控。

semantica/graph_store/graph_store.py525:562
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 透传,因为图库的连接参数更多。

semantica/graph_store/graph_store.py564:597
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 覆盖全局配置,意味着调用方在构造时传的参数优先级最高。neptuneamazon_neptuneageapache_age 的别名合并在一个分支,用 or 连接,说明后端命名允许两种叫法。未知后端走 ValidationError,报错时没有「支持列表」的提示,这一点比向量门面弱。

QueryEngine 是图门面里最有价值的一个独立类,完整定义:

semantica/graph_store/graph_store.py240:310
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 侧先看置信度数据类型这一处注释,它解释了为什么一个数值要费这么大劲:

semantica/export/rdf_exporter.py53:64
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 铸造:

semantica/export/rdf_exporter.py123:133
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 序列化把实体转成三元组子句:

semantica/export/rdf_exporter.py833:860
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 也走了内容寻址:

semantica/export/json_exporter.py37:57
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 一致,合并时自然去重。

图节点用这个函数命名:

semantica/export/json_exporter.py609:614
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 导出最容易翻车的地方:

semantica/export/lpg_exporter.py229:267
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)

idname 两个固定属性先行转义,其余标量属性过滤后逐个 _format_cypher_value,非标量(dict、list)直接丢弃,因为 Cypher 属性值只收标量。类型直接当节点标签用,不做映射。这里有个隐患会在边界条件里展开:node_type 没转义,如果实体类型名里带空格或特殊字符,生成的 Cypher 就是非法语句。

关系侧:

semantica/export/lpg_exporter.py281:327
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 是静默跳过,坏边不报错,只在结果里少一条。

转义函数本身:

semantica/export/lpg_exporter.py329:341
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()。这套转义覆盖了「值里带引号」的常见翻车,但覆盖不了标签和关系类型的注入。

最后看出口层的总调度:

semantica/export/methods.py924:953
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 链。

sequenceDiagram participant C as 调用方 participant V as VectorStore participant B as 后端适配器 C->>V: store_vectors V->>V: SUPPORTED_BACKENDS 白名单校验 V->>B: hasattr 探测 add 或 add_vectors B-->>V: 返回向量 ID V-->>C: 统一 SearchResult 列表
flowchart LR A[多路结果列表] --> B[逐条算 1 除以 k 加 rank] B --> C[按 id 累加分数] C --> D[按分数降序] D --> E[回填原始结果并写回 score]
flowchart TD A[kg 字典] --> B{扩展名} B -->|.ttl| C[RDFExporter\nTurtle] B -->|.jsonld| D[JSONExporter\nJSON-LD] B -->|.cypher| E[LPGExporter\nCypher] C --> F[写文件] D --> F E --> F

设计决策分析

存储层最核心的决策是「鸭子类型收口,不强制统一接口」。Semantica 没有给后端定义一个 AbstractVectorStore 抽象基类,而是让 VectorStorehasattr 逐个探测方法名。这样做的代价是每次委托都要写一条 if/elif 链,好处是新后端只要实现了「长得像」的方法就能挂进来,不需要继承任何东西。docs/architecture.md 第 152 行把这条哲学写成了原则:

docs/architecture.md152: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」是这句话的落点。换后端时调用方一行代码不改,这句也写进了性能特性表:

docs/architecture.md184:184
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_extensionpg_extension 表发现没有 vector 行,抛 ProcessingError 并附安装链接,这是第 210 到 227 行。扩展在但表不存在时,_ensure_table_existsCREATE 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 是一个抽象基类,把接口钉死在抽象方法上:

packages/graphrag-vectors/graphrag_vectors/vector_store.py56:66
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,实现 connectcreate_indexload_documentssimilarity_search_by_vector 等抽象方法。工厂用 match 语句按枚举类型注册:

packages/graphrag-vectors/graphrag_vectors/vector_store_factory.py73:92
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),注释没有,语义直接丢给调用方:

packages/graphrag-vectors/graphrag_vectors/lancedb.py216:225
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 抽象处理的是「键值文件存储」:

packages/graphrag-storage/graphrag_storage/storage.py13:25
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/ 下用 rdfjsonldjson-ldturtleowl 作关键词搜 Python 文件,除 knowledgegraphrag 这类误命中外无结果;find packages -iname '*export*' 无文件。GraphRAG 没有出口层,它用 graphrag-storageParquetTableProviderCsvTableProvidergraphrag-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_irijson.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 里标签带空格后语法是否还能过。

逻辑轨迹面板伪代码(右侧标真实行号):

text
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 一个函数就是全部:规范序列化加哈希,二十行内能落地。

哪些是过度设计: VectorStorestore_decisionsearch_decisionsbuild_decision_contextexplain_decision 这整块决策追踪能力,对纯存储场景是负担,它把向量库和一个「决策记忆」概念绑在一起,而门面的核心职责只是读写向量。GraphStorebuild_from_conversationsadd_nodesadd_edges 这类「兼容方法」里塞满了 # Let's check...# Assuming... 的注释,是典型的未完成代码,抄之前要确认自己真需要从对话直接建图。RDF 侧的 export_shaclRDFValidator 的 namespace 校验,属于规范完整性,单团队内部分享图谱时用不上。

不依赖 Python 的一条结论: 内容寻址标识符可以用任何有哈希函数的语言实现,核心是「先规范化再哈希」。规范化要满足三件事:键排序、类型稳定、可序列化值兜底。这三件事分别对应 json.dumps(sort_keys=True)default=str、以及只哈希内容不哈希时间。这条从 Semantica 的两个 IRI 铸造函数里都能独立抽出来。

思考题

  1. VectorStore.search_vectors 委托给有 search_vectors 方法的后端时,数量参数按位置传而不按关键字传。如果某个后端把这个位置参数实现成了「返回向量本身的最大维度」,会出什么问题?结合第 704 到 708 行的注释,说说位置绑定牺牲了什么。

  2. _content_irijson.dumps(payload, sort_keys=True) 铸 IRI,其中 payload 是整个 kg 字典。如果 kg 里有一个键的值是 dict,且该 dict 的两个键顺序在一次代码升级后互换,IRI 会不会变?为什么 sort_keys=True 能兜住这一层、却兜不住「键本身被重命名」?

  3. 动手验证:在 semantica 仓库根目录起 Python 交互环境,运行:

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 消息里出现了哪八个后端名。

  1. LPGExporter._create_nodes_batchnode_type 不做转义,对 node_idlabel 做了转义。如果实体 type"Person Name",生成的 Cypher 里标签 Person Name 含空格,导入 Neo4j 时会发生什么?为什么作者只转义了值、没转义标签和关系类型,这是遗漏还是「标签由受控本体生成」的隐含假设?