Commit 2f04d179 authored by mayinghao's avatar mayinghao

1工厂上传,2 支持多线程多个文件同时上传

parent 37ad29d0
Pipeline #24 canceled with stages
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
# -*- coding: utf-8 -*-
import time
import requests
from PyQt5.QtCore import QThread, pyqtSignal, QWaitCondition, QMutex, QObject, pyqtSlot
import os.path
from PyQt5.QtWidgets import QApplication,QMainWindow
import sys
import oss2
import json
from oss2 import SizedFileAdapter, determine_part_size
import time
import alibabacloud_oss_v2 as oss
import uuid
from oss2.models import PartInfo
# 切片上传大小
part_size = 5 * 1024 * 1024
name_prefix = "https://qgongye.oss-cn-shanghai.aliyuncs.com\\upload/order"
callback_url = "https://uat.emku.laidecloud.com/console/productOrder/filePthSave.ajax"
# 未完成任务列表
auth = oss2.Auth("LTAI3qeEHR5SP2ph", "zsV3FnjVgBDhaAm23J2vdAmMuX0YFY")
# 填写Bucket所在地域对应的Endpoint。以华东1(上海)为例,Endpoint填写为https://oss-cn-hangzhou.aliyuncs.com。
endpoint = "https://oss-cn-shanghai.aliyuncs.com"
# 填写Endpoint对应的Region信息,例如cn-hangzhou。注意,v4签名下,必须填写该参数
region = "cn-shanghai"
# yourBucketName填写存储空间名称。
bucket = oss2.Bucket(auth, endpoint, "qgongye", region=region)
# 未完成任务类
class Resume_Thread(QObject):
# 使用自定义信号和UI主线程通讯,参数是发送信号时附带参数的数据类型,可以是str、int、list等
finishSignal = pyqtSignal(str)
updateSignal = pyqtSignal(tuple)
def __init__(self, thread_num=0, job_list=()):
QObject.__init__(self)
# 要发送的文件地址
self._filename = ''
self.cond = QWaitCondition()
self.mutex = QMutex()
self._isPause = True
self.thread_num = thread_num
self.job_list = job_list
self.send_data = 0
self.leftBytes = 0
self.total_size = 0
self.oss_prefix = "upload/order/"
self.type_info = job_list[3]
self.order_id = job_list[2]
self.uuid_str = job_list[4]
self.total_size = job_list[5]
def pause(self):
self._isPause = True
def resume(self):
self._isPause = False
self.cond.wakeAll()
@property
def filename(self):
return self._filename
@filename.setter
def filename(self, value):
self._filename = value
# 进度条回调函数
def upload_callback(self, bytes_uploaded, total_bytes):
if bytes_uploaded >= total_bytes and (total_bytes != self.leftBytes):
self.send_data += 1
if bytes_uploaded == part_size:
rate = int(100 * (float(self.send_data * part_size) / float(self.total_size)))
else:
rate = int(100 * (float(bytes_uploaded + self.send_data * part_size) / float(self.total_size)))
#print(f"{rate}% {bytes_uploaded + self.send_data * part_size} {self.send_data} {bytes_uploaded}")
self.updateSignal.emit((self.thread_num, rate))
def resume_thread(self):
print(self.thread_num)
self._isPause = not self._isPause
if not self._isPause:
print("thread resume")
self.cond.wakeAll()
else:
print("thread stop")
@pyqtSlot()
def run(self):
# 下拉列表当前选项
part_finish_list = {}
local_path = self.job_list[0]
print(local_path)
currentItem = self.job_list[1]
oss_key = ""
# 如果有未完成的下载
if currentItem != "":
parts = []
print(currentItem)
key = local_path.split("/")[-1]
extension = key.split(".")[-1]
oss_key = self.oss_prefix + self.uuid_str + "." + extension
# 收集阿里云已经上传完毕的数据包
for part_info in oss2.PartIterator(bucket, oss_key, currentItem):
part_finish_list[part_info.part_number] = part_info.etag
print('part_number:', part_info.part_number)
parts.append(PartInfo(part_info.part_number, part_info.etag))
self.send_data = len(part_finish_list)-1
if self.send_data<0:
self.send_data = 0
rate = int(100 * (float(self.send_data * part_size) / float(self.total_size)))
self.updateSignal.emit((self.thread_num, rate))
# 计算总循环次数和 最后一个包的数据
total_size = os.path.getsize(local_path)
print(total_size)
circle = total_size // part_size
# 最后一包数据
left_bytes = total_size % part_size
if left_bytes != 0:
circle += 1
with open(local_path, "rb") as f:
for j in range(1, circle + 1):
if j not in list(part_finish_list.keys()):
self.mutex.lock()
if self._isPause:
self.cond.wait(self.mutex)
# 定位到到指定文件位置
f.seek((j - 1) * part_size)
read_size = min(part_size, total_size - part_size * (j - 1))
content = f.read(read_size)
result = bucket.upload_part(oss_key, currentItem, j,
content, progress_callback=self.upload_callback)
parts.append(PartInfo(j, result.etag))
self.mutex.unlock()
print(j, end=",")
headers = dict()
# 设置文件访问权限ACL。此处设置为OBJECT_ACL_PRIVATE,表示私有权限。
headers["x-oss-object-acl"] = oss2.OBJECT_ACL_PRIVATE
bucket.complete_multipart_upload(oss_key, currentItem, parts, headers=headers)
key = local_path.split('/')[-1]
tmp_key_json = "./upload_tag/tmp_" + key + ".json"
print(f"文件存在:{os.path.exists(tmp_key_json)} {tmp_key_json}")
if os.path.exists(tmp_key_json):
os.remove(tmp_key_json)
datas = {"id": str(self.order_id), "name": self.uuid_str + "." + extension, "path": name_prefix,
"info": key,
"type": self.type_info}
response = requests.post(callback_url, data=datas)
res = response.json()
print(res["obj"])
# -*- coding: utf-8 -*-
import time
import requests
from PyQt5.QtCore import QThread, pyqtSignal, QWaitCondition, QMutex, QObject, pyqtSlot
import os.path
from PyQt5.QtWidgets import QApplication,QMainWindow
import sys
import oss2
import json
from oss2 import SizedFileAdapter, determine_part_size
import uuid
import alibabacloud_oss_v2 as oss
from oss2.credentials import EnvironmentVariableCredentialsProvider
from oss2.models import PartInfo
# 切片上传大小
part_size = 5 * 1024 * 1024
name_prefix = "https://qgongye.oss-cn-shanghai.aliyuncs.com\\upload/order"
callback_url = "https://uat.emku.laidecloud.com/console/productOrder/filePthSave.ajax"
# 未完成任务列表
job_list = {}
auth = oss2.Auth("LTAI3qeEHR5SP2ph", "zsV3FnjVgBDhaAm23J2vdAmMuX0YFY")
# 填写Bucket所在地域对应的Endpoint。以华东1(上海)为例,Endpoint填写为https://oss-cn-hangzhou.aliyuncs.com。
endpoint = "https://oss-cn-shanghai.aliyuncs.com"
# 填写Endpoint对应的Region信息,例如cn-hangzhou。注意,v4签名下,必须填写该参数
region = "cn-shanghai"
# yourBucketName填写存储空间名称。
bucket = oss2.Bucket(auth, endpoint, "qgongye", region=region)
# self.thread = QThread()
# self.upload = Upload_Thread()
# self.upload.filename = filename
# self.upload.moveToThread(self.thread)
# self.upload.updateSignal.connect(self.changeProcessbar)
# self.thread.started.connect(self.upload.run)
# self.thread.start()
# 定义一个线程类
class Upload_Thread(QObject):
#自定义信号声明
# 使用自定义信号和UI主线程通讯,参数是发送信号时附带参数的数据类型,可以是str、int、list等
finishSignal = pyqtSignal(str)
updateSignal = pyqtSignal(tuple)
def __init__(self, thread_num=0, order_id=None, type_info=None):
QObject.__init__(self)
# 要发送的文件地址
self._filename = ''
self.cond = QWaitCondition()
self.mutex = QMutex()
self._isPause = False
self.thread_num = thread_num
self.send_data = 0
self.leftBytes = 0
self.total_size = 0
self.oss_prefix = "upload/order/"
self.uuid_str = str(uuid.uuid4())
self.type_info = type_info
self.order_id = order_id
def pause(self):
self._isPause = True
def resume(self):
self._isPause = False
self.cond.wakeAll()
@property
def filename(self):
return self._filename
@filename.setter
def filename(self, value):
self._filename = value
def resume_thread(self):
print(self.thread_num)
self._isPause = not self._isPause
if not self._isPause:
print("thread resume")
self.cond.wakeAll()
else:
print("thread stop")
def upload_callback(self, bytes_uploaded, total_bytes):
if bytes_uploaded >= total_bytes and (total_bytes != self.leftBytes):
self.send_data += 1
if bytes_uploaded == part_size:
rate = int(100 * (float(self.send_data * part_size) / float(self.total_size)))
else:
rate = int(100 * (float(bytes_uploaded + self.send_data * part_size) / float(self.total_size)))
#print(f"{rate}% {bytes_uploaded + self.send_data * part_size} {self.send_data} {bytes_uploaded}")
self.updateSignal.emit((self.thread_num, rate))
@pyqtSlot()
def run(self):
print("start upload thread")
key = self.filename.split('/')[-1]
extension = key.split(".")[-1]
oss_key = self.oss_prefix + self.uuid_str+"."+extension
# exist = False
# 返回值为true表示文件存在,false表示文件不存在。
if not True:
print("server has this file")
#self.ui.upload_status.setText(f"{self.filename} 已经存在于阿里云OSS")
else:
# 判断是否有碎片
print('object not exist')
#
# if filename[-4] == '.':
# json_name = filename[:-4]
# 如果没有上传完毕就把{filename:upload_id}的json文件放在upload_tag里
# if os.path.exists("./upload_tag/"+json_name+".json"):
# #找碎片
# pass
# else:
# pass
parts = []
total_size = os.path.getsize(self.filename)
self.total_size = total_size
self.leftBytes = total_size % part_size
self.send_data = 0
upload_id = bucket.init_multipart_upload(oss_key).upload_id
print(upload_id)
with open("./upload_tag/tmp_" + key + ".json", "w+") as f:
json.dump({upload_id: [self.filename, self.order_id, self.type_info, self.uuid_str, self.total_size]}, f)
if os.path.exists(self.filename):
total_size = os.path.getsize(self.filename)
print(total_size)
circle = total_size // part_size
# 最后一包数据
left_bytes = total_size % part_size
if left_bytes != 0:
circle += 1
# 逐个上传分片。
with open(self.filename, 'rb') as fileobj:
part_number = 1
offset = 0
# self.ui.progressBar1.setVisible(True)
while offset < total_size:
self.mutex.lock()
if self._isPause:
self.cond.wait(self.mutex)
num_to_upload = min(part_size, total_size - offset)
# self.ui.upload_status.setText(f'正在上传第{part_number}/{circle}个分片包')
# 调用SizedFileAdapter(fileobj, size)方法会生成一个新的文件对象,重新计算起始追加位置。
result = bucket.upload_part(oss_key, upload_id, part_number,
SizedFileAdapter(fileobj, num_to_upload), progress_callback=self.upload_callback
)
parts.append(PartInfo(part_number, result.etag))
offset += num_to_upload
part_number += 1
# if part_number == 2:
# self._isPause = True
self.mutex.unlock()
# 完成分片上传。
# 如需在完成分片上传时设置相关Headers,请参考如下示例代码。
headers = dict()
# 设置文件访问权限ACL。此处设置为OBJECT_ACL_PRIVATE,表示私有权限。
# headers["x-oss-object-acl"] = oss2.OBJECT_ACL_PRIVATE
bucket.complete_multipart_upload(oss_key, upload_id, parts, headers=headers)
if os.path.exists("./upload_tag/tmp_" + key + ".json"):
os.remove("./upload_tag/tmp_" + key + ".json")
datas = {"id": str(self.order_id), "name": self.uuid_str + "." + extension, "path": name_prefix, "info": key,
"type": self.type_info}
response = requests.post(callback_url, data=datas)
res = response.json()
print(res["obj"])
<RCC>
<qresource prefix="icons">
<file>icons/mini.svg</file>
<file>icons/quit.svg</file>
<file>icons/dog.svg</file>
<file>icons/baidu.svg</file>
<file>icons/bili.svg</file>
<file>icons/close.svg</file>
<file>icons/login.svg</file>
<file>icons/register.svg</file>
</qresource>
</RCC>
This diff is collapsed.
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment