220 lines
8.0 KiB
Python
220 lines
8.0 KiB
Python
import os
|
||
import time
|
||
import requests
|
||
import csv
|
||
from requests.adapters import HTTPAdapter
|
||
from urllib3.util.retry import Retry
|
||
from tqdm import tqdm
|
||
import concurrent.futures
|
||
|
||
# --- 配置项 ---
|
||
csv_file_path = '/home/tplink/code/autool-dispatcher/saved_csv/download_error_2026-4-21.csv'
|
||
DOWNLOAD_FOLDER = r"/srv/samba/disk2/dpi/mumu_apk"
|
||
TEMPLATE_APKPURE_DOWNLOAD = "https://d.apkpure.net/b/XAPK/{package_name}?version=latest"
|
||
MAX_WORKERS = 10
|
||
# 最大线程数
|
||
RETRY_DELAY_SECONDS = 5 # 每次重试的等待时间
|
||
|
||
# 禁用由 verify=False 引起的 InsecureRequestWarning 警告
|
||
import urllib3
|
||
|
||
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
|
||
|
||
# --- 创建下载目录 ---
|
||
if not os.path.exists(DOWNLOAD_FOLDER):
|
||
print(f"创建下载目录: {DOWNLOAD_FOLDER}")
|
||
try:
|
||
os.makedirs(DOWNLOAD_FOLDER)
|
||
except OSError as e:
|
||
print(f"错误: 无法创建下载目录 '{DOWNLOAD_FOLDER}'。请检查权限或路径是否正确。")
|
||
exit()
|
||
|
||
|
||
def requests_session_with_retries(retries=3, backoff_factor=0.3, status_forcelist=(500, 502, 504), session=None):
|
||
session = session or requests.Session()
|
||
retry = Retry(total=retries, read=retries, connect=retries, backoff_factor=backoff_factor,
|
||
status_forcelist=status_forcelist)
|
||
adapter = HTTPAdapter(max_retries=retry)
|
||
session.mount('http://', adapter)
|
||
session.mount('https://', adapter)
|
||
return session
|
||
|
||
|
||
def get_file_size(url):
|
||
"""尝试获取文件总大小,如果HEAD请求失败,则尝试GET请求。"""
|
||
try:
|
||
response = requests.head(url, verify=False, timeout=10)
|
||
response.raise_for_status()
|
||
content_length = response.headers.get('Content-Length')
|
||
if content_length:
|
||
return int(content_length)
|
||
except requests.RequestException:
|
||
pass
|
||
|
||
try:
|
||
response = requests.get(url, stream=True, verify=False, timeout=10)
|
||
response.raise_for_status()
|
||
content_length = response.headers.get('Content-Length')
|
||
if content_length:
|
||
response.close()
|
||
return int(content_length)
|
||
except requests.RequestException as e:
|
||
tqdm.write(f"警告: 无法获取文件大小。原因: {e}")
|
||
return 0
|
||
|
||
|
||
def download_file(package_name):
|
||
"""
|
||
从指定URL下载文件,支持断点续传和进度显示。
|
||
此函数现在是线程安全的,使用tqdm.write进行输出。
|
||
"""
|
||
download_url = TEMPLATE_APKPURE_DOWNLOAD.format(package_name=package_name)
|
||
print(download_url)
|
||
package_dir = os.path.join(DOWNLOAD_FOLDER, package_name)
|
||
os.makedirs(package_dir, exist_ok=True)
|
||
final_filename = os.path.join(package_dir, f"{package_name}.xapk")
|
||
partial_filename = final_filename + ".part"
|
||
|
||
if os.path.exists(final_filename):
|
||
tqdm.write(f"文件 '{os.path.basename(final_filename)}' 已完整下载,跳过。")
|
||
return package_name, True
|
||
|
||
total_size = get_file_size(download_url)
|
||
|
||
resumable_size = 0
|
||
headers = {
|
||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/108.0.0.0 Safari/537.36',
|
||
'Referer': f'https://apkpure.net/1/{package_name}'
|
||
}
|
||
if os.path.exists(partial_filename):
|
||
resumable_size = os.path.getsize(partial_filename)
|
||
headers['Range'] = f'bytes={resumable_size}-'
|
||
|
||
try:
|
||
tqdm.write(f"\n--- 开始下载: {package_name} ---")
|
||
session = requests_session_with_retries()
|
||
with session.get(download_url, headers=headers, stream=True, timeout=60, verify=False) as response:
|
||
response.raise_for_status()
|
||
|
||
mode = 'ab' if resumable_size > 0 else 'wb'
|
||
|
||
with open(partial_filename, mode) as f, tqdm(
|
||
total=total_size, initial=resumable_size,
|
||
unit='B', unit_scale=True, unit_divisor=1024,
|
||
desc=f"{package_name}", ascii=True, leave=False, miniters=1
|
||
) as pbar:
|
||
for chunk in response.iter_content(chunk_size=1024 * 1024):
|
||
if not chunk:
|
||
break
|
||
f.write(chunk)
|
||
pbar.update(len(chunk))
|
||
|
||
if os.path.getsize(partial_filename) != total_size and total_size > 0:
|
||
raise Exception("文件下载不完整,大小不匹配。")
|
||
|
||
os.rename(partial_filename, final_filename)
|
||
tqdm.write(f"下载完成,文件已保存为: {os.path.basename(final_filename)}")
|
||
return package_name, True
|
||
|
||
except Exception as e:
|
||
tqdm.write(f"\n请求异常: {package_name} - {e}")
|
||
return package_name, False
|
||
|
||
|
||
def run_downloads_in_parallel(packages_to_download):
|
||
"""
|
||
使用多线程并发下载文件列表。
|
||
"""
|
||
if not packages_to_download:
|
||
print("没有需要下载的包。")
|
||
return []
|
||
|
||
print(f"\n--- 开始多线程下载 {len(packages_to_download)} 个包 (最大线程数: {MAX_WORKERS}) ---")
|
||
|
||
failed_packages = []
|
||
|
||
with concurrent.futures.ThreadPoolExecutor(max_workers=MAX_WORKERS) as executor:
|
||
future_to_package = {executor.submit(download_file, pkg): pkg for pkg in packages_to_download}
|
||
|
||
main_pbar = tqdm(total=len(packages_to_download), desc="总下载进度")
|
||
|
||
for future in concurrent.futures.as_completed(future_to_package):
|
||
package_name, success = future.result()
|
||
if not success:
|
||
failed_packages.append(package_name)
|
||
main_pbar.update(1)
|
||
|
||
main_pbar.close()
|
||
|
||
return failed_packages
|
||
|
||
|
||
def retry_failed_downloads(failed_packages):
|
||
"""
|
||
反复重试下载失败的包,直到所有包都下载成功。
|
||
"""
|
||
if not failed_packages:
|
||
print("\n所有包都已成功下载!")
|
||
return
|
||
|
||
while failed_packages:
|
||
print(f"\n--- 以下包下载失败,开始重试: {len(failed_packages)} 个 ---")
|
||
for pkg in failed_packages:
|
||
print(f"- {pkg}")
|
||
|
||
print(f"等待 {RETRY_DELAY_SECONDS} 秒后开始下一次重试...")
|
||
time.sleep(RETRY_DELAY_SECONDS)
|
||
|
||
failed_packages = run_downloads_in_parallel(failed_packages)
|
||
|
||
print("\n--- 所有包已成功下载!程序退出。---")
|
||
|
||
|
||
|
||
# --- 主程序入口 ---
|
||
if __name__ == "__main__":
|
||
packages_from_csv = []
|
||
if not os.path.exists(csv_file_path):
|
||
print(f"错误: 找不到文件 '{csv_file_path}'。请确保文件存在于脚本的同一目录下。")
|
||
exit()
|
||
|
||
with open(csv_file_path, 'r', encoding='utf-8') as f:
|
||
reader = csv.DictReader(f)
|
||
for row in reader:
|
||
pkg = row.get('package_name', '').strip()
|
||
if pkg:
|
||
packages_from_csv.append(pkg)
|
||
|
||
success_packages = set()
|
||
success_csv_path = 'success_tasks.csv'
|
||
if os.path.exists(success_csv_path):
|
||
with open(success_csv_path, 'r', encoding='utf-8') as f:
|
||
reader = csv.DictReader(f)
|
||
for row in reader:
|
||
pkg = row.get('包名', '').strip()
|
||
if pkg:
|
||
success_packages.add(pkg)
|
||
print(f"已成功处理的包: {len(success_packages)} 个")
|
||
|
||
packages_from_csv = [pkg for pkg in packages_from_csv if pkg not in success_packages]
|
||
|
||
existing_files = set()
|
||
for name in os.listdir(DOWNLOAD_FOLDER):
|
||
subdir = os.path.join(DOWNLOAD_FOLDER, name)
|
||
if os.path.isdir(subdir):
|
||
for f in os.listdir(subdir):
|
||
if f.endswith('.xapk'):
|
||
existing_files.add(f.replace('.xapk', ''))
|
||
print(len(packages_from_csv))
|
||
print(existing_files)
|
||
print(existing_files.intersection(set(packages_from_csv)))
|
||
packages_to_download = [pkg for pkg in packages_from_csv if pkg not in existing_files]
|
||
packages_to_download = set(packages_to_download)
|
||
if not packages_to_download:
|
||
print("所有包都已存在,无需下载。")
|
||
else:
|
||
# 首次下载尝试
|
||
failed_downloads = run_downloads_in_parallel(packages_to_download)
|
||
|
||
# 进入重试循环,直到所有包都成功
|
||
retry_failed_downloads(failed_downloads) |