Skip to content

Commit 164d138

Browse files
committed
feat: Migrate the milvus data from MinIO to S3
1 parent 3e2b6fa commit 164d138

17 files changed

Lines changed: 766 additions & 0 deletions
Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,19 @@
1+
#!/usr/bin/env bash
2+
# 步骤 0:安装 milvus-backup 并下载配置模板。
3+
set -euo pipefail
4+
source "$(dirname "$0")/env.sh"
5+
6+
cd "${WORKDIR}"
7+
mkdir -p configs logs
8+
9+
PKG="milvus-backup_${BACKUP_OS}_${BACKUP_ARCH}.tar.gz"
10+
echo ">> 下载 ${PKG} (${BACKUP_VERSION})"
11+
curl -fL -o milvus-backup.tar.gz \
12+
"https://github.com/zilliztech/milvus-backup/releases/download/${BACKUP_VERSION}/${PKG}"
13+
14+
tar -xzf milvus-backup.tar.gz
15+
chmod +x milvus-backup
16+
17+
echo ">> 校验"
18+
"${BACKUP_BIN}" --help >/dev/null
19+
echo "milvus-backup 安装完成: ${BACKUP_BIN}"

ch4/minio-s3-demo/01-seed-data.py

Lines changed: 67 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,67 @@
1+
#!/usr/bin/env python3
2+
"""步骤 1:在源 Milvus 创建 collection、灌入测试数据、建索引、flush。
3+
4+
通过环境变量读取配置(见 env.sh)。运行:
5+
python3 01-seed-data.py
6+
"""
7+
import os
8+
import random
9+
10+
from pymilvus import (
11+
connections,
12+
utility,
13+
Collection,
14+
CollectionSchema,
15+
FieldSchema,
16+
DataType,
17+
)
18+
19+
HOST = os.environ.get("SRC_MILVUS_HOST", "localhost")
20+
PORT = os.environ.get("SRC_MILVUS_PORT", "19530")
21+
USER = os.environ.get("SRC_MILVUS_USER", "root")
22+
PASS = os.environ.get("SRC_MILVUS_PASS", "Milvus")
23+
NAME = os.environ.get("TEST_COLLECTION", "migration_demo")
24+
DIM = int(os.environ.get("TEST_DIM", "768"))
25+
ROWS = int(os.environ.get("TEST_ROWS", "10000"))
26+
27+
connections.connect(alias="default", host=HOST, port=PORT, user=USER, password=PASS)
28+
29+
if utility.has_collection(NAME):
30+
print(f"collection {NAME} 已存在,先删除重建")
31+
utility.drop_collection(NAME)
32+
33+
schema = CollectionSchema(
34+
fields=[
35+
FieldSchema("id", DataType.INT64, is_primary=True, auto_id=False),
36+
FieldSchema("label", DataType.INT64),
37+
FieldSchema("embedding", DataType.FLOAT_VECTOR, dim=DIM),
38+
],
39+
description="MinIO->S3 migration demo",
40+
)
41+
coll = Collection(NAME, schema=schema)
42+
print(f"创建 collection: {NAME} (dim={DIM})")
43+
44+
# 分批灌入,避免单次请求过大
45+
BATCH = 2000
46+
inserted = 0
47+
while inserted < ROWS:
48+
n = min(BATCH, ROWS - inserted)
49+
ids = list(range(inserted, inserted + n))
50+
labels = [random.randint(0, 9) for _ in range(n)]
51+
vecs = [[random.random() for _ in range(DIM)] for _ in range(n)]
52+
coll.insert([ids, labels, vecs])
53+
inserted += n
54+
print(f"已灌入 {inserted}/{ROWS}")
55+
56+
print(">> flush")
57+
coll.flush()
58+
59+
print(">> 建索引 (IVF_FLAT / L2)")
60+
coll.create_index(
61+
field_name="embedding",
62+
index_params={"index_type": "IVF_FLAT", "metric_type": "L2", "params": {"nlist": 128}},
63+
)
64+
utility.wait_for_index_building_complete(NAME)
65+
66+
coll.load()
67+
print(f"完成。{NAME} num_entities={coll.num_entities} indexes={[i.index_name for i in coll.indexes]}")

