Haystack 与 DataStax Astra DB 集成实战:AstraDocumentStore 与 AstraEmbeddingRetriever 构建向量检索流水线
【免费下载链接】haystackOpen-source AI orchestration framework for building context-engineered, production-ready LLM applications. Design modular pipelines and agent workflows with explicit control over retrieval, routing, memory, and generation. Built for scalable agents, RAG, multimodal applications, semantic search, and conversational systems.项目地址: https://gitcode.com/GitHub_Trending/ha/haystack
本文以 Haystack 的 Astra 集成(astra-haystack)为核心,系统讲解如何用AstraDocumentStore将海量文档写入基于 Apache Cassandra 的 serverless 向量数据库 DataStax Astra DB,并通过AstraEmbeddingRetriever在 RAG、语义搜索与抽取式问答流水线中完成向量检索。读完本文,你将掌握 Astra 文档库的初始化参数、去重策略、元数据过滤与文档管理 API 的完整用法,并能在 HaystackPipeline中端到端搭建一套"文本嵌入 → 向量检索 → 生成"的生产级检索链路。
一、为什么选择 Astra DB 作为 Haystack 的向量存储
DataStax Astra DB 是一款构建于 Apache Cassandra 之上的 serverless 向量数据库,原生支持向量搜索与自动扩缩容,可部署在 AWS、GCP 或 Azure,并能轻松扩展至多个云区域,以获得多区域可用性、低延迟数据访问和数据主权,同时避免云厂商锁定。
在 Haystack 中,Astra 集成以haystack_integrations包的形式提供,核心包含两个组件:
AstraDocumentStore(文档存储):负责与 Astra DB 建立连接、写入/查询/删除文档,并基于向量相似度执行检索;AstraEmbeddingRetriever(向量检索器):接收查询向量,从AstraDocumentStore中召回与查询最相关的文档。
两者对应的 API 参考文档见 版本化 API 参考,完整使用指南见 AstraDocumentStore 使用文档 与 AstraEmbeddingRetriever 使用文档。
二、安装集成与获取连接凭证
2.1 安装
在已创建 Astra DB 账号与数据库的前提下,安装astra-haystack集成:
pip install astra-haystack如果希望直接运行文中的嵌入示例(不依赖云端嵌入 API),可以一并安装 sentence-transformers:
pip install sentence-transformers2.2 获取凭证
在 AstraDB 的 Web UI 中,你需要准备两类信息:
- 数据库 ID / API Endpoint:在 Astra 控制台的 Connect 标签页中,选择JSON API并点击Generate Configuration,即可生成 API Endpoint;
- Application Token:同样在 Connect 页面生成,作为访问数据库的认证令牌。
此外,你还需要一个collection 名称与一个namespace。创建 collection 时,必须同步指定嵌入向量维度(embedding dimensions)和相似度度量(similarity metric)。其中 namespace 用于在数据库中组织数据,在 Apache Cassandra 术语中称为keyspace。
2.3 通过环境变量管理密钥
Haystack 强烈建议通过环境变量传递认证数据,而不是把密钥硬编码进代码。运行示例前,请先填充:
export ASTRA_DB_API_ENDPOINT="https://<database-id>-<region>.apps.astra.datastax.com" export ASTRA_DB_APPLICATION_TOKEN="AstraCS:..."这两个环境变量名与AstraDocumentStore构造函数中Secret.from_env_var(...)的默认读取来源完全对应,详见下文。
三、AstraDocumentStore:初始化与核心参数
AstraDocumentStore是 Astra 集成在 Haystack 侧的"数据面",连接通过Astra DB JSON API建立与管理。最简初始化方式是利用环境变量:
from haystack import Document from haystack_integrations.document_stores.astra import AstraDocumentStore document_store = AstraDocumentStore() document_store.write_documents( [Document(content="This is first"), Document(content="This is second")], ) print(document_store.count_documents())3.1 完整构造函数签名
__init__( api_endpoint: Secret = Secret.from_env_var("ASTRA_DB_API_ENDPOINT"), token: Secret = Secret.from_env_var("ASTRA_DB_APPLICATION_TOKEN"), collection_name: str = "documents", embedding_dimension: int = 768, duplicates_policy: DuplicatePolicy = DuplicatePolicy.NONE, similarity: str = "cosine", namespace: str | None = None, ) -> None3.2 参数说明
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
api_endpoint | Secret | 环境变量ASTRA_DB_API_ENDPOINT | Astra DB 的 JSON API 端点,可在控制台 Connect 页面生成 |
token | Secret | 环境变量ASTRA_DB_APPLICATION_TOKEN | Astra DB 应用令牌(Application Token) |
collection_name | str | "documents" | 当前 Astra DB 中 keyspace 内使用的 collection 名称 |
embedding_dimension | int | 768 | 嵌入向量的维度,必须与写入向量的维度一致 |
duplicates_policy | DuplicatePolicy | DuplicatePolicy.NONE | 处理重复文档的策略,取值见下节 |
similarity | str | "cosine" | 用于比较文档向量的相似度函数 |
namespace | str \| None | None | 数据所在命名空间(Cassandra keyspace),不传时使用默认 keyspace |
注意:如果 API endpoint 或 token 未设置,构造函数会抛出
ValueError。
从源码结构看,Secret机制是 Haystack 统一的敏感信息管理方式:Secret.from_env_var(...)让密钥只在真正发起请求时才从环境变量读取,从而避免密钥出现在序列化后的 YAML/JSON 配置中。DuplicatePolicy与FilterPolicy等类型定义在核心库的 文档存储类型目录 下,其中 policy.py 定义了去重枚举,filter_policy.py 定义了过滤合并策略。
四、写入文档与重复文档处理策略
write_documents用于将文档索引入库,供后续查询使用:
write_documents( documents: list[Document], policy: DuplicatePolicy = DuplicatePolicy.NONE ) -> int- 入参:
documents为 HaystackDocument对象列表;policy指定重复文档处理策略。 - 返回:实际写入的文档数量(
int)。 - 异常:
ValueError—— 传入的文档既不是Document也不是dict;DuplicateDocumentError—— 已存在相同 ID 的文档且策略为FAIL;Exception—— 文档 ID 不是字符串,或文档中同时包含id与_id字段。
4.1 DuplicatePolicy 四种策略
DuplicatePolicy枚举定义在 haystack/document_stores/types/policy.py(核心库),Astra 文档库直接复用该枚举:
| 策略 | 行为 |
|---|---|
DuplicatePolicy.NONE | 默认策略。如果相同 ID 的文档已存在,则跳过、不写入 |
DuplicatePolicy.SKIP | 如果相同 ID 的文档已存在,则跳过、不写入 |
DuplicatePolicy.OVERWRITE | 如果相同 ID 的文档已存在,则覆盖写入 |
DuplicatePolicy.FAIL | 如果相同 ID 的文档已存在,则抛出错误 |
注意NONE与SKIP的文档行为一致,二者区别在于语义定位:NONE是构造函数与write_documents的默认值,SKIP更适合在明确"允许跳过重复"的业务场景中显式声明。该枚举同样被 Haystack 核心文档存储(如 InMemoryDocumentStore)与文档写入组件 DocumentWriter 复用,是整个框架统一的去重语义。
五、文档管理 API:查询、过滤、更新与统计
AstraDocumentStore实现了 Haystack 文档存储协议(协议定义见 haystack/document_stores/types/protocol.py),并提供如下完整的数据管理方法:
5.1 查询类
filter_documents(filters: dict[str, Any] | None = None) -> list[Document]:返回最多 1000 条匹配过滤条件的文档;过滤条件非法或不受支持时抛出AstraDocumentStoreFilterError。get_documents_by_id(ids: list[str]) -> list[Document]:按 ID 列表批量获取文档。get_document_by_id(document_id: str) -> Document:按单个 ID 获取文档;未找到时抛出MissingDocumentError。search(query_embedding: list[float], top_k: int, filters: dict[str, Any] | None = None) -> list[Document]:基于查询向量执行相似度检索,返回top_k条匹配文档。这是向量检索的底层实现,AstraEmbeddingRetriever在run时最终会调用到它。
5.2 计数类
count_documents() -> int:统计文档库中的文档总数。count_documents_by_filter(filters: dict[str, Any]) -> int:应用过滤条件后统计匹配文档数。count_unique_metadata_by_filter(filters, metadata_fields: list[str]) -> dict[str, int]:对匹配文档的每个元数据字段统计唯一值个数,返回{字段名: 唯一值个数}。
5.3 删除类
delete_documents(document_ids: list[str]) -> None:按 ID 删除文档;若提供了 ID 但没有任何文档被删除,抛出MissingDocumentError。delete_all_documents() -> None:清空文档库。delete_by_filter(filters: dict[str, Any]) -> int:删除匹配过滤条件的文档,返回删除数量;过滤非法时抛出AstraDocumentStoreFilterError。
5.4 更新类
update_by_filter(filters: dict[str, Any], meta: dict[str, Any]) -> int:将匹配过滤条件文档的元数据更新为meta指定字段(与已有元数据合并),返回更新数量。
5.5 元数据探查类
get_metadata_fields_info() -> dict[str, dict[str, str]]:返回元数据字段及其类型映射,形如{"field_name": {"type": "..."}}。get_metadata_field_min_max(metadata_field: str) -> dict[str, Any]:返回指定字段的min与max值。get_metadata_field_unique_values(metadata_field, search_term=None, from_=0, size=10, filters=None) -> tuple[list[Any], int]:检索某字段的唯一值,支持大小写不敏感的模糊搜索(search_term)、分页(from_/size)与过滤(filters),返回(分页后的值列表, 总数)。
关于get_metadata_field_unique_values的类型注意点:不同类型但值相等的元数据会被视为不同的唯一值(例如整数1、布尔True与字符串"1"会被分别返回),但有一个例外——Astra DB 的 Data API 在存储时会无条件将整数值的浮点数(如1.0)规范化(canonicalize)为整数,因此1.0(float)写入后读回一定是整数1(int),而带小数部分的浮点数(如1.5)不受影响,可以正常往返。
六、AstraEmbeddingRetriever:向量检索组件
AstraEmbeddingRetriever是嵌入检索组件,它比较查询向量与文档向量的相似度,并根据结果从AstraDocumentStore中召回最相关的文档。其 API 参考见 版本化 API 参考 的haystack_integrations.components.retrievers.astra.retriever部分。
6.1 初始化
__init__( document_store: AstraDocumentStore, filters: dict[str, Any] | None = None, top_k: int = 10, filter_policy: str | FilterPolicy = FilterPolicy.REPLACE, ) -> None| 参数 | 说明 |
|---|---|
document_store | 一个AstraDocumentStore实例(必填) |
filters | 用于缩小搜索空间的过滤条件字典 |
top_k | 最多检索的文档数量,默认 10 |
filter_policy | 过滤条件应用策略(REPLACE或MERGE),默认REPLACE |
6.2 run 与 run_async
run( query_embedding: list[float], filters: dict[str, Any] | None = None, top_k: int | None = None, ) -> dict[str, list[Document]]query_embedding:查询文本的向量表示(float 列表)。filters:运行时应用的过滤条件。运行时过滤条件的具体应用方式取决于初始化时选择的filter_policy。top_k:运行时覆盖的最大检索数量(可选)。- 返回:
{"documents": [...]},即从AstraDocumentStore检索到的文档列表。
run_async的签名与返回结构与run完全一致,用于异步场景;从文档说明看,它是将同步搜索放到线程池中执行,从而避免阻塞事件循环,适合在异步 Pipeline 中与run_async链路配合使用。
6.3 序列化
to_dict() -> dict[str, Any]:将组件序列化为字典,便于通过 YAML/JSON 配置持久化。from_dict(data: dict[str, Any]) -> AstraEmbeddingRetriever:从字典反序列化重建组件。
这两个方法让AstraEmbeddingRetriever可以无缝嵌入 Haystack 的声明式 Pipeline 定义(序列化为 YAML/JSON)中。
七、filter_policy:初始化过滤与运行时过滤如何合并
filter_policy控制检索器初始化时设置的过滤条件(init filters)与run调用时传入的过滤条件(runtime filters)之间的关系。核心枚举与合并逻辑定义在 haystack/document_stores/types/filter_policy.py:
FilterPolicy.REPLACE(默认):运行时传入的过滤条件替换初始化时设置的过滤条件。FilterPolicy.MERGE:运行时过滤条件与初始化过滤条件合并,若字段重叠,运行时值覆盖初始化值。
从 filter_policy.py 中apply_filter_policy的实现可以看出,MERGE策略会依据初始化/运行时过滤条件各自是"比较型过滤"(包含field、operator、value键)还是"逻辑型过滤"(包含operator与conditions键)进行四种组合:
| 初始化过滤 | 运行时过滤 | 合并结果 |
|---|---|---|
| 比较型 | 比较型 | 以AND合并为逻辑过滤(同字段时运行时覆盖初始化) |
| 比较型 | 逻辑型 | 当逻辑运算符一致时,把比较条件并入conditions |
| 逻辑型 | 比较型 | 当逻辑运算符一致时,把比较条件并入conditions(同字段时运行时覆盖) |
| 逻辑型 | 逻辑型 | 运算符相同则合并conditions,否则忽略初始化过滤并告警 |
例如,初始化时设置{"field": "meta.type", "operator": "==", "value": "article"},运行时传入{"operator": "AND", "conditions": [{"field": "meta.rating", "operator": ">=", "value": 3}]},在MERGE+AND策略下会合并为一个同时包含两个条件的逻辑过滤。该机制让"公共过滤条件固化在组件里、个性化过滤条件每次查询动态传入"成为可能,是构建多租户或分类检索场景的关键开关。
八、端到端示例:在 Pipeline 中使用 AstraEmbeddingRetriever
下面是一个完整的 RAG 语义检索流水线示例(节选自 astraretriever.mdx 使用文档),演示了"文档嵌入写入 + 查询嵌入检索"的完整闭环。
from haystack import Document, Pipeline from haystack.components.embedders import ( SentenceTransformersTextEmbedder, SentenceTransformersDocumentEmbedder, ) from haystack_integrations.components.retrievers.astra import AstraEmbeddingRetriever from haystack_integrations.document_stores.astra import AstraDocumentStore document_store = AstraDocumentStore() model = "sentence-transformers/all-mpnet-base-v2" documents = [ Document(content="There are over 7,000 languages spoken around the world today."), Document( content="Elephants have been observed to behave in a way that indicates a high level of self-awareness, such as recognizing themselves in mirrors.", ), Document( content="In certain parts of the world, like the Maldives, Puerto Rico, and San Diego, you can witness the phenomenon of bioluminescent waves.", ), ] document_embedder = SentenceTransformersDocumentEmbedder(model=model) document_embedder.warm_up() documents_with_embeddings = document_embedder.run(documents) document_store.write_documents( documents_with_embeddings.get("documents"), policy=DuplicatePolicy.SKIP, ) query_pipeline = Pipeline() query_pipeline.add_component( "text_embedder", SentenceTransformersTextEmbedder(model=model), ) query_pipeline.add_component( "retriever", AstraEmbeddingRetriever(document_store=document_store), ) query_pipeline.connect("text_embedder.embedding", "retriever.query_embedding") query = "How many languages are there?" result = query_pipeline.run({"text_embedder": {"text": query}}) print(result["retriever"]["documents"][0])示例输出(分数与向量维度取决于所用嵌入模型):
Document(id=cfe93bc1c274908801e6670440bf2bbba54fad792770d57421f85ffa2a4fcc94, content: 'There are over 7,000 languages spoken around the world today.', score: 0.8929937, embedding: vector of size 768)8.1 典型流水线位置
根据 astraretriever.mdx 的说明,AstraEmbeddingRetriever最常见的三种流水线位置是:
- RAG 流水线:位于 Text Embedder 之后、
PromptBuilder之前; - 语义搜索流水线:作为查询流水线的最后一个组件直接输出结果;
- 抽取式问答流水线:位于 Text Embedder 之后、
ExtractiveReader之前。
使用时请确保索引流水线中已有Document Embedder、查询流水线中已有Text Embedder,为检索器提供查询向量与文档向量。
九、索引警告(Indexing Warnings)的原因与处理
创建AstraDocumentStore时,你可能会看到如下两类警告之一:
Astra DB collection
...is detected as having indexing turned on for all fields (either created manually or by older versions of this plugin). This implies stricter limitations on the amount of text each string in a document can store. Consider indexing anew on a fresh collection to be able to store longer texts.
或者:
Astra DB collection
...is detected as having the following indexing policy:{...}. This does not match the requested indexing policy for this object:{...}. In particular, there may be stricter limitations on the amount of text each string in a document can store. Consider indexing anew on a fresh collection to be able to store longer texts.
9.1 为什么会出现该警告
collection 已经存在,且被配置为对所有字段开启索引(可能由你之前手动创建,或由旧版本插件创建)。而 Haystack 创建 collection 时会应用一套针对其用途优化的索引策略:该策略允许存储更长的文本,并避免索引那些你不需要过滤的字段,从而降低写入开销。
9.2 常见原因
- 你在 Haystack 之外创建了 collection(例如在 Astra UI 中手动创建,或通过 AstraPy 的
Database.create_collection()创建); - 你使用旧版本的插件创建了 collection。
9.3 影响与解决方案
影响:这只是一个警告。除非你尝试存储非常长的文本字段(此时 Astra DB 会返回索引错误),否则应用可以正常运行。
解决方案:
- 推荐做法:如果能够重新填充数据,删除并重建 collection,然后重新运行你的 Haystack 应用,让插件以优化后的索引策略重新创建 collection;
- 忽略警告:如果你确定不会存储很长的文本字段,可以直接忽略该警告继续使用。
十、异常体系:Astra 文档存储的错误类型
Astra 集成定义了三级错误体系(见 版本化 API 参考 的haystack_integrations.document_stores.astra.errors部分),与 Haystack 核心错误类型衔接:
| 异常类 | 基类 | 触发场景 |
|---|---|---|
AstraDocumentStoreError | DocumentStoreError | 所有 AstraDocumentStore 错误的父类 |
AstraDocumentStoreFilterError | FilterError | 向 AstraDocumentStore 传入了非法过滤条件 |
AstraDocumentStoreConfigError | AstraDocumentStoreError | 向 AstraDocumentStore 传入了非法配置 |
在编写健壮的检索代码时,建议对AstraDocumentStoreFilterError(过滤条件书写错误)与AstraDocumentStoreConfigError(初始化配置错误)分别捕获,以便快速定位是查询侧问题还是存储侧问题。
十一、结语
Astra 集成让 Haystack 可以直接利用 DataStax Astra DB 的 serverless 向量数据库能力:AstraDocumentStore承担文档写入、过滤、统计、更新与向量检索等全部数据面操作,AstraEmbeddingRetriever则作为标准 Haystack 组件无缝嵌入 RAG、语义搜索与抽取式问答流水线。结合本文介绍的DuplicatePolicy去重策略、filter_policy过滤合并机制与索引警告处理方案,你可以在生产环境中稳定运行基于 Astra DB 的向量检索应用。更多配套资料可参阅 AstraDocumentStore 使用文档、AstraEmbeddingRetriever 使用文档 以及 Haystack 文档存储类型定义。
【免费下载链接】haystackOpen-source AI orchestration framework for building context-engineered, production-ready LLM applications. Design modular pipelines and agent workflows with explicit control over retrieval, routing, memory, and generation. Built for scalable agents, RAG, multimodal applications, semantic search, and conversational systems.项目地址: https://gitcode.com/GitHub_Trending/ha/haystack
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考