304 lines
10 KiB
Python
304 lines
10 KiB
Python
"""
|
||
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()
|