ch4/minio-s3-demo/02-backup.sh

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,72 @@
1+
#!/usr/bin/env bash
2+
# 步骤 2:从源 MinIO 备份到目标 S3(crossStorage)。
3+
set -euo pipefail
4+
source "$(dirname "$0")/env.sh"
5+
cd "${WORKDIR}"
6+
7+
# 可选静态密钥块(IAM 角色时留空)
8+
S3_KEYS=""
9+
if [[ -n "${S3_AK}" ]]; then
10+
S3_KEYS=$(printf ' backupAccessKeyID: "%s"\n backupSecretAccessKey: "%s"' "${S3_AK}" "${S3_SK}")
11+
fi
12+
13+
cat > configs/backup-minio-to-s3.yaml <<YAML
14+
log:
15+
level: info
16+
console: true
17+
file:
18+
filename: "logs/backup.log"
19+
20+
milvus:
21+
address: ${SRC_MILVUS_HOST}
22+
port: ${SRC_MILVUS_PORT}
23+
user: "${SRC_MILVUS_USER}"
24+
password: "${SRC_MILVUS_PASS}"
25+
tlsMode: 0
26+
27+
minio:
28+
# ---- 源 Milvus 当前存储:MinIO ----
29+
storageType: "minio"
30+
address: ${MINIO_HOST}
31+
port: ${MINIO_PORT}
32+
accessKeyID: ${MINIO_AK}
33+
secretAccessKey: ${MINIO_SK}
34+
useSSL: false
35+
useIAM: false
36+
bucketName: "${MINIO_BUCKET}"
37+
rootPath: "${MINIO_ROOTPATH}"
38+
39+
# ---- 备份目标:AWS S3 ----
40+
backupStorageType: "aws"
41+
backupAddress: "${S3_ENDPOINT}"
42+
backupRegion: "${S3_REGION}"
43+
backupPort: ${S3_PORT}
44+
backupUseSSL: true
45+
backupBucketName: "${S3_BUCKET}"
46+
backupRootPath: "${S3_BACKUP_ROOTPATH}"
47+
${S3_KEYS}
48+
49+
# MinIO -> S3 跨对象存储拷贝,必须开启
50+
crossStorage: true
51+
52+
backup:
53+
parallelism:
54+
copydata: 128
55+
backupCollection: 4
56+
backupSegment: 1024
57+
restoreCollection: 2
58+
importJob: 256
59+
YAML
60+
61+
CFG="configs/backup-minio-to-s3.yaml"
62+
echo ">> 连通性检查"
63+
"${BACKUP_BIN}" check --config "${CFG}"
64+
65+
BACKUP_NAME="minio_to_s3_$(date +%Y%m%d_%H%M%S)"
66+
echo "${BACKUP_NAME}" > "${BACKUP_NAME_FILE}"
67+
echo ">> 创建备份: ${BACKUP_NAME} (含索引/rbac)"
68+
"${BACKUP_BIN}" create --config "${CFG}" -n "${BACKUP_NAME}" --rebuild_index --rbac
69+
70+
echo ">> 备份列表"
71+
"${BACKUP_BIN}" list --config "${CFG}"
72+
echo "备份已写入: s3://${S3_BUCKET}/${S3_BACKUP_ROOTPATH}/${BACKUP_NAME}/"

ch4/minio-s3-demo/03-restore.sh

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,79 @@
1+
#!/usr/bin/env bash
2+
# 步骤 3:在连接 S3 的全新目标 Milvus 上恢复(含索引)。
3+
# 前置:目标 Milvus 已部署、etcd 干净、对象存储指向 S3 的 ${S3_DATA_ROOTPATH}。
4+
set -euo pipefail
5+
source "$(dirname "$0")/env.sh"
6+
cd "${WORKDIR}"
7+
8+
if [[ ! -f "${BACKUP_NAME_FILE}" ]]; then
9+
echo "找不到备份名文件 ${BACKUP_NAME_FILE},先跑 02-backup.sh 或手动 export BACKUP_NAME" >&2
10+
exit 1
11+
fi
12+
BACKUP_NAME="$(cat "${BACKUP_NAME_FILE}")"
13+
14+
# 目标 Milvus 存储密钥块(IAM 时留空 -> useIAM:true)
15+
USE_IAM="true"
16+
S3_KEYS=""
17+
if [[ -n "${S3_AK}" ]]; then
18+
USE_IAM="false"
19+
S3_KEYS=$(printf ' accessKeyID: "%s"\n secretAccessKey: "%s"' "${S3_AK}" "${S3_SK}")
20+
fi
21+
BK_KEYS=""
22+
if [[ -n "${S3_AK}" ]]; then
23+
BK_KEYS=$(printf ' backupAccessKeyID: "%s"\n backupSecretAccessKey: "%s"' "${S3_AK}" "${S3_SK}")
24+
fi
25+
26+
cat > configs/restore-s3.yaml <<YAML
27+
log:
28+
level: info
29+
console: true
30+
file:
31+
filename: "logs/restore.log"
32+
33+
milvus:
34+
address: ${DST_MILVUS_HOST}
35+
port: ${DST_MILVUS_PORT}
36+
user: "${DST_MILVUS_USER}"
37+
password: "${DST_MILVUS_PASS}"
38+
tlsMode: 0
39+
40+
minio:
41+
# ---- 目标 Milvus 当前存储:S3 (数据路径) ----
42+
storageType: "aws"
43+
address: "${S3_ENDPOINT}"
44+
region: "${S3_REGION}"
45+
port: ${S3_PORT}
46+
useSSL: true
47+
useIAM: ${USE_IAM}
48+
bucketName: "${S3_BUCKET}"
49+
rootPath: "${S3_DATA_ROOTPATH}"
50+
${S3_KEYS}
51+
52+
# ---- 备份文件位置:同一 S3 (备份路径) ----
53+
backupStorageType: "aws"
54+
backupAddress: "${S3_ENDPOINT}"
55+
backupRegion: "${S3_REGION}"
56+
backupPort: ${S3_PORT}
57+
backupUseSSL: true
58+
backupBucketName: "${S3_BUCKET}"
59+
backupRootPath: "${S3_BACKUP_ROOTPATH}"
60+
${BK_KEYS}
61+
62+
# 目标与备份同一存储,无需跨存储
63+
crossStorage: false
64+
65+
backup:
66+
parallelism:
67+
copydata: 128
68+
restoreCollection: 2
69+
importJob: 256
70+
YAML
71+
72+
CFG="configs/restore-s3.yaml"
73+
echo ">> 连通性检查"
74+
"${BACKUP_BIN}" check --config "${CFG}"
75+
76+
echo ">> 恢复备份: ${BACKUP_NAME} (含索引)"
77+
"${BACKUP_BIN}" restore --config "${CFG}" -n "${BACKUP_NAME}" --restore_index
78+
79+
echo "恢复完成。目标数据应出现在: s3://${S3_BUCKET}/${S3_DATA_ROOTPATH}/"

ch4/minio-s3-demo/04-verify.py

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
#!/usr/bin/env python3
2+
"""步骤 4:验证数据完整性 —— 对比源/目标 Milvus 的 collection、行数、索引,并做检索冒烟。
3+
4+
运行:
5+
python3 04-verify.py
6+
"""
7+
import os
8+
import sys
9+
10+
from pymilvus import connections, utility, Collection
11+
12+
DIM = int(os.environ.get("TEST_DIM", "768"))
13+
TARGET = os.environ.get("TEST_COLLECTION", "migration_demo")
14+
15+
SRC = dict(
16+
host=os.environ.get("SRC_MILVUS_HOST", "localhost"),
17+
port=os.environ.get("SRC_MILVUS_PORT", "19530"),
18+
user=os.environ.get("SRC_MILVUS_USER", "root"),
19+
password=os.environ.get("SRC_MILVUS_PASS", "Milvus"),
20+
)
21+
DST = dict(
22+
host=os.environ.get("DST_MILVUS_HOST", "localhost"),
23+
port=os.environ.get("DST_MILVUS_PORT", "19530"),
24+
user=os.environ.get("DST_MILVUS_USER", "root"),
25+
password=os.environ.get("DST_MILVUS_PASS", "Milvus"),
26+
)
27+
28+
29+
def snapshot(alias, conn):
30+
connections.connect(alias=alias, **conn)
31+
out = {}
32+
for name in utility.list_collections(using=alias):
33+
c = Collection(name, using=alias)
34+
c.load()
35+
out[name] = {
36+
"entities": c.num_entities,
37+
"partitions": len(c.partitions),
38+
"indexes": sorted(i.index_name for i in c.indexes),
39+
"fields": sorted(f.name for f in c.schema.fields),
40+
}
41+
return out
42+
43+
44+
def main():
45+
src = snapshot("src", SRC)
46+
dst = snapshot("dst", DST)
47+
48+
print("=== 源 Milvus ===")
49+
for k, v in src.items():
50+
print(f" {k}: {v}")
51+
print("=== 目标 Milvus ===")
52+
for k, v in dst.items():
53+
print(f" {k}: {v}")
54+
55+
ok = True
56+
# 1) collection 集合一致
57+
if set(src) != set(dst):
58+
print(f"[FAIL] collection 列表不一致: 仅源 {set(src)-set(dst)} / 仅目标 {set(dst)-set(src)}")
59+
ok = False
60+
61+
# 2) 每个 collection 行数/分区/索引/字段一致
62+
for name in set(src) & set(dst):
63+
s, d = src[name], dst[name]
64+
for key in ("entities", "partitions", "indexes", "fields"):
65+
if s[key] != d[key]:
66+
print(f"[FAIL] {name}.{key}: 源={s[key]} 目标={d[key]}")
67+
ok = False
68+
69+
# 3) 目标检索冒烟
70+
if TARGET in dst:
71+
c = Collection(TARGET, using="dst")
72+
c.load()
73+
res = c.search(
74+
data=[[0.0] * DIM],
75+
anns_field="embedding",
76+
param={"metric_type": "L2", "params": {"nprobe": 10}},
77+
limit=5,
78+
)
79+
hits = len(res[0])
80+
print(f"[search] {TARGET} 返回 {hits} 条")
81+
if hits == 0:
82+
print(f"[FAIL] {TARGET} 检索返回 0 条")
83+
ok = False
84+
85+
print("\n结果:", "PASS ✅" if ok else "FAIL ❌")
86+
sys.exit(0 if ok else 1)
87+
88+
89+
if __name__ == "__main__":
90+
main()

0 commit comments

Comments
 (0)