""" MinIO 对象存储客户端封装 提供 bucket 管理、文件上传/下载/删除/预览等标准操作 """ import io import uuid from datetime import datetime, timedelta from pathlib import Path from typing import BinaryIO, Optional, Union from django.conf import settings from minio import Minio from minio.error import S3Error from loguru import logger class MinioStorage: """MinIO 对象存储管理器""" def __init__(self): self._client: Optional[Minio] = None @property def client(self) -> Minio: if self._client is None: self._client = Minio( endpoint=settings.MINIO_ENDPOINT, access_key=settings.MINIO_ACCESS_KEY, secret_key=settings.MINIO_SECRET_KEY, secure=settings.MINIO_SECURE, ) return self._client # ── Bucket 管理 ────────────────────────────────────── def bucket_exists(self, bucket_name: str) -> bool: """检查 bucket 是否存在""" return self.client.bucket_exists(bucket_name) def create_bucket(self, bucket_name: str) -> bool: """创建 bucket,已存在则跳过""" if not self.bucket_exists(bucket_name): self.client.make_bucket(bucket_name) logger.info(f"MinIO bucket created: {bucket_name}") return True return False def list_buckets(self) -> list[str]: """列出所有 bucket""" return [b.name for b in self.client.list_buckets()] def remove_bucket(self, bucket_name: str) -> bool: """删除空 bucket""" if self.bucket_exists(bucket_name): self.client.remove_bucket(bucket_name) logger.info(f"MinIO bucket removed: {bucket_name}") return True return False # ── 对象管理 ──────────────────────────────────────── def list_objects( self, bucket_name: str, prefix: str = "", recursive: bool = True ) -> list[dict]: """列出 bucket 中的对象 Returns: [{name, size, last_modified, etag, content_type}, ...] """ objects = self.client.list_objects(bucket_name, prefix=prefix, recursive=recursive) return [ { "name": obj.object_name, "size": obj.size, "last_modified": obj.last_modified.isoformat() if obj.last_modified else None, "etag": obj.etag, "content_type": obj.content_type, "is_dir": obj.is_dir, } for obj in objects ] def upload_file( self, bucket_name: str, object_name: str, file_path: Union[str, Path], content_type: Optional[str] = None, ) -> str: """上传本地文件到 MinIO Returns: object_name """ ensure_bucket(self, bucket_name) self.client.fput_object( bucket_name=bucket_name, object_name=object_name, file_path=str(file_path), content_type=content_type, ) logger.info(f"MinIO upload: {bucket_name}/{object_name}") return object_name def upload_bytes( self, bucket_name: str, object_name: str, data: Union[bytes, BinaryIO], length: int = -1, content_type: Optional[str] = None, ) -> str: """上传字节数据到 MinIO Returns: object_name """ ensure_bucket(self, bucket_name) if isinstance(data, bytes): stream = io.BytesIO(data) actual_length = len(data) else: stream = data actual_length = length self.client.put_object( bucket_name=bucket_name, object_name=object_name, data=stream, length=actual_length, content_type=content_type, ) logger.info(f"MinIO upload bytes: {bucket_name}/{object_name}") return object_name def upload_stream( self, bucket_name: str, object_name: str, data: BinaryIO, length: int, content_type: Optional[str] = None, ) -> str: """流式上传文件对象到 MinIO""" ensure_bucket(self, bucket_name) self.client.put_object( bucket_name=bucket_name, object_name=object_name, data=data, length=length, content_type=content_type, ) logger.info(f"MinIO upload stream: {bucket_name}/{object_name}") return object_name def download_file( self, bucket_name: str, object_name: str, local_path: Union[str, Path] ) -> str: """下载对象到本地文件 Returns: local_path """ self.client.fget_object(bucket_name, object_name, str(local_path)) logger.info(f"MinIO download: {bucket_name}/{object_name} -> {local_path}") return str(local_path) def download_bytes(self, bucket_name: str, object_name: str) -> bytes: """下载对象为 bytes""" response = self.client.get_object(bucket_name, object_name) try: return response.read() finally: response.close() response.release_conn() def get_object(self, bucket_name: str, object_name: str): """获取对象流响应,调用方负责关闭连接""" return self.client.get_object(bucket_name, object_name) def get_object_info(self, bucket_name: str, object_name: str) -> dict: """获取对象元信息""" stat = self.client.stat_object(bucket_name, object_name) return { "name": stat.object_name, "size": stat.size, "etag": stat.etag, "content_type": stat.content_type, "last_modified": stat.last_modified.isoformat() if stat.last_modified else None, } def object_exists(self, bucket_name: str, object_name: str) -> bool: """检查对象是否存在""" try: self.client.stat_object(bucket_name, object_name) return True except S3Error: return False def delete_object(self, bucket_name: str, object_name: str) -> bool: """删除单个对象""" self.client.remove_object(bucket_name, object_name) logger.info(f"MinIO delete: {bucket_name}/{object_name}") return True def delete_objects(self, bucket_name: str, object_names: list[str]) -> list[str]: """批量删除对象""" errors = [] for name in object_names: try: self.client.remove_object(bucket_name, name) except S3Error as e: errors.append(name) logger.error(f"MinIO delete error {bucket_name}/{name}: {e}") logger.info(f"MinIO batch delete: {bucket_name}, {len(object_names) - len(errors)}/{len(object_names)}") return errors # ── 预签名 URL ────────────────────────────────────── def presign_get( self, bucket_name: str, object_name: str, expires: int = 3600, response_headers: Optional[dict] = None, ) -> str: """生成预签名下载 URL Args: bucket_name: bucket 名称 object_name: 对象路径 expires: 过期时间(秒),默认 1 小时 """ presign_endpoint = getattr(settings, 'MINIO_PRESIGN_ENDPOINT', settings.MINIO_ENDPOINT) if presign_endpoint != settings.MINIO_ENDPOINT: presign_client = Minio( endpoint=presign_endpoint, access_key=settings.MINIO_ACCESS_KEY, secret_key=settings.MINIO_SECRET_KEY, secure=settings.MINIO_SECURE, ) return presign_client.presigned_get_object( bucket_name, object_name, expires=timedelta(seconds=expires), response_headers=response_headers, ) return self.client.presigned_get_object( bucket_name, object_name, expires=timedelta(seconds=expires), response_headers=response_headers, ) def presign_upload( self, bucket_name: str, object_name: str, expires: int = 3600 ) -> str: """生成预签名上传 URL""" return self.client.presigned_put_object( bucket_name, object_name, expires=timedelta(seconds=expires) ) # ── 辅助方法 ──────────────────────────────────────── def generate_object_name(self, filename: str, prefix: str = "") -> str: """生成唯一对象名,避免覆盖 Format: {prefix}/YYYY-MM-DD/{uuid}_{original_filename} """ date_str = datetime.now().strftime("%Y-%m-%d") uid = uuid.uuid4().hex[:12] safe_name = Path(filename).name parts = [p for p in [prefix, date_str, f"{uid}_{safe_name}"] if p] return "/".join(parts) def copy_object( self, src_bucket: str, src_object: str, dst_bucket: str, dst_object: str ) -> str: """复制对象(同 bucket 内或跨 bucket)""" from minio.commonconfig import CopySource ensure_bucket(self, dst_bucket) self.client.copy_object( bucket_name=dst_bucket, object_name=dst_object, source=CopySource(src_bucket, src_object), ) logger.info(f"MinIO copy: {src_bucket}/{src_object} -> {dst_bucket}/{dst_object}") return dst_object def ensure_bucket(storage: MinioStorage, bucket_name: str): """确保 bucket 存在""" if not storage.bucket_exists(bucket_name): storage.create_bucket(bucket_name) # 全局单例 minio_storage = MinioStorage()