Indexify:开源的数据框架,可以轻松抽取任意规模的非结构化数据

核心:分布式的数据平台,扩展性强,支持对非结构化数据(多模态数据)进行统一处理

创始人对核心架构设计的解释

HackerNews: https://news.ycombinator.com/item?id=40458923 (opens in a new tab)

嗨,HN,我是 Tensorlake 的创始人。制作 LLM 应用程序原型已经变得容易多了,但在生产环境中构建能够处理不断更新数据的决策型 LLM 应用程序仍然非常具有挑战性。我们看到人们面临的系统工程问题是 -

  1. 如果应用程序对信息的新鲜度很敏感,则可靠地实时处理摄取的内容。
  2. 能够引入任何类型的模型,并在 GPU 和 CPU 上运行管道的不同部分。
  3. 对摄取高峰、计算基础设施故障的容错能力。
  4. 随着数据量的增长,扩展计算、读取和写入。

我们建立并开源了 Indexify(https://github.com/tensorlakeai/indexify),为在数据频繁更新或新数据不断创建的动态环境中工作的 (opens in a new tab) LLM 应用程序提供计算引擎和数据框架。

开发人员描述了一个声明式提取图,其中的各个阶段用于提取或转换非结构化数据。数据从一个阶段传递到另一个阶段,最终到达接收器,如矢量数据库、Blob 存储或 Postgres 等结构化数据存储。

示例:

  1. 进行视频理解的图表可以是:摄取 -> 音频提取 -> 转录 -> NER 和嵌入。另一条路径,摄取 -> 关键帧提取 -> 对象和场景描述 ( https://github.com/tensorlakeai/indexify/blob/main/docs/docs (opens in a new tab)... )
  2. PDF 上的结构化提取和搜索:PDF -> Markdown -> 分块 -> 嵌入,NER ( https://github.com/tensorlakeai/indexify/blob/main/docs/docs (opens in a new tab)... )

应用层: Indexify 可作为 LLM 应用程序堆栈中的检索器,因此您可以轻松地将其与现有应用程序配合使用。通过 HTTP 调用检索器 API 以获取从 Indexify 提取的数据,这几乎就是您搜索或检索数据所需的所有集成。

您可以使用可组合的提取器并将它们链接在一起以构建可处理任何非结构化数据的复杂实时数据管道。

因为这是 HN,所以我可以自由地谈论一些技术细节 :)

它如何实现实时性?我们使用 Raft 构建了一个复制状态机,每秒处理数十或数千个摄取事件。存储和网络层经过优化,可使调度程序在 2 毫秒内创建任务。 调度程序的架构与 Google 的 Borg 和 Hashicorp 的 Nomad 非常相似。我们的架构可以扩展到多台机器上的并行调度,并拥有像 Nomad 这样的集中式排序器。

存储系统:由于重点是非结构化数据,我们希望能够支持存储和提取大型文件,并能够随着数据量的增长而水平扩展。 Indexify 在后台使用 blob 存储来存储非结构化数据。如果图形创建了嵌入,它们会自动存储在向量存储中,而结构化数据则会存储在 Postgres 等结构化存储中。 在后台,我们在摄取服务器和数据存储之间拥有 Rust 特性,因此我们可以轻松实现对其他向量存储的支持。

同步向量和结构化存储 : 如果 Indexify 检测到图中同时存在结构化数据和向量存储,它还会将结构化数据与向量存储同步。这样就可以使用预过滤功能来缩小搜索范围,以获得更好的结果。

API: Indexify 在向量存储上公开语义搜索 API,并在半结构化数据上公开读取 SQL 查询。我们可以自动找出结构化数据的架构并在其上公开 SQL 接口。在后台,我们解析 SQL,并有一个层来扫描和读取数据库以对行进行切片和切块。因此,BI 工具应该可以立即在提取的数据上工作。我们有 Python 和 Typescript 库,使人们可以轻松地构建新应用程序或将其集成到现有应用程序中。

I'm curious about how Indexify handles fault tolerance during ingestion spikes and compute infrastructure failures. Could you please provide some examples?

Sorry just seeing this! There are various aspects of how it handles ingestion spikes -

  1. The ingestion api writes to blob stores, which are horizontally scalable. Only when an ingestion finishes, we write the metadata to the replicated state machine.
  2. The replicated state machine takes 100k IOPs on most commodity machines, and they can vertically scale.
  3. The extractors can autoscale based on the number of tasks in the system.
  4. The ingestion server cluster can autoscale also based on the amount of IOPs they are doing.

它处理摄取峰值的方式有很多种 -

  1. 摄取 API 写入 Blob 存储,这些存储是水平可扩展的。只有当摄取完成时,我们才会将元数据写入复制状态机。
  2. 复制状态机在大多数商用机器上需要 100k IOP,并且可以垂直扩展。
  3. 提取器可以根据系统中的任务数量自动扩展。
  4. 摄取服务器集群也可以根据它们正在执行的 IOP 数量自动扩展。

索引文档与由于系统出现暂时性错误而将其丢弃之间的区别在于,用户满意与失去用户信任。 Indexify 的提取状态机允许将计算表示为管道,其中每个步骤都是持久的,并且会自动重试失败。 控制平面本身可以在地理区域内复制,从而使其能够抵御计算节点故障或整个数据中心故障。

整体架构

  1. 存储的内容格式
  2. 数据存储
  3. 数据检索
  4. 任务拆分是如何进行的

梳理核心流程以及模块之间的交互形式;同时,确定核心模块的整体设计

Indexify Server

Indexify Server 包括 Ingestion Server 和 Coordinator。在大规模生成环境中,Indexify可以进行水平扩展以实现高可用。

Coordinator的设计

Coordinator 是一个速度超快的任务调度程序,每当数据被提取时,它每秒都会创建数千个任务。它评估提取策略并将任务分配给提取器。 Coordinator 没有任何外部依赖项,并使用后台的复制状态机在多台机器上复制自身。 这种设计使我们能够构建一个完全反应式的调度程序来评估数千个提取策略和内容并安排任务。

底层使用了 OpenRaft

整体架构的核心处理流程

  1. Document上传,转换为ContentMetadata,由coordinatorClient调用create_content(content_metadata),存储coordinator的share_state中。在state_change中增加internal_api::ChangeType::NewContent 对应的update_entry,同时写入raft中
  2. coordinator如何分配share_state中的任务?coordinator scheduler监听state_changes,对于不同的ChangeType并执行相应的action(函数).
  3. scheduler通过content metadata来创建任务;由 AllocationPlanner 来规划任务,构造 tasks与executors之间的映射关系
  4. executor通过heartbeat与coordinator进行通信,获取对应的tasks。每次只获取10个任务。
  5. executor独立管理这些任务,如果任务完成,主动向coordinator更新任务状态

VectorIndexManager:

Add TaskAllocator

This commit introduces the TaskAllocator, a dedicated module for task allocation in the Coordinator. Key aspects of this integration include:

Centralizing Task Allocation: The TaskAllocator encapsulates task allocation logic, which was previously part of the Coordinator. This separation of concerns allows for a more focused and manageable approach to task distribution.

TaskAllocator Design Rationale: The TaskAllocator was designed to provide a clear and isolated domain for managing task allocation. This isolation makes it easier to modify and extend task allocation strategies without impacting the broader system.

Round Robin Allocation: Introduced round robin allocation, leveraging the 'priority-queue' crate for optimized task distribution.

Ingestion Server

摄取服务器完全无状态。它可以根据数据摄取率水平扩展。它严格只执行 I/O,因此可以在生产环境中在廉价硬件上运行。摄取服务器使用协调器来存储有关摄取内容的元数据。

Ingestion Server 公开 HTTP API,供客户端将内容上传到 Indexify。它们还公开检索 API,供客户端检索提取的数据,或搜索向量索引或使用 SQL 查询结构化数据。 Ingestion Server 还从提取器中提取嵌入结构化数据,并将其存储在适当的存储系统中。Ingestion API 位于存储系统之上

  • Blob存储
  • 向量存储
  • 结构化存储

Extractor

提取器是转换非结构化数据或从中提取信息的计算函数。任何用于处理非结构化数据的模型或算法都可以通过实现抽象类(提取器 SDK 的一部分)来实现为提取器。它们可以在任何硬件上运行,单个 Indexify 部署可以在单个集群中支持数十甚至数千个提取器。

它们通过双向 Grpc 流与协调器通信。启动时,它们会向协调器注册其功能并定期发送心跳。当协调器有一些需要在提取器上运行的任务时,它会通过心跳流将任务发送给提取器。 提取器从存储系统下载内容,然后运行其计算功能。任务完成后,任何提取的数据都会上传回摄取服务器,任务结果会通过心跳流发送给协调器。

Extractor:用来进行结构化提取。提取器使用Content包含非结构化数据的原始字节,并从中生成内容和特征的列表。

Content:表示原始的非结构化的数据

Extractor(Content) -> List[Content]
Extractor(Content) -> List[Feature(Type=Metadata)]
Extractor(Content) -> List[Feature(Type=Embedding)]
Extractor(Content) -> List[Feature... Content ...]

例如,PDF 文档可以分解为 - 图像、多个文本块、编码为 JSON 的表格数据、文本块的嵌入,每个都将被编码为内容并由提取器发出。Indexify 将原始字节存储到 blob 存储中,并将嵌入或 JSON 文档等特征存储到索引中以供检索。

提取器通常使用 AI 模型和一些额外的内容预处理和后处理来构建。

竞品

参考: https://news.ycombinator.com/item?id=40458923 (opens in a new tab)

竞品:trellis https://docs.runtrellis.com/docs/getting-started (opens in a new tab) 未开源。