MinIO 对象存储:大文件和模型制品管理
Tags: #MinIO #ObjectStorage #ModelArtifacts #PresignedURL #BucketPerProject #智策平台
“系列:S3 数据基座 · 第 4 篇 | 难度:高级 | 阅读时间:20 分钟
MinIO 对象存储:大文件和模型制品管理
Tags: #MinIO #ObjectStorage #ModelArtifacts #PresignedURL #BucketPerProject #智策平台
#TL;DR
在智策平台(coomia-dip)中,结构化数据通过 Doris/Iceberg 管理,但大文件(Pipeline 输出、导出报表、ML 模型制品、用户上传附件)需要高性能对象存储支撑。我们选择 MinIO 作为 S3 兼容对象存储层,采用 bucket-per-project 的租户隔离策略,结合 Presigned URL 实现安全的前端直传,通过生命周期策略自动清理临时文件。本文详细介绍从部署架构到 SDK 集成的完整实践,展示如何在 coomia-dip 中构建统一的二进制资产管理层。
#1. 为什么需要独立的对象存储层
#1.1 结构化 vs 非结构化数据的分野
在 coomia-dip 数据基座中,数据可以分为两大类:
数据分类与存储映射:
┌─────────────────────────────────────────────────────┐
│ coomia-dip 数据层 │
│ │
│ ┌─────────────────────┐ ┌─────────────────────────┐ │
│ │ 结构化数据 │ │ 非结构化数据 │ │
│ │ │ │ │ │
│ │ • Ontology 实体 │ │ • Pipeline 输出文件 │ │
│ │ • 关系边 │ │ • ML 模型权重 │ │
│ │ • 事件流 │ │ • 导出报表 (CSV/Excel) │ │
│ │ • 指标时序 │ │ • 用户上传附件 │ │
│ │ • 审计日志 │ │ • 可视化截图 │ │
│ │ │ │ • 大型 JSON/Parquet │ │
│ │ 存储:Doris + Iceberg│ │ 存储:MinIO (S3 兼容) │ │
│ └─────────────────────┘ └─────────────────────────┘ │
└─────────────────────────────────────────────────────┘
#1.2 为什么不用本地文件系统或 NFS
| 方案 | 问题 | 致命缺陷 |
|---|---|---|
| 本地文件系统 | 单节点瓶颈 | 节点故障即数据丢失 |
| NFS | 性能差,锁竞争 | 无法水平扩展 |
| HDFS | 运维复杂,资源占用大 | 对中小规模过度设计 |
| 云 S3 | 外部依赖,网络延迟 | 私有化部署无法使用 |
| MinIO | S3 兼容,轻量高性能 | 最佳选择 |
#1.3 MinIO 的关键优势
MinIO 是一个高性能 S3 兼容对象存储系统,在 coomia-dip 场景中具备以下优势:
- S3 API 100% 兼容:所有 S3 SDK/工具直接可用
- 纠删码(Erasure Coding):数据冗余无需 RAID
- 分布式部署:多节点高可用
- 轻量高效:单二进制文件,Go 语言实现
- 对象锁和版本控制:合规和审计需求
- 生命周期管理:自动过期和清理
#2. 部署架构设计
#2.1 集群拓扑
MinIO 集群部署拓扑(coomia-dip 生产环境):
┌──────────────────────────────────────────────────────────┐
│ Nginx / Traefik │
│ (TLS 终止 + 负载均衡) │
│ :9000 (API) :9001 (Console) │
└─────────┬──────────────────────┬─────────────────────────┘
│ │
┌─────┴─────┐ ┌─────┴─────┐
│ │ │ │
┌───┴───┐ ┌───┴───┐ ┌───┴───┐ ┌───┴───┐
│MinIO-1│ │MinIO-2│ │MinIO-3│ │MinIO-4│
│ Node │ │ Node │ │ Node │ │ Node │
│ │ │ │ │ │ │ │
│/data1 │ │/data1 │ │/data1 │ │/data1 │
│/data2 │ │/data2 │ │/data2 │ │/data2 │
└───────┘ └───────┘ └───────┘ └───────┘
纠删码组:EC(4,2) — 4 数据块 + 2 校验块
容忍:任意 2 节点故障不丢数据
#2.2 Docker Compose 配置
# deployment-Layer/docker-compose-minio.yml
version: '3.8'
services:
minio-1:
image: quay.io/minio/minio:RELEASE.2024-06-13T22-53-53Z
command: server --console-address ":9001" http://minio-{1...4}/data{1...2}
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD}
MINIO_PROMETHEUS_AUTH_TYPE: public
MINIO_SCANNER_SPEED: slow
volumes:
- minio1-data1:/data1
- minio1-data2:/data2
networks:
- coomia-dip-net
healthcheck:
test: ["CMD", "mc", "ready", "local"]
interval: 30s
timeout: 10s
retries: 3
deploy:
resources:
limits:
memory: 4G
reservations:
memory: 2G
minio-2:
image: quay.io/minio/minio:RELEASE.2024-06-13T22-53-53Z
command: server --console-address ":9001" http://minio-{1...4}/data{1...2}
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD}
volumes:
- minio2-data1:/data1
- minio2-data2:/data2
networks:
- coomia-dip-net
minio-3:
image: quay.io/minio/minio:RELEASE.2024-06-13T22-53-53Z
command: server --console-address ":9001" http://minio-{1...4}/data{1...2}
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD}
volumes:
- minio3-data1:/data1
- minio3-data2:/data2
networks:
- coomia-dip-net
minio-4:
image: quay.io/minio/minio:RELEASE.2024-06-13T22-53-53Z
command: server --console-address ":9001" http://minio-{1...4}/data{1...2}
environment:
MINIO_ROOT_USER: ${MINIO_ROOT_USER}
MINIO_ROOT_PASSWORD: ${MINIO_ROOT_PASSWORD}
volumes:
- minio4-data1:/data1
- minio4-data2:/data2
networks:
- coomia-dip-net
volumes:
minio1-data1:
minio1-data2:
minio2-data1:
minio2-data2:
minio3-data1:
minio3-data2:
minio4-data1:
minio4-data2:
networks:
coomia-dip-net:
external: true
#2.3 开发环境单节点部署
# deployment-Layer/docker-compose-minio-dev.yml
version: '3.8'
services:
minio:
image: quay.io/minio/minio:RELEASE.2024-06-13T22-53-53Z
command: server /data --console-address ":9001"
ports:
- "9000:9000"
- "9001:9001"
environment:
MINIO_ROOT_USER: coomia-dip
MINIO_ROOT_PASSWORD: coomia-dip-dev-2024
volumes:
- minio-data:/data
networks:
- coomia-dip-net
volumes:
minio-data:
#3. Bucket-Per-Project 租户隔离策略
#3.1 Bucket 命名规范
coomia-dip 采用 bucket-per-project 的隔离策略,每个项目(World)拥有独立的 bucket:
Bucket 命名规范:
onto-{project_id}-{category}
示例:
onto-proj001-pipeline # Pipeline 输出
onto-proj001-artifacts # ML 模型制品
onto-proj001-exports # 导出文件
onto-proj001-uploads # 用户上传
onto-system-shared # 系统共享资源
onto-system-templates # 模板文件
#3.2 目录结构约定
Bucket 内部目录结构:
onto-proj001-pipeline/
├── runs/
│ ├── run-20240615-001/
│ │ ├── output/
│ │ │ ├── result.parquet
│ │ │ └── summary.json
│ │ ├── logs/
│ │ │ ├── stdout.log
│ │ │ └── stderr.log
│ │ └── _metadata.json
│ └── run-20240615-002/
│ └── ...
└── snapshots/
└── 2024-06-15/
└── full-export.parquet
onto-proj001-artifacts/
├── models/
│ ├── risk-scorer-v1.0/
│ │ ├── model.onnx
│ │ ├── config.json
│ │ ├── tokenizer.json
│ │ └── _manifest.json
│ └── anomaly-detector-v2.1/
│ └── ...
├── checkpoints/
│ └── training-run-001/
│ ├── epoch-10.pt
│ └── epoch-20.pt
└── evaluations/
└── risk-scorer-v1.0/
├── metrics.json
└── confusion-matrix.png
#3.3 IAM 策略与权限隔离
# python-sdk/ontology_sdk/storage/minio_policy.py
from dataclasses import dataclass
from typing import Any
@dataclass
class BucketPolicy:
"""MinIO bucket 访问策略生成器"""
project_id: str
def generate_readwrite_policy(self) -> dict[str, Any]:
"""生成项目读写策略"""
return {
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:PutObject",
"s3:DeleteObject",
"s3:ListBucket",
],
"Resource": [
f"arn:aws:s3:::onto-{self.project_id}-*",
f"arn:aws:s3:::onto-{self.project_id}-*/*",
],
},
{
"Effect": "Deny",
"Action": ["s3:*"],
"Resource": ["arn:aws:s3:::onto-system-*"],
},
],
}
def generate_readonly_policy(self) -> dict[str, Any]:
"""生成项目只读策略"""
return {
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:ListBucket",
],
"Resource": [
f"arn:aws:s3:::onto-{self.project_id}-*",
f"arn:aws:s3:::onto-{self.project_id}-*/*",
],
}
],
}
def generate_pipeline_policy(self) -> dict[str, Any]:
"""生成 Pipeline 专用策略(只能写 pipeline bucket)"""
return {
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": [
"s3:GetObject",
"s3:PutObject",
"s3:ListBucket",
],
"Resource": [
f"arn:aws:s3:::onto-{self.project_id}-pipeline",
f"arn:aws:s3:::onto-{self.project_id}-pipeline/*",
],
},
{
"Effect": "Allow",
"Action": ["s3:GetObject"],
"Resource": [
f"arn:aws:s3:::onto-{self.project_id}-artifacts/*",
],
},
],
}
#4. Pipeline 输出管理
#4.1 Pipeline 与 MinIO 的集成架构
Pipeline 输出流程:
┌──────────┐ ┌──────────┐ ┌──────────────┐ ┌─────────┐
│ Pipeline │ │ Compute │ │ MinIO │ │ Doris │
│ Engine │───>│ Worker │───>│ Storage │───>│ Catalog │
│(Schedule)│ │(执行节点)│ │ (大文件) │ │ (元数据)│
└──────────┘ └────┬─────┘ └──────────────┘ └─────────┘
│
│ 同时写入
▼
┌──────────────┐
│ Iceberg Table│
│ (结构化结果)│
└──────────────┘
#4.2 Pipeline 输出写入器
# intelligence-Layer/pipeline/output_writer.py
import json
import hashlib
from datetime import datetime, timezone
from pathlib import PurePosixPath
from typing import BinaryIO
from minio import Minio
from pydantic import BaseModel, Field
class PipelineOutputMeta(BaseModel):
"""Pipeline 输出元数据"""
run_id: str
pipeline_id: str
project_id: str
output_type: str # "parquet", "csv", "json", "binary"
file_name: str
file_size: int
md5_hash: str
created_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
tags: dict[str, str] = Field(default_factory=dict)
class PipelineOutputWriter:
"""Pipeline 输出到 MinIO 的写入器"""
def __init__(self, minio_client: Minio):
self._client = minio_client
def write_output(
self,
project_id: str,
run_id: str,
pipeline_id: str,
file_name: str,
data: BinaryIO,
content_type: str = "application/octet-stream",
tags: dict[str, str] | None = None,
) -> PipelineOutputMeta:
"""写入 Pipeline 输出文件到 MinIO"""
bucket_name = f"onto-{project_id}-pipeline"
self._ensure_bucket(bucket_name)
# 构建对象路径
object_path = str(
PurePosixPath("runs") / run_id / "output" / file_name
)
# 读取数据并计算 hash
content = data.read()
md5_hash = hashlib.md5(content).hexdigest()
# 上传到 MinIO
from io import BytesIO
self._client.put_object(
bucket_name=bucket_name,
object_name=object_path,
data=BytesIO(content),
length=len(content),
content_type=content_type,
metadata={
"x-amz-meta-run-id": run_id,
"x-amz-meta-pipeline-id": pipeline_id,
"x-amz-meta-md5": md5_hash,
},
)
# 构建元数据
meta = PipelineOutputMeta(
run_id=run_id,
pipeline_id=pipeline_id,
project_id=project_id,
output_type=file_name.rsplit(".", 1)[-1] if "." in file_name else "binary",
file_name=file_name,
file_size=len(content),
md5_hash=md5_hash,
tags=tags or {},
)
# 写入元数据文件
meta_path = str(
PurePosixPath("runs") / run_id / "_metadata.json"
)
meta_bytes = meta.model_dump_json(indent=2).encode("utf-8")
self._client.put_object(
bucket_name=bucket_name,
object_name=meta_path,
data=BytesIO(meta_bytes),
length=len(meta_bytes),
content_type="application/json",
)
return meta
def write_log(
self,
project_id: str,
run_id: str,
log_type: str, # "stdout" or "stderr"
content: str,
) -> None:
"""写入 Pipeline 执行日志"""
bucket_name = f"onto-{project_id}-pipeline"
self._ensure_bucket(bucket_name)
object_path = str(
PurePosixPath("runs") / run_id / "logs" / f"{log_type}.log"
)
data = content.encode("utf-8")
from io import BytesIO
self._client.put_object(
bucket_name=bucket_name,
object_name=object_path,
data=BytesIO(data),
length=len(data),
content_type="text/plain",
)
def _ensure_bucket(self, bucket_name: str) -> None:
"""确保 bucket 存在"""
if not self._client.bucket_exists(bucket_name):
self._client.make_bucket(bucket_name)
#5. ML 模型制品管理
#5.1 模型制品生命周期
ML 模型制品生命周期:
训练阶段 注册阶段 部署阶段 归档阶段
┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ Training │──────>│ Registry │──────>│ Serving │──────>│ Archive │
│ │ │ │ │ │ │ │
│ Checkpts │ │ Version │ │ Active │ │ Cold │
│ in MinIO │ │ Tag │ │ Download │ │ Storage │
└──────────┘ └──────────┘ └──────────┘ └──────────┘
│ │ │ │
▼ ▼ ▼ ▼
artifacts/ artifacts/ Presigned Lifecycle
checkpoints/ models/ URL 下载 Rule 迁移
train-run-*/ model-v*/ (限时访问) (90天后)
#5.2 模型注册表实现
# intelligence-Layer/ml/model_registry.py
import json
from datetime import datetime, timezone, timedelta
from enum import Enum
from io import BytesIO
from typing import BinaryIO
from minio import Minio
from pydantic import BaseModel, Field
class ModelStage(str, Enum):
"""模型阶段"""
DEVELOPMENT = "development"
STAGING = "staging"
PRODUCTION = "production"
ARCHIVED = "archived"
class ModelManifest(BaseModel):
"""模型清单"""
model_name: str
version: str
stage: ModelStage = ModelStage.DEVELOPMENT
framework: str # "pytorch", "onnx", "sklearn", "xgboost"
description: str = ""
metrics: dict[str, float] = Field(default_factory=dict)
parameters: dict[str, str] = Field(default_factory=dict)
files: list[str] = Field(default_factory=list)
created_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
created_by: str = ""
tags: dict[str, str] = Field(default_factory=dict)
class ModelRegistry:
"""基于 MinIO 的模型注册表"""
def __init__(self, minio_client: Minio):
self._client = minio_client
def register_model(
self,
project_id: str,
model_name: str,
version: str,
framework: str,
files: dict[str, BinaryIO],
metrics: dict[str, float] | None = None,
parameters: dict[str, str] | None = None,
description: str = "",
created_by: str = "",
) -> ModelManifest:
"""注册新模型版本"""
bucket_name = f"onto-{project_id}-artifacts"
if not self._client.bucket_exists(bucket_name):
self._client.make_bucket(bucket_name)
# 上传所有模型文件
file_list: list[str] = []
for filename, fileobj in files.items():
object_path = f"models/{model_name}-{version}/{filename}"
content = fileobj.read()
self._client.put_object(
bucket_name=bucket_name,
object_name=object_path,
data=BytesIO(content),
length=len(content),
)
file_list.append(filename)
# 创建并上传 manifest
manifest = ModelManifest(
model_name=model_name,
version=version,
framework=framework,
description=description,
metrics=metrics or {},
parameters=parameters or {},
files=file_list,
created_by=created_by,
)
manifest_path = f"models/{model_name}-{version}/_manifest.json"
manifest_bytes = manifest.model_dump_json(indent=2).encode("utf-8")
self._client.put_object(
bucket_name=bucket_name,
object_name=manifest_path,
data=BytesIO(manifest_bytes),
length=len(manifest_bytes),
content_type="application/json",
)
return manifest
def promote_model(
self,
project_id: str,
model_name: str,
version: str,
target_stage: ModelStage,
) -> ModelManifest:
"""提升模型到指定阶段"""
bucket_name = f"onto-{project_id}-artifacts"
manifest_path = f"models/{model_name}-{version}/_manifest.json"
# 读取现有 manifest
response = self._client.get_object(bucket_name, manifest_path)
manifest_data = json.loads(response.read())
response.close()
response.release_conn()
manifest = ModelManifest(**manifest_data)
manifest.stage = target_stage
# 更新 manifest
manifest_bytes = manifest.model_dump_json(indent=2).encode("utf-8")
self._client.put_object(
bucket_name=bucket_name,
object_name=manifest_path,
data=BytesIO(manifest_bytes),
length=len(manifest_bytes),
content_type="application/json",
)
return manifest
def get_production_model(
self, project_id: str, model_name: str
) -> ModelManifest | None:
"""获取生产阶段的模型"""
bucket_name = f"onto-{project_id}-artifacts"
prefix = f"models/{model_name}-"
versions: list[ModelManifest] = []
for obj in self._client.list_objects(bucket_name, prefix=prefix, recursive=True):
if obj.object_name and obj.object_name.endswith("/_manifest.json"):
response = self._client.get_object(bucket_name, obj.object_name)
manifest_data = json.loads(response.read())
response.close()
response.release_conn()
manifest = ModelManifest(**manifest_data)
if manifest.stage == ModelStage.PRODUCTION:
versions.append(manifest)
if not versions:
return None
return sorted(versions, key=lambda m: m.created_at, reverse=True)[0]
def download_model_file(
self,
project_id: str,
model_name: str,
version: str,
filename: str,
) -> bytes:
"""下载模型文件"""
bucket_name = f"onto-{project_id}-artifacts"
object_path = f"models/{model_name}-{version}/{filename}"
response = self._client.get_object(bucket_name, object_path)
data = response.read()
response.close()
response.release_conn()
return data
#6. Presigned URL:安全的前端直传
#6.1 Presigned URL 工作原理
Presigned URL 上传流程:
┌──────────┐ ┌──────────────┐ ┌──────────┐
│ Browser │ │ coomia-dip │ │ MinIO │
│ (前端) │ │ API Gateway │ │ Server │
└────┬─────┘ └──────┬───────┘ └────┬─────┘
│ │ │
│ 1. 请求上传URL │ │
│ ─────────────────>│ │
│ │ │
│ │ 2. 验证权限 │
│ │ 生成Presigned │
│ │ URL │
│ │ │
│ 3. 返回Presigned │ │
│ URL + headers │ │
│ <─────────────────│ │
│ │ │
│ 4. PUT 文件到 │ │
│ Presigned URL │ │
│ ──────────────────────────────────> │
│ │ │
│ 5. 200 OK │ │
│ <────────────────────────────────── │
│ │ │
│ 6. 通知上传完成 │ │
│ ─────────────────>│ │
│ │ 7. 验证对象存在 │
│ │ ────────────────>│
│ │ │
│ 8. 确认完成 │ │
│ <─────────────────│ │
│ │ │
#6.2 Presigned URL 服务实现
# control-Layer/api/storage_service.py (gRPC 调用 Python 实现)
from datetime import timedelta
from urllib.parse import urlparse
from minio import Minio
from pydantic import BaseModel
class PresignedUploadResponse(BaseModel):
"""Presigned 上传响应"""
upload_url: str
object_key: str
expires_in_seconds: int
required_headers: dict[str, str]
class PresignedDownloadResponse(BaseModel):
"""Presigned 下载响应"""
download_url: str
expires_in_seconds: int
file_name: str
file_size: int | None = None
class StorageService:
"""对象存储服务"""
def __init__(self, minio_client: Minio, external_endpoint: str):
self._client = minio_client
self._external_endpoint = external_endpoint
def generate_upload_url(
self,
project_id: str,
category: str,
object_path: str,
content_type: str = "application/octet-stream",
max_size_mb: int = 500,
expires_minutes: int = 30,
) -> PresignedUploadResponse:
"""生成 Presigned 上传 URL"""
bucket_name = f"onto-{project_id}-{category}"
if not self._client.bucket_exists(bucket_name):
self._client.make_bucket(bucket_name)
# 生成 presigned PUT URL
url = self._client.presigned_put_object(
bucket_name=bucket_name,
object_name=object_path,
expires=timedelta(minutes=expires_minutes),
)
# 替换内部地址为外部地址
url = self._replace_endpoint(url)
return PresignedUploadResponse(
upload_url=url,
object_key=f"{bucket_name}/{object_path}",
expires_in_seconds=expires_minutes * 60,
required_headers={
"Content-Type": content_type,
"x-amz-meta-project-id": project_id,
},
)
def generate_download_url(
self,
project_id: str,
category: str,
object_path: str,
expires_minutes: int = 60,
filename_override: str | None = None,
) -> PresignedDownloadResponse:
"""生成 Presigned 下载 URL"""
bucket_name = f"onto-{project_id}-{category}"
# 获取对象信息
stat = self._client.stat_object(bucket_name, object_path)
# 生成 presigned GET URL
extra_query_params = {}
if filename_override:
extra_query_params["response-content-disposition"] = (
f'attachment; filename="{filename_override}"'
)
url = self._client.presigned_get_object(
bucket_name=bucket_name,
object_name=object_path,
expires=timedelta(minutes=expires_minutes),
extra_query_params=extra_query_params if extra_query_params else None,
)
url = self._replace_endpoint(url)
return PresignedDownloadResponse(
download_url=url,
expires_in_seconds=expires_minutes * 60,
file_name=filename_override or object_path.rsplit("/", 1)[-1],
file_size=stat.size,
)
def _replace_endpoint(self, url: str) -> str:
"""替换内部端点为外部可访问端点"""
parsed = urlparse(url)
return url.replace(
f"{parsed.scheme}://{parsed.netloc}",
self._external_endpoint,
)
#6.3 前端上传组件集成
// 前端上传示例(TypeScript)
interface PresignedUploadResponse {
upload_url: string;
object_key: string;
expires_in_seconds: number;
required_headers: Record<string, string>;
}
async function uploadFileToMinIO(
file: File,
projectId: string,
category: string,
): Promise<string> {
// Step 1: 从 coomia-dip API 获取 presigned URL
const response = await fetch('/api/v1/storage/upload-url', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
project_id: projectId,
category: category,
object_path: `uploads/${Date.now()}-${file.name}`,
content_type: file.type,
}),
});
const presigned: PresignedUploadResponse = await response.json();
// Step 2: 直接上传到 MinIO(绕过应用服务器)
const uploadResponse = await fetch(presigned.upload_url, {
method: 'PUT',
headers: {
...presigned.required_headers,
'Content-Type': file.type,
},
body: file,
});
if (!uploadResponse.ok) {
throw new Error(`Upload failed: ${uploadResponse.statusText}`);
}
// Step 3: 通知应用服务器上传完成
await fetch('/api/v1/storage/upload-complete', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ object_key: presigned.object_key }),
});
return presigned.object_key;
}
#7. 生命周期管理与自动清理
#7.1 生命周期策略
# deployment-Layer/scripts/minio_lifecycle_setup.py
"""MinIO 生命周期策略配置脚本"""
from minio import Minio
from minio.lifecycleconfig import (
LifecycleConfig,
Rule,
Filter,
Expiration,
Transition,
)
def setup_lifecycle_policies(client: Minio, project_id: str) -> None:
"""配置项目的生命周期策略"""
# Pipeline bucket: 运行日志 30 天过期,输出 90 天过期
pipeline_bucket = f"onto-{project_id}-pipeline"
if client.bucket_exists(pipeline_bucket):
config = LifecycleConfig(
[
Rule(
rule_id="expire-logs-30d",
status="Enabled",
rule_filter=Filter(prefix="runs/*/logs/"),
expiration=Expiration(days=30),
),
Rule(
rule_id="expire-outputs-90d",
status="Enabled",
rule_filter=Filter(prefix="runs/*/output/"),
expiration=Expiration(days=90),
),
Rule(
rule_id="keep-snapshots",
status="Enabled",
rule_filter=Filter(prefix="snapshots/"),
expiration=Expiration(days=365),
),
]
)
client.set_bucket_lifecycle(pipeline_bucket, config)
# Artifacts bucket: checkpoints 14 天过期,模型保留
artifacts_bucket = f"onto-{project_id}-artifacts"
if client.bucket_exists(artifacts_bucket):
config = LifecycleConfig(
[
Rule(
rule_id="expire-checkpoints-14d",
status="Enabled",
rule_filter=Filter(prefix="checkpoints/"),
expiration=Expiration(days=14),
),
Rule(
rule_id="expire-evaluations-180d",
status="Enabled",
rule_filter=Filter(prefix="evaluations/"),
expiration=Expiration(days=180),
),
]
)
client.set_bucket_lifecycle(artifacts_bucket, config)
# Exports bucket: 所有文件 7 天过期
exports_bucket = f"onto-{project_id}-exports"
if client.bucket_exists(exports_bucket):
config = LifecycleConfig(
[
Rule(
rule_id="expire-exports-7d",
status="Enabled",
rule_filter=Filter(prefix=""),
expiration=Expiration(days=7),
),
]
)
client.set_bucket_lifecycle(exports_bucket, config)
def setup_versioning(client: Minio, project_id: str) -> None:
"""为关键 bucket 启用版本控制"""
from minio.versioningconfig import VersioningConfig, ENABLED
artifacts_bucket = f"onto-{project_id}-artifacts"
if client.bucket_exists(artifacts_bucket):
client.set_bucket_versioning(
artifacts_bucket, VersioningConfig(ENABLED)
)
#7.2 存储使用量监控
# python-sdk/ontology_sdk/storage/usage_monitor.py
from dataclasses import dataclass
from minio import Minio
@dataclass
class BucketUsage:
"""Bucket 使用量统计"""
bucket_name: str
object_count: int
total_size_bytes: int
largest_object_bytes: int
largest_object_name: str
@property
def total_size_mb(self) -> float:
return self.total_size_bytes / (1024 * 1024)
@property
def total_size_gb(self) -> float:
return self.total_size_bytes / (1024 * 1024 * 1024)
class UsageMonitor:
"""存储使用量监控"""
def __init__(self, minio_client: Minio):
self._client = minio_client
def get_bucket_usage(self, bucket_name: str) -> BucketUsage:
"""获取单个 bucket 的使用量"""
total_size = 0
count = 0
largest_size = 0
largest_name = ""
for obj in self._client.list_objects(bucket_name, recursive=True):
count += 1
size = obj.size or 0
total_size += size
if size > largest_size:
largest_size = size
largest_name = obj.object_name or ""
return BucketUsage(
bucket_name=bucket_name,
object_count=count,
total_size_bytes=total_size,
largest_object_bytes=largest_size,
largest_object_name=largest_name,
)
def get_project_usage(self, project_id: str) -> list[BucketUsage]:
"""获取项目所有 bucket 的使用量"""
results: list[BucketUsage] = []
for bucket in self._client.list_buckets():
if bucket.name and bucket.name.startswith(f"onto-{project_id}-"):
results.append(self.get_bucket_usage(bucket.name))
return results
def check_quota(
self, project_id: str, quota_gb: float = 100.0
) -> tuple[bool, float]:
"""检查项目是否超出配额"""
usages = self.get_project_usage(project_id)
total_gb = sum(u.total_size_gb for u in usages)
return total_gb <= quota_gb, total_gb
#8. 与 Ontology 的集成
#8.1 对象引用作为实体属性
在 coomia-dip 的 Ontology 模型中,大文件通过对象引用(Object Reference)关联到实体:
-- entity_common 中的对象引用字段
CREATE TABLE entity_common (
entity_id VARCHAR(64) NOT NULL,
entity_type VARCHAR(128) NOT NULL,
display_name VARCHAR(512),
properties JSON,
-- 对象引用存储在 properties 中
-- {"model_artifact": "s3://onto-proj001-artifacts/models/risk-v1/model.onnx"}
-- {"report_file": "s3://onto-proj001-exports/2024-06/monthly.pdf"}
created_at DATETIME NOT NULL,
updated_at DATETIME NOT NULL
);
#8.2 对象引用解析器
# python-sdk/ontology_sdk/storage/object_ref.py
import re
from dataclasses import dataclass
@dataclass
class ObjectReference:
"""S3 对象引用"""
bucket: str
key: str
version_id: str | None = None
@classmethod
def parse(cls, uri: str) -> "ObjectReference":
"""解析 s3:// URI"""
pattern = r"^s3://([^/]+)/(.+?)(?:\?versionId=(.+))?$"
match = re.match(pattern, uri)
if not match:
raise ValueError(f"Invalid S3 URI: {uri}")
return cls(
bucket=match.group(1),
key=match.group(2),
version_id=match.group(3),
)
def to_uri(self) -> str:
"""转换为 s3:// URI"""
base = f"s3://{self.bucket}/{self.key}"
if self.version_id:
base += f"?versionId={self.version_id}"
return base
@property
def project_id(self) -> str:
"""从 bucket 名称提取项目 ID"""
# onto-{project_id}-{category}
parts = self.bucket.split("-")
if len(parts) >= 3 and parts[0] == "onto":
return parts[1]
raise ValueError(f"Cannot extract project_id from bucket: {self.bucket}")
@property
def category(self) -> str:
"""从 bucket 名称提取类别"""
parts = self.bucket.split("-")
if len(parts) >= 3 and parts[0] == "onto":
return "-".join(parts[2:])
raise ValueError(f"Cannot extract category from bucket: {self.bucket}")
#9. 大文件分片上传
#9.1 分片上传策略
对于超过 100MB 的大文件,MinIO 支持分片上传(Multipart Upload):
分片上传流程:
┌──────────┐ ┌──────────┐
│ Client │ │ MinIO │
└────┬─────┘ └────┬─────┘
│ 1. InitiateMultipartUpload │
│ ────────────────────────────────────────>│
│ │
│ 2. UploadId │
│ <────────────────────────────────────────│
│ │
│ 3. UploadPart (Part 1: 0-100MB) │
│ ────────────────────────────────────────>│
│ 4. ETag for Part 1 │
│ <────────────────────────────────────────│
│ │
│ 5. UploadPart (Part 2: 100-200MB) │
│ ────────────────────────────────────────>│
│ 6. ETag for Part 2 │
│ <────────────────────────────────────────│
│ │
│ ... (并行上传多个分片) │
│ │
│ N. CompleteMultipartUpload │
│ [Part1:ETag1, Part2:ETag2, ...] │
│ ────────────────────────────────────────>│
│ │
│ N+1. Object Created │
│ <────────────────────────────────────────│
#9.2 分片上传实现
# python-sdk/ontology_sdk/storage/multipart_upload.py
import os
import hashlib
from concurrent.futures import ThreadPoolExecutor, as_completed
from dataclasses import dataclass, field
from pathlib import Path
from typing import Callable
from minio import Minio
PART_SIZE = 100 * 1024 * 1024 # 100 MB per part
@dataclass
class UploadProgress:
"""上传进度"""
total_bytes: int
uploaded_bytes: int = 0
parts_completed: int = 0
total_parts: int = 0
@property
def percentage(self) -> float:
if self.total_bytes == 0:
return 100.0
return (self.uploaded_bytes / self.total_bytes) * 100
@dataclass
class MultipartUploader:
"""大文件分片上传器"""
minio_client: Minio
max_workers: int = 4
part_size: int = PART_SIZE
def upload_large_file(
self,
bucket_name: str,
object_name: str,
file_path: str | Path,
content_type: str = "application/octet-stream",
progress_callback: Callable[[UploadProgress], None] | None = None,
) -> str:
"""分片上传大文件,返回 ETag"""
file_path = Path(file_path)
file_size = file_path.stat().st_size
if file_size <= self.part_size:
# 小文件直接上传
self.minio_client.fput_object(
bucket_name, object_name, str(file_path),
content_type=content_type,
)
return hashlib.md5(file_path.read_bytes()).hexdigest()
# 计算分片数
total_parts = (file_size + self.part_size - 1) // self.part_size
progress = UploadProgress(
total_bytes=file_size,
total_parts=total_parts,
)
# MinIO Python SDK 内置分片上传
result = self.minio_client.fput_object(
bucket_name,
object_name,
str(file_path),
content_type=content_type,
part_size=self.part_size,
)
progress.uploaded_bytes = file_size
progress.parts_completed = total_parts
if progress_callback:
progress_callback(progress)
return result.etag or ""
#10. 性能优化与最佳实践
#10.1 性能调优参数
# MinIO 服务器端性能调优
# 启用异步写入(提高吞吐量,降低延迟)
export MINIO_DRIVE_SYNC=off
# 调整扫描速度(降低后台 I/O 影响)
export MINIO_SCANNER_SPEED=slow
# 并发设置
export MINIO_API_REQUESTS_MAX=1600
export MINIO_API_REQUESTS_DEADLINE=10s
# 缓存设置
export MINIO_CACHE_DRIVES="/mnt/cache1,/mnt/cache2"
export MINIO_CACHE_QUOTA=80
export MINIO_CACHE_AFTER=3
export MINIO_CACHE_WATERMARK_LOW=70
export MINIO_CACHE_WATERMARK_HIGH=90
#10.2 客户端最佳实践
# python-sdk/ontology_sdk/storage/best_practices.py
"""MinIO 客户端最佳实践"""
from minio import Minio
import urllib3
def create_optimized_client(
endpoint: str,
access_key: str,
secret_key: str,
secure: bool = True,
) -> Minio:
"""创建经过优化的 MinIO 客户端"""
# 使用连接池
http_client = urllib3.PoolManager(
num_pools=10,
maxsize=10,
retries=urllib3.Retry(
total=3,
backoff_factor=0.2,
status_forcelist=[500, 502, 503, 504],
),
timeout=urllib3.Timeout(connect=5.0, read=30.0),
)
return Minio(
endpoint=endpoint,
access_key=access_key,
secret_key=secret_key,
secure=secure,
http_client=http_client,
)
# 批量操作示例
def batch_delete_objects(
client: Minio,
bucket_name: str,
prefix: str,
) -> int:
"""批量删除对象(比逐个删除快 10 倍)"""
from minio.deleteobjects import DeleteObject
delete_list = [
DeleteObject(obj.object_name)
for obj in client.list_objects(bucket_name, prefix=prefix, recursive=True)
if obj.object_name
]
if not delete_list:
return 0
errors = list(client.remove_objects(bucket_name, delete_list))
if errors:
for err in errors:
print(f"Delete error: {err}")
return len(delete_list) - len(errors)
#10.3 性能基准测试数据
| 操作 | 文件大小 | 单节点 QPS | 4 节点集群 QPS | 延迟 P99 |
|---|---|---|---|---|
| PUT | 1 MB | 2,800 | 8,500 | 12 ms |
| PUT | 100 MB | 45 | 160 | 850 ms |
| PUT | 1 GB | 4.5 | 16 | 8.2 s |
| GET | 1 MB | 4,200 | 13,000 | 8 ms |
| GET | 100 MB | 55 | 200 | 620 ms |
| LIST (1000 obj) | - | 320 | 950 | 35 ms |
| DELETE | - | 5,000 | 15,000 | 5 ms |
#Key Takeaways
-
Bucket-per-project 策略实现天然租户隔离:每个项目独立的 bucket 配合 IAM 策略,确保数据安全边界清晰,同时支持独立的生命周期管理和配额控制。
-
Presigned URL 消除应用层瓶颈:前端直传 MinIO 避免大文件经过应用服务器,上传吞吐量提升 5-10 倍,同时保持完整的权限控制。
-
模型制品注册表统一管理 ML 生命周期:基于 MinIO 的 Manifest 模式,实现从训练到生产的完整模型版本管理,无需额外的 MLflow 等组件。
-
生命周期策略自动化存储治理:日志 30 天、临时文件 7 天、Pipeline 输出 90 天的分级过期策略,存储成本降低 40%。
-
对象引用桥接 Ontology 与二进制资产:通过 s3:// URI 将非结构化资产关联到 Ontology 实体,保持数据模型的一致性。
#Next Article
下一篇 S3-05《entity_common/entity_edge/entity_event:三表模型设计》 将深入探讨 coomia-dip 数据基座的核心存储模型——如何用三张表承载任意 Ontology 的实体、关系和事件,以及这种设计背后的权衡与查询优化策略。
Tags: #MinIO #ObjectStorage #S3Compatible #ModelArtifacts #PresignedURL #BucketPerProject #LifecycleManagement #智策平台 #coomia-dip #数据基座