跳到主要内容
版本:3.0.0

Milvus

Milvus 源连接器

引擎支持​

Spark
Flink
SeaTunnel Zeta

描述​

Milvus 源连接器用于从 Milvus 或 Zilliz Cloud 读取数据。它可以读取一个集合, 也可以读取某个数据库下的所有集合,并且会把 Milvus 的分区信息、向量索引信息等元数据传给下游, 下游连接器支持时可以继续使用这些信息。

常见用法:

  • 配置 collection 读取一个 Milvus 集合。
  • 不配置 collection 时,读取指定数据库下的所有集合。
  • 从 Milvus 复制到 Milvus,并保留向量字段、分区元数据和索引元数据。
  • 读取 FLOAT_VECTOR、BINARY_VECTOR、FLOAT16_VECTOR、BFLOAT16_VECTOR 和 SPARSE_FLOAT_VECTOR 字段。
  • 遇到限流或 gRPC 限制时自动重试。

主要特性​

数据类型映射​

Milvus 数据类型SeaTunnel 数据类型
INT8TINYINT
INT16SMALLINT
INT32INT
INT64BIGINT
FLOATFLOAT
DOUBLEDOUBLE
BOOLBOOLEAN
JSONSTRING
ARRAYARRAY
VARCHARSTRING
FLOAT_VECTORFLOAT_VECTOR
BINARY_VECTORBINARY_VECTOR
FLOAT16_VECTORFLOAT16_VECTOR
BFLOAT16_VECTORBFLOAT16_VECTOR
SPARSE_FLOAT_VECTORSPARSE_FLOAT_VECTOR

源选项​

名称类型是否必传默认值描述
urlString是-Milvus 或 Zilliz Cloud 的连接地址,例如 http://127.0.0.1:19530。
tokenString是-Milvus 认证令牌。本地 Milvus 通常使用 username:password。
databaseString否default源数据库。
collectionString否-源集合。配置后只读取这个集合;不配置时读取 database 下的所有集合。旧别名 collection_name 也仍然支持。
batch_sizeInteger否1000每次从 Milvus 拉取的记录数。值越大吞吐越高,但内存占用也越大;记录中包含较大向量负载时可以适当调小。
rate_limitInteger否1000000Source 每秒最多向 Milvus 请求的记录数。用于在共享 Milvus 配额(QPS)或 gRPC 消息大小限制下对流任务进行限速。设为 -1 关闭限速。

注意事项​

  • database 默认是 default,本地 Milvus 的简单任务通常不用配置。
  • collection 是可选项。只想读一个集合时再配置。
  • batch_size 控制单次拉取的页面大小,与 reader 的并行度无关。需要配合 parallelism 一起调整,以平衡吞吐和内存。
  • rate_limit 是 Milvus 服务端的提示,用于在大批量向量读取时规避 GRPC limit 错误。除非日志里出现限速或 gRPC 报错,否则保持默认值即可。
  • 不配置 collection 时,源端会发现 database 下的所有集合,并把每个集合作为一张独立的 SeaTunnel 表输出。
  • 源端会按 Milvus 分区拆分读取任务。有分区键的集合使用一个 split 读取;没有分区键的集合会按分区名拆分,并分配给多个 reader。
  • 源端读取带分区的集合时,下游 Milvus 接收器可以利用这些元数据在目标集合创建相同分区名。
  • 源端读取到向量索引信息时,下游 Milvus 接收器可以配合 create_index = true 创建相同向量索引。
  • Milvus 源是 BOUNDED(有界)源:作业在所有分区(split)扫描完成后会自然结束,不会像 Kafka、Fluss 那样提供按记录级别 offset 持续增量读取。检查点/恢复以 split(分区)为粒度——已经扫描完成的分区不会被重读,但作业失败时正在扫描的分区会从分区开头重新扫描;如果希望持续摄入新增向量,需要在外部(例如业务写入端)配合周期性重新提交 SeaTunnel 作业。

任务示例​

