Weaviate Blog

Import & Vectorize Data with Weaviate at Scale

8.5内容质量
Import & Vectorize Data with Weaviate at Scale

TL;DR · AI 摘要

Weaviate官方博客分享了大规模数据导入和向量化实践,重点介绍服务器端批处理、错误处理及媒体处理策略,解决速率限制和批量失败问题。

核心要点

  • 使用Weaviate服务器端批处理可动态调整批次大小,避免速率限制
  • HTTP 200状态码可能隐藏单个对象导入失败问题
  • 媒体文件分块处理可防止Python内存溢出

结构提纲

按章节快速跳转。

  1. 揭示向量数据库导入阶段常见的四个生产级问题

  2. 通过动态调整批次大小实现自动负载均衡

  3. 需要监控失败对象避免重复计算和资源浪费

  4. 分块处理百万级媒体文件防止内存溢出

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • 大规模数据导入实践
    • 服务器端批处理
      • 动态调整批次大小
      • 自动负载均衡
    • 错误处理
      • 监控失败对象
      • 避免重复计算
    • 媒体处理
      • 分块处理
      • 内存优化

金句 / Highlights

值得收藏与分享的关键句。

#Weaviate#向量数据库#数据导入#错误处理
打开原文

使用 Weaviate 批量导入和向量化数据 | Weaviate

使用 Weaviate 批量导入和向量化数据

2026年6月18日

·

12分钟阅读

Ivan Despot

开发者体验工程师

Tommy Smith

多语言工程师

大多数向量数据库试点项目在数据导入阶段失败,而非搜索阶段。你构建了一个精妙的检索流水线,看着它在千份文档上运行良好,但当有人交给你五千万行数据时,接下来的两周会消失在速率限制、部分失败和三次批量逻辑重写中。

这篇博文是我第一次将真实数据集导入 Weaviate 时希望拥有的指南。它涵盖服务器端批量处理、错误处理、最容易出错的数据类型决策,以及如何在不搭建 OCR 流水线的情况下导入媒体和 PDF。

阅读时可尝试以下操作

最快捷的跟进方式是使用 Weaviate Cloud 免费试用版配合 Weaviate Embeddings。无需管理基础设施,无需嵌入 API 密钥,试用版提供免费向量化。以下所有代码片段均可在该设置下直接运行。

没有人警告你的导入问题

一个可用的原型无法告诉你规模化后的表现。困扰生产团队的问题几乎从不会在教程中出现。

以下四个问题最为棘手:

  • 来自嵌入服务提供商的速率限制:大多数团队在真实导入开始一小时内就会触发这些限制,然后编写重试代码,接着还要为重试代码编写重试逻辑。
  • HTTP 200 响应欺骗你:成功的批量响应并不表示每个对象都已正确导入。单个对象可能在绿色状态码下静默失败。
  • 重试时的重复工作:如果你每次重新运行时都生成新 ID,就会重复向量化相同文档并支付双倍费用。
  • 媒体导致的内存爆炸:将一百万张产品图片加载到 Python 列表中,会在批量逻辑运行前就结束你的脚本。

本文其余部分将针对这些问题提供解决方案。

服务器端批量处理

服务器端批量处理是一种流式导入模式,Weaviate 服务器会根据自身的当前负载情况告诉客户端下一步应发送多少数据。你无需猜测批次大小和并发级别,服务器会测量队列深度并在持久连接上应用反向压力。

这很重要,因为正确的批次大小不是固定值。它取决于你拥有的属性数量、文本字段的大小、是否实时向量化、向量化器在后台执行的操作,以及集群正在处理的其他任务。手动调整容易出错。服务器已经掌握了所有这些信息。

以下是 Python 客户端的实现方式:

code
import
weaviate
from
weaviate
.
classes
.
init
import
Auth
client
=
weaviate
.
connect_to_weaviate_cloud
(
cluster_url
=
WCD_URL
,
auth_credentials
=
Auth
.
api_key
(
WCD_API_KEY
)
,
)
collection
=
client
.
collections
.
get
(
"Products"
)
with
collection
.
batch
.
stream
(
)
as
batch
:
for
row
in
iter_rows
(
)
:
batch
.
add_object
(
properties
=
row
)
if
batch
.
number_errors
>
10
:
print
(
"Too many errors, stopping."
)
break
if
collection
.
batch
.
failed_objects
:
print
(
f"
{
len
(
collection
.
batch
.
failed_objects
)
}
objects failed."
)

这就是完整模式。无需指定 batch_size,无需设置 concurrent_requests,无需调整参数。stream() 上下文管理器打开持久连接,服务器设定速率,错误信息异步返回且不中断流程。

错误处理与重试机制 ​

生产环境导入脚本中最常见的错误是将 200 响应视为成功证明。这并不正确。200 状态码表示请求已到达服务器并被接受,但并不表示所有对象都已写入。向量器错误、模式不匹配以及上游嵌入 API 的限流响应都会以单个对象失败的形式出现在整体成功的批量响应中。

