Team Ai
Apppublic

PCNUSMSE/transcript_service

sourceHugging Faceapache-2.0updated 1mo agoView on Hugging Face
0likes
oss_service.py293 linesDownload Raw Back to services
1"""OSS云存储服务模块2 3提供阿里云OSS文件上传、下载和管理功能。4"""5 6import os7import uuid8from datetime import datetime, timedelta9from pathlib import Path10from typing import List, Optional, Tuple11import asyncio12import aiohttp13import oss214from oss2.exceptions import OssError15 16from ..core.config import get_config17from ..utils.logger import get_task_logger18 19 20class OSSService:21    """OSS云存储服务"""22    23    def __init__(self):24        """初始化OSS服务"""25        self.config = get_config()26        self.oss_config = self.config.oss27        28        # 初始化OSS客户端29        auth = oss2.Auth(30            self.oss_config.access_key_id,31            self.oss_config.access_key_secret32        )33        self.bucket = oss2.Bucket(34            auth,35            self.oss_config.endpoint,36            self.oss_config.bucket_name37        )38        39        self.logger = get_task_logger(logger_name="transcript_service.oss")40    41    def _generate_object_key(self, filename: str, task_id: str) -> str:42        """生成OSS对象键名43        44        Args:45            filename: 原始文件名46            task_id: 任务ID47            48        Returns:49            OSS对象键名50        """51        now = datetime.now()52        date_path = now.strftime("%Y/%m/%d")53        timestamp = now.strftime("%Y%m%d_%H%M%S")54        55        # 获取文件扩展名56        file_ext = Path(filename).suffix57        safe_filename = f"{timestamp}_{task_id}_{uuid.uuid4().hex[:8]}{file_ext}"58        59        return f"{self.oss_config.temp_prefix}/{date_path}/{safe_filename}"60    61    async def upload_file(self, file_path: Path, task_id: str) -> Tuple[bool, str, Optional[str]]:62        """上传文件到OSS63        64        Args:65            file_path: 本地文件路径66            task_id: 任务ID67            68        Returns:69            (是否成功, 公网URL或错误信息, 对象键名)70        """71        try:72            self.logger.info(f"开始上传文件到OSS: {file_path.name}")73            74            # 生成对象键名75            object_key = self._generate_object_key(file_path.name, task_id)76            77            # 上传文件并设置公共读取权限78            try:79                # 首先上传文件80                self.bucket.put_object_from_file(object_key, str(file_path))81                82                # 设置对象ACL为公共读取83                self.bucket.put_object_acl(object_key, oss2.OBJECT_ACL_PUBLIC_READ)84                85                # 生成公网访问URL86                url = self._generate_public_url(object_key)87                self.logger.info(f"文件上传成功: {object_key}, URL: {url}")88                return True, url, object_key89                90            except oss2.exceptions.OssError as oss_err:91                # 如果设置ACL失败,尝试使用签名URL92                if 'public-read' in str(oss_err).lower():93                    self.logger.warning(f"ACL设置失败,使用签名URL: {oss_err}")94                    url = self._generate_signed_url(object_key)95                    self.logger.info(f"文件上传成功: {object_key}, URL: {url}")96                    return True, url, object_key97                else:98                    raise99                100        except OssError as e:101            error_msg = f"OSS错误: {str(e)}"102            self.logger.error(error_msg)103            return False, error_msg, None104        except Exception as e:105            error_msg = f"上传文件时发生未知错误: {str(e)}"106            self.logger.exception(error_msg)107            return False, error_msg, None108    109    async def upload_multiple_files(self, file_paths: List[Path], task_id: str) -> List[Tuple[str, bool, str, Optional[str]]]:110        """批量上传文件到OSS111        112        Args:113            file_paths: 本地文件路径列表114            task_id: 任务ID115            116        Returns:117            [(文件名, 是否成功, URL或错误信息, 对象键名), ...]118        """119        results = []120        121        # 创建异步任务122        tasks = []123        for file_path in file_paths:124            task = self._upload_single_file_async(file_path, task_id)125            tasks.append((file_path.name, task))126        127        # 等待所有上传完成128        for filename, task in tasks:129            success, url_or_error, object_key = await task130            results.append((filename, success, url_or_error, object_key))131        132        return results133    134    async def _upload_single_file_async(self, file_path: Path, task_id: str) -> Tuple[bool, str, Optional[str]]:135        """异步上传单个文件"""136        return await asyncio.get_event_loop().run_in_executor(137            None, 138            lambda: asyncio.run(self.upload_file(file_path, task_id))139        )140    141    def _generate_public_url(self, object_key: str) -> str:142        """生成公网访问URL143        144        Args:145            object_key: OSS对象键名146            147        Returns:148            公网访问URL149        """150        # 生成简单的公网访问URL(不带签名)151        # 正确的格式: https://bucket-name.endpoint/object-key152        # 注意: endpoint不能包含协议前缀153        endpoint = self.oss_config.endpoint154        if endpoint.startswith('http://'):155            endpoint = endpoint[7:]156        elif endpoint.startswith('https://'):157            endpoint = endpoint[8:]158        159        # 构造公网URL - 注意这里的格式必须正确160        url = f"https://{self.oss_config.bucket_name}.{endpoint}/{object_key}"161        162        # 记录生成的URL以便调试163        self.logger.debug(f"生成公网URL: {url}")164        165        return url166    167    def _generate_signed_url(self, object_key: str) -> str:168        """生成签名URL(备用方案)169        170        Args:171            object_key: OSS对象键名172            173        Returns:174            签名URL175        """176        # 生成有时效性的签名URL177        expire_time = int((datetime.now() + timedelta(hours=self.oss_config.url_expire_hours)).timestamp())178        url = self.bucket.sign_url('GET', object_key, expire_time)179        return url180    181    def delete_file(self, object_key: str) -> bool:182        """删除OSS文件183        184        Args:185            object_key: OSS对象键名186            187        Returns:188            是否删除成功189        """190        try:191            self.bucket.delete_object(object_key)192            self.logger.info(f"文件删除成功: {object_key}")193            return True194        except OssError as e:195            self.logger.error(f"删除文件失败: {object_key}, 错误: {str(e)}")196            return False197        except Exception as e:198            self.logger.exception(f"删除文件时发生未知错误: {object_key}, 错误: {str(e)}")199            return False200    201    def cleanup_old_files(self, days: Optional[int] = None) -> int:202        """清理过期的临时文件203        204        Args:205            days: 保留天数,默认使用配置中的值206            207        Returns:208            删除的文件数量209        """210        cleanup_days = days or self.oss_config.auto_cleanup_days211        cutoff_date = datetime.now() - timedelta(days=cleanup_days)212        213        deleted_count = 0214        prefix = self.oss_config.temp_prefix215        216        try:217            # 列出所有临时文件218            for obj in oss2.ObjectIterator(self.bucket, prefix=prefix):219                # 检查文件最后修改时间220                if obj.last_modified.replace(tzinfo=None) < cutoff_date:221                    if self.delete_file(obj.key):222                        deleted_count += 1223            224            self.logger.info(f"清理完成,删除了 {deleted_count} 个过期文件")225            return deleted_count226            227        except Exception as e:228            self.logger.exception(f"清理过期文件时发生错误: {str(e)}")229            return deleted_count230    231    def get_file_info(self, object_key: str) -> Optional[dict]:232        """获取文件信息233        234        Args:235            object_key: OSS对象键名236            237        Returns:238            文件信息字典239        """240        try:241            info = self.bucket.head_object(object_key)242            return {243                'size': info.content_length,244                'last_modified': info.last_modified,245                'etag': info.etag,246                'content_type': info.content_type247            }248        except OssError as e:249            self.logger.error(f"获取文件信息失败: {object_key}, 错误: {str(e)}")250            return None251    252    def check_bucket_exists(self) -> bool:253        """检查存储桶是否存在254        255        Returns:256            存储桶是否存在257        """258        try:259            return self.bucket.bucket_exists()260        except Exception as e:261            self.logger.error(f"检查存储桶失败: {str(e)}")262            return False263    264    def get_bucket_info(self) -> Optional[dict]:265        """获取存储桶信息266        267        Returns:268            存储桶信息269        """270        try:271            info = self.bucket.get_bucket_info()272            return {273                'name': info.name,274                'location': info.location,275                'creation_date': info.creation_date,276                'storage_class': info.storage_class277            }278        except Exception as e:279            self.logger.error(f"获取存储桶信息失败: {str(e)}")280            return None281 282 283# 全局OSS服务实例284oss_service = OSSService()285 286 287def get_oss_service() -> OSSService:288    """获取OSS服务实例289    290    Returns:291        OSS服务实例292    """293    return oss_service