返回博客

MinIO 对象存储:大文件和模型制品管理

Tags: #MinIO #ObjectStorage #ModelArtifacts #PresignedURL #BucketPerProject #智策平台

Coomia发布于 2025年7月13日22 分钟阅读
分享本文Twitter / X

系列: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 数据基座中,数据可以分为两大类:

Code
数据分类与存储映射:

┌─────────────────────────────────────────────────────┐
│                  coomia-dip 数据层                     │
│                                                       │
│  ┌─────────────────────┐  ┌─────────────────────────┐ │
│  │   结构化数据          │  │   非结构化数据            │ │
│  │                       │  │                           │ │
│  │  • Ontology 实体      │  │  • Pipeline 输出文件      │ │
│  │  • 关系边             │  │  • ML 模型权重            │ │
│  │  • 事件流             │  │  • 导出报表 (CSV/Excel)   │ │
│  │  • 指标时序           │  │  • 用户上传附件           │ │
│  │  • 审计日志           │  │  • 可视化截图             │ │
│  │                       │  │  • 大型 JSON/Parquet      │ │
│  │  存储:Doris + Iceberg│  │  存储:MinIO (S3 兼容)    │ │
│  └─────────────────────┘  └─────────────────────────┘ │
└─────────────────────────────────────────────────────┘

#1.2 为什么不用本地文件系统或 NFS

方案问题致命缺陷
本地文件系统单节点瓶颈节点故障即数据丢失
NFS性能差,锁竞争无法水平扩展
HDFS运维复杂,资源占用大对中小规模过度设计
云 S3外部依赖,网络延迟私有化部署无法使用
MinIOS3 兼容,轻量高性能最佳选择

#1.3 MinIO 的关键优势

MinIO 是一个高性能 S3 兼容对象存储系统,在 coomia-dip 场景中具备以下优势:

  • S3 API 100% 兼容:所有 S3 SDK/工具直接可用
  • 纠删码(Erasure Coding):数据冗余无需 RAID
  • 分布式部署:多节点高可用
  • 轻量高效:单二进制文件,Go 语言实现
  • 对象锁和版本控制:合规和审计需求
  • 生命周期管理:自动过期和清理

#2. 部署架构设计

#2.1 集群拓扑

Code
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 配置

YAML
# 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 开发环境单节点部署

YAML
# 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:

Code
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 目录结构约定

Code
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
# 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 的集成架构

Code
Pipeline 输出流程:

┌──────────┐    ┌──────────┐    ┌──────────────┐    ┌─────────┐
│ Pipeline │    │ Compute  │    │   MinIO      │    │  Doris  │
│ Engine   │───>│ Worker   │───>│   Storage    │───>│ Catalog │
│(Schedule)│    │(执行节点)│    │  (大文件)    │    │ (元数据)│
└──────────┘    └────┬─────┘    └──────────────┘    └─────────┘
                     │
                     │ 同时写入
                     ▼
              ┌──────────────┐
              │ Iceberg Table│
              │  (结构化结果)│
              └──────────────┘

#4.2 Pipeline 输出写入器

Python
# 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 模型制品生命周期

Code
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 模型注册表实现

Python
# 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 工作原理

Code
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 服务实现

Python
# 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
// 前端上传示例(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 生命周期策略

Python
# 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
# 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)关联到实体:

SQL
-- 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
# 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):

Code
分片上传流程:

┌──────────┐                              ┌──────────┐
│  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
# 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 性能调优参数

Bash
# 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
# 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 性能基准测试数据

操作文件大小单节点 QPS4 节点集群 QPS延迟 P99
PUT1 MB2,8008,50012 ms
PUT100 MB45160850 ms
PUT1 GB4.5168.2 s
GET1 MB4,20013,0008 ms
GET100 MB55200620 ms
LIST (1000 obj)-32095035 ms
DELETE-5,00015,0005 ms

#Key Takeaways

  1. Bucket-per-project 策略实现天然租户隔离:每个项目独立的 bucket 配合 IAM 策略,确保数据安全边界清晰,同时支持独立的生命周期管理和配额控制。

  2. Presigned URL 消除应用层瓶颈:前端直传 MinIO 避免大文件经过应用服务器,上传吞吐量提升 5-10 倍,同时保持完整的权限控制。

  3. 模型制品注册表统一管理 ML 生命周期:基于 MinIO 的 Manifest 模式,实现从训练到生产的完整模型版本管理,无需额外的 MLflow 等组件。

  4. 生命周期策略自动化存储治理:日志 30 天、临时文件 7 天、Pipeline 输出 90 天的分级过期策略,存储成本降低 40%。

  5. 对象引用桥接 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 #数据基座