oss.py
4.24 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
import requests
import oss2
import math
import time
from urllib.parse import urlparse
from requests.adapters import HTTPAdapter
from .config import (
HTTP_POOL_MAXSIZE, HTTP_TRANSFER_RETRIES, HTTP_TRANSFER_RETRY_BACKOFF_SECONDS, OSS_CONFIG,
)
from .utils import compute_audio_md5
_HTTP_SESSION = requests.Session()
_HTTP_ADAPTER = HTTPAdapter(pool_connections=HTTP_POOL_MAXSIZE, pool_maxsize=HTTP_POOL_MAXSIZE)
_HTTP_SESSION.mount('http://', _HTTP_ADAPTER)
_HTTP_SESSION.mount('https://', _HTTP_ADAPTER)
def _http_get(url: str, timeout: int = 30):
return _HTTP_SESSION.get(url, timeout=timeout)
def _clean_url(url) -> str:
if url is None:
return ''
if isinstance(url, float) and math.isnan(url):
return ''
text = str(url).strip()
if not text or text.lower() in {'nan', 'none', 'null'}:
return ''
parsed = urlparse(text)
if parsed.scheme not in {'http', 'https'} or not parsed.netloc:
return ''
return text
def _host(url: str | None) -> str:
return urlparse(url or '').netloc.lower()
def _is_target_oss_url(url: str | None, base_url: str) -> bool:
return bool(url) and _host(url) == _host(base_url)
def _download_url(url: str, base_url: str) -> str:
"""Optionally rewrite target-bucket downloads to an internal/CNAME base URL."""
download_base_url = OSS_CONFIG.get('download_base_url')
rewrite_from = OSS_CONFIG.get('download_rewrite_from_base_url') or base_url
if not download_base_url or not url.startswith(rewrite_from.rstrip('/') + '/'):
return url
return f"{download_base_url.rstrip('/')}{url[len(rewrite_from.rstrip('/')):]}"
def _download_content(url: str, base_url: str) -> bytes:
"""下载资源;针对超时、连接中断和服务端错误进行有限重试。"""
download_url = _download_url(url, base_url)
attempts = max(1, HTTP_TRANSFER_RETRIES)
for attempt in range(attempts):
try:
resp = _http_get(download_url, timeout=30)
resp.raise_for_status()
return resp.content
except requests.RequestException as exc:
status_code = getattr(getattr(exc, 'response', None), 'status_code', None)
# 4xx(429 除外)是确定性失败,不浪费时间重试。
if status_code is not None and 400 <= status_code < 500 and status_code != 429:
raise
if attempt == attempts - 1:
raise
time.sleep(HTTP_TRANSFER_RETRY_BACKOFF_SECONDS * (2 ** attempt))
def transfer_url(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> str:
"""
将 url 指向的文件转移到 archive-dev OSS 的 oss_key 路径。
若 url 为空或已在目标 bucket,直接返回原 url(不上传)。
返回新的公开访问 URL。
"""
url = _clean_url(url)
if not url:
return ''
if _is_target_oss_url(url, base_url):
return url
content = _download_content(url, base_url)
bucket.put_object(oss_key, content)
return f"{base_url.rstrip('/')}/{oss_key}"
def download_text_url(url: str | None, base_url: str) -> str:
"""下载歌词等文本 URL;即使源文件已在目标 OSS,也实际读取其内容。"""
url = _clean_url(url)
if not url:
return ''
content = _download_content(url, base_url)
return content.decode('utf-8-sig')
def transfer_url_with_md5(url: str | None, oss_key: str, bucket: oss2.Bucket, base_url: str) -> tuple[str, str]:
"""
将音频 URL 转存到 OSS,并基于下载到的音频字节计算 MD5。
已在目标 OSS 的 URL 无需重新下载,无法可靠计算 MD5,返回空 MD5。
"""
url = _clean_url(url)
if not url:
return '', ''
if _is_target_oss_url(url, base_url):
return url, ''
content = _download_content(url, base_url)
audio_md5 = compute_audio_md5(content)
bucket.put_object(oss_key, content)
return f"{base_url.rstrip('/')}/{oss_key}", audio_md5
def build_oss_key(platform: str, category: str, filename: str) -> str:
"""
构造 archive-dev 内的存储路径。
platform: 'qq' | 'kugou' | 'netease'
category: 'audio' | 'cover' | 'singer' | 'album'
filename: 带扩展名的文件名
"""
return f"crawler/{category}/{platform}/{filename}"