Python 客户端在每次批量操作时会暴露三个关键信息:

  • batch.failed_objects — 每个失败的对象,附带错误信息
  • batch.failed_references — 每个失败的交叉引用
  • batch.number_errors — 在上下文管理器中的运行计数

将 failed_objects 视为一个队列进行处理。将其写入文件、重试,如果在相同错误下再次失败,应将其移动到死信队列,避免整个导入过程因单个损坏行而停滞。

以下是生产环境中可复用的重试-检查点模式:

python
import json
from weaviate.util import generate_uuid5

with collection.batch.stream() as batch:
    for row in iter_rows():
        batch.add_object(
            properties=row,
            uuid=generate_uuid5(row["source_id"]),
        )

with open("failed.jsonl", "a") as f:
    for obj in collection.batch.failed_objects:
        f.write(
            json.dumps({
                "properties": obj.object_.properties,
                "error": obj.message,
            }) + "\n"
        )

该模式可安全重复运行的两个关键点:

  1. generate_uuid5 会为相同 source_id 生成相同 UUID,重试时会覆盖而非重复创建
  2. 错误会记录到可重新导入的文件中,修复底层问题后可恢复数据,避免数据丢失或嵌入服务重复计费

常见故障模式及应对方案:

| 症状 | 可能原因 | 解决方案 | |----------------------|------------------------------|----------------------------------| | HTTP 200 但对象缺失 | 向量器请求被限流 | 检查 failed_objects,重试失败子集 | | 客户端内存暴涨 | 在流式传输前加载完整数据集 | 从磁盘或数据库游标流式读取 | | 重试后出现重复对象 | 每次运行生成随机 UUID | 使用 generate_uuid5 从稳定源键生成 | | 导入后向量为空 | 集合未配置向量器模块 | 重新运行前检查集合配置 |

通过 MCP 服务器导入数据

Weaviate 内置了 MCP 服务器(预览版,v1.37.1 新增),允许 LLM 或 IDE 助手(如 Claude Code、Cursor、VS Code)通过 Model Context Protocol 读写实例。启用方式:设置 MCP_SERVER_ENABLED=true,启用写入功能设置 MCP_SERVER_WRITE_ACCESS_ENABLED=true。该服务器提供 weaviate-objects-upsert 工具,可在对话中创建或更新对象。它与 REST API 使用相同端口,遵循 RBAC 策略,无需额外部署。

该工具适用于代理需要在工作过程中写入少量记录的场景:持久化代理记忆、同步小型集合、或在不离开编辑器的情况下修复少量对象。

这不是一个数据导入流水线。每个对象都由模型组装并通过工具调用传递,因此受上下文窗口和单次调用延迟限制,无法使用上述章节中的背压、流式传输或重试-检查点机制。当需要处理数十个以上对象时,应使用 collection.batch.stream()(或客户端的批量 API),将 MCP 服务器留给其设计的对话场景和小批量写入。

导入前选择数据类型

Schema 决策在导入运行后修复的成本是原来的十倍。第一次就做到正确。

导入时最重要的几点:

  • 使用合适的分词方式处理文本。分词决定了哪些 BM25 查询匹配哪些记录。默认方式适用于英文散文。对于产品编码、URL 或任何需要精确匹配字面字符串的场景,请切换到字段分词方式。分词教程会逐步讲解不同方式的权衡。
  • 使用 uuid 作为外键。通过索引实现快速过滤,在插入时进行验证,并在客户端以真正的 UUID 形式呈现,而非字符串。
  • int 与 number 的区分。使用 int 表示计数和 ID,使用 number 表示价格和比率。混淆使用会导致每个查询都需要类型转换。
  • 对实际会进行过滤的关联关系使用引用类型。如果要独立查询相关字段,不要将所有内容扁平化为一个嵌套对象。

完整列表请参见数据类型参考文档。

blobHash:存储嵌入向量,跳过原始字节

如果正在导入媒体文件(图片、音频、视频、PDF 等),这是你需要了解的数据类型。

常规 blob 类型会将完整的 base64 负载存储在磁盘上。blobHash 则不会。它在导入时将原始字节传递给向量化器,使模型能够看到实际媒体内容,然后仅保留 SHA-256 哈希值。向量索引保留嵌入向量,blob 存储保留 32 字节的指纹。

json
{
  "properties": [
    {
      "name": "product_image",
      "dataType": ["blobHash"]
    }
  ]
}

实际影响:10TB 的图片库可以缩小到几个 GB 的哈希值加上向量索引。相似性搜索的表现与使用 blob 时完全相同。你只需支付向量索引的存储成本,原始字节存储在 Weaviate 中的成本被完全避免。将这些原始字节保留在对象存储中才是它们应有的位置。

还有一个额外优势:当更新对象时,新的 base64 会先被哈希并和存储的哈希进行比对,只有在哈希不匹配时才会触发重新向量化。仅凭这一点,当有人错误地重新运行导入流水线时,就能立即获得回报。

无需 OCR 管道的 PDF 向量化

对大多数用户来说,实际问题是:"如何在不编写 OCR 管道的情况下导入 PDF 文件夹?" 最简答案是启动 Weaviate Cloud 试用版并使用 Weaviate Embeddings。

Weaviate Embeddings 集成了专为图像文档检索设计的多模态模型。你只需提供页面图像,它就能生成向量。无需 OCR 步骤,无需布局检测,无需文本提取。表格、图表、扫描表单、多语言文档都可以以相同方式处理。这是云服务专属方案,但它是尝试在真实数据集上进行 PDF 检索的最低摩擦路径,适用于最多几十万页的文档集合,无需你做任何架构决策。

另一个现成方案是 Google 的 multi2vec-google(使用 gemini-embedding-2,3072 维度),该方案在 Weaviate Cloud 中默认启用。它遵循相同的工作流程:你将页面渲染为图像并嵌入这些图像。该模块接受图像输入而非原始 PDF 文件,因此栅格化步骤与 Weaviate Embeddings 完全一致。

如果你在大规模自托管场景中需要处理文档布局,可以查看多向量 ColPali 方案。它使用视觉语言模型为每页生成多个向量,并完全跳过分块处理。虽然涉及更多组件,但这是处理视觉丰富文档检索的最先进方案。

所有三种方案都位于同一位置——模型提供商参考文档。

多模态数据摄入:文本、图像、音频、视频

在 Weaviate 中,多模态功能不是独立的产品。它是集合上的一个向量化模块。您声明集合使用的模型,通过媒体属性导入对象(理想情况下以 blobHash 形式),并使用已有的客户端跨模态查询。

各提供商的覆盖范围:

| 提供商 | 模块 | 文本 | 图像 | 音频 | 视频 | |----------------|-----------------------|------|------|------|------| | Weaviate Embeddings | 原生(WCD) | ✓ | | | | | Google | multi2vec-google | ✓ | ✓ | ✓ | ✓ | | Voyage AI | multi2vec-voyageai | ✓ | ✓ | ✓ | ✓ | | Jina AI | multi2vec-jinaai | ✓ | ✓ | ✓ | ✓ | | Cohere | multi2vec-cohere | ✓ | ✓ | ✓ | ✓ | | NVIDIA | multi2vec-nvidia | ✓ | ✓ | ✓ | ✓ | | CLIP(自托管) | multi2vec-clip | ✓ | ✓ | | | | ImageBind(自托管)| multi2vec-bind | ✓ | ✓ | | |

一个具体场景。您正在为电子商务目录构建搜索功能。每个产品都有名称、描述、三张照片和一个15秒的演示视频。您希望一个查询——"具有主动降噪功能的紧凑无线耳塞"——无论相关信号出现在文本、照片还是视频中,都能找到正确的商品。

您声明一个集合,使用 multi2vec-google 跨命名向量:名称和描述的文本向量,加上产品图片的独立 blobHash 向量和演示视频的另一个向量。每个 blobHash 属性都需要自己的命名向量——一旦 Weaviate 将原始字节替换为哈希值,它们就无法与其他字段一起重新向量化,因此模式将它们隔离。然后单个多目标查询可以同时跨三个向量进行排序:一个集合、一个查询、三种模态——媒体字节在 Weaviate 中永远不会存储两次。

每个提供商的完整设置细节请参见模型提供商参考文档。

在点击开始之前,请检查以下事项

  • 有意识地选择数据类型。对于不需要检索原始字节的媒体使用 blobHash,对于字面字符串重要的文本使用字段分词。
  • 在集合级别而非导入脚本中选择向量化器。Weaviate Embeddings 是 Weaviate Cloud 上最简单的默认选项。
  • 使用确定性 UUID(从稳定源键生成的 generate_uuid5),使重试操作具有幂等性。
  • 如果使用 Python 客户端,请通过 collection.batch.stream() 使用服务器端批处理。
  • 将失败对象记录到死信文件中。不要仅仅打印它们。
  • 对于大型任务,请检查进度,这样崩溃不会导致整个运行失败。
  • 如果大规模导入媒体,请将源字节保留在对象存储中,让 Weaviate 保存哈希值而非文件。

做到这七点,您就不会成为下个季度第三次重写导入脚本的团队。

最快尝试这些方法的方式是免费的 Weaviate Cloud 试用版配合 Weaviate Embeddings。无需管理基础设施,无需嵌入式 API 密钥,本文中的每个代码片段都可以在无需修改的情况下运行。一旦集群部署完成,导入教程是一个很好的下一步选择。

准备开始构建了吗?

查看快速入门教程,或注册免费的 Weaviate Cloud 账户。

GitHub

论坛

X(Twitter)

不想错过其他博客文章?

注册我们的双周通讯以保持更新!

通过提交,我同意

服务条款

.