读取一个集合​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "default"
collection = "simple_example"
}
}

sink {
Console {}
}

读取一个数据库下的所有集合​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "default"
}
}

sink {
Console {}
}

复制一个 Milvus 集合到另一个数据库​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
collection = "simple_example"
}
}

sink {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "test"
collection = "simple_example"
}
}

复制集合并重建向量索引​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
collection = "simple_example"
}
}

sink {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "test_index_preservation"
collection = "simple_example_preservation"
create_index = true
}
}

配合检查点周期性地重跑​

Milvus 源是 BOUNDED 的,所有分区扫描完成后作业会自然结束。本示例以 STREAMING 模式运行,并设置较短的检查点间隔——如果希望持续摄入新增向量, 可以让外部调度器按需重新提交作业;恢复时以 split(分区)为粒度,已经完整扫描 的分区不会被重读,正在扫描时失败的分区会从分区开头重新读取。下游 Sink 设置 enable_upsert = true,配合主键去重避免重复写入。

env {
parallelism = 1
job.mode = "STREAMING"
checkpoint.interval = 30000
}

source {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "streaming_test"
collection = "simple_example"
batch_size = 500
rate_limit = 200000
}
}

sink {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "streaming_test"
enable_upsert = true
batch_size = 1000
}
}

限制读取速度以保护共享集群​

当 Milvus 集群同时被其他任务共享时,可以调小 rate_limit 和 batch_size, 避免 Source 超过集群的 gRPC 消息大小限制。

env {
parallelism = 2
job.mode = "BATCH"
}

source {
Milvus {
url = "http://127.0.0.1:19530"
token = "username:password"
database = "shared"
batch_size = 200
rate_limit = 100000
}
}

sink {
Console {}
}

常见问题​

Milvus Source 能否一次性读取数据库中的所有 Collection?​

可以。如果省略 collection 参数或将其留空,Milvus 源连接器将读取配置的 database 下的所有集合。

支持哪些向量数据类型?​

连接器支持 FLOAT_VECTOR、BINARY_VECTOR、FLOAT16_VECTOR、BFLOAT16_VECTOR 以及 SPARSE_FLOAT_VECTOR,并可将分区与索引元数据透传给下游连接器。

Source 如何处理 gRPC 消息限制或限流错误?​

可以通过调整 batch_size 和 rate_limit 参数来控制读取吞吐。当遇到集群限流或 gRPC 限制时,连接器内置了自动重试与退避机制。

变更日志​

Change Log
ChangeCommitVersion
[Fix][Connector-V2] Accept any JSON root in Milvus sink JSON fields (#11526)https://github.com/apache/seatunnel/commit/9f6b790ff3.0.0
[Improve][Connector-V2] Validate Milvus sink batch size (#11504)https://github.com/apache/seatunnel/commit/dff1b25b13.0.0
[Feature][Connector-V2][Milvus] Support scalar field write null (#11216)https://github.com/apache/seatunnel/commit/e21e4f7ed3.0.0
[Fix][Connector-V2] Fix wrong exception reference in MilvusSourceReader query error handling (#10975)https://github.com/apache/seatunnel/commit/bf96310c73.0.0
[Fix][Connector-V2] Fix SMALLINT switch fall-through in Milvus source converter (#10966)https://github.com/apache/seatunnel/commit/67e74dfa23.0.0
[Improve][connector-milvus] Improved milvus source enumerator splits allocation algorithm for subtasks (#10868)https://github.com/apache/seatunnel/commit/d991e6bbb3.0.0
[Bug][Connector-V2] Fix Milvus sink collection_name target handling (#10893)https://github.com/apache/seatunnel/commit/08f26538c3.0.0
[Fix][API] Support vector index in table schema config parsing (#10582)https://github.com/apache/seatunnel/commit/9fb1a34893.0.0
[Fix][Connect-V2][Milvus] Fix partition creation (#10589)https://github.com/apache/seatunnel/commit/582f9c7183.0.0