-
Notifications
You must be signed in to change notification settings - Fork 28
/
file_mutipart_upload.go
389 lines (346 loc) · 11.5 KB
/
file_mutipart_upload.go
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
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
package ufsdk
import (
"bytes"
"crypto/md5"
"encoding/json"
"fmt"
"io"
"math"
"net/http"
"net/url"
"strconv"
"strings"
"sync"
)
//MultipartState 用于保存分片上传的中间状态
type MultipartState struct {
BlkSize int //服务器返回的分片大小
uploadID string
mimeType string
keyName string
etags map[int]string
mux sync.Mutex
}
//UnmarshalJSON custom unmarshal json
func (m *MultipartState) UnmarshalJSON(bytes []byte) error {
tmp := struct {
BlkSize int `json:"BlkSize"`
UploadID string `json:"UploadId"`
}{}
err := json.Unmarshal(bytes, &tmp)
if err != nil {
return err
}
m.BlkSize = tmp.BlkSize
m.uploadID = tmp.UploadID
return nil
}
type Part struct {
PartNum uint32 `json:"PartNum"`
LastModified int64 `json:"LastModified"`
Etag string `json:"Etag"`
Size uint32 `json:"Size"`
}
type ListPartsResponse struct {
Bucket string `json:"Bucket"`
Key string `json:"Key"`
StorageClass string `json:"StorageClass"`
UploadID string `json:"UploadId"`
Status int32 `json:"Status"`
NextPartNumberMarker uint32 `json:"NextPartNumberMarker"`
MaxParts uint32 `json:"MaxParts"`
IsTruncated bool `json:"IsTruncated"`
Parts []Part `json:"Parts"`
}
type uploadChan struct {
etag string
err error
}
//MPut 分片上传一个文件,filePath 是本地文件所在的路径,内部会自动对文件进行分片上传,上传的方式是同步一片一片的上传。
//mimeType 如果为空的话,会调用 net/http 里面的 DetectContentType 进行检测。
//keyName 表示传到 ufile 的文件名。
//大于 100M 的文件推荐使用本接口上传。
func (u *UFileRequest) MPut(filePath, keyName, mimeType string) error {
file, err := openFile(filePath)
if err != nil {
return err
}
defer file.Close()
if mimeType == "" {
mimeType = getMimeType(file)
}
state, err := u.InitiateMultipartUpload(keyName, mimeType)
if err != nil {
return err
}
chunk := make([]byte, state.BlkSize)
var pos int
for {
bytesRead, fileErr := file.Read(chunk)
if fileErr == io.EOF || bytesRead == 0 { //后面直接读到了结尾
break
}
buf := bytes.NewBuffer(chunk[:bytesRead])
err := u.UploadPart(buf, state, pos)
if err != nil {
u.AbortMultipartUpload(state)
return err
}
pos++
}
return u.FinishMultipartUpload(state)
}
//AsyncMPut 异步分片上传一个文件,filePath 是本地文件所在的路径,内部会自动对文件进行分片上传,上传的方式是使用异步的方式同时传多个分片的块。
//mimeType 如果为空的话,会调用 net/http 里面的 DetectContentType 进行检测。
//keyName 表示传到 ufile 的文件名。
//大于 100M 的文件推荐使用本接口上传。
//同时并发上传的分片数量为10
func (u *UFileRequest) AsyncMPut(filePath, keyName, mimeType string) error {
return u.AsyncUpload(filePath, keyName, mimeType, 10)
}
//AsyncUpload AsyncMPut 的升级版, jobs 表示同时并发的数量。
func (u *UFileRequest) AsyncUpload(filePath, keyName, mimeType string, jobs int) error {
if jobs <= 0 {
jobs = 1
}
if jobs >= 30 {
jobs = 10
}
file, err := openFile(filePath)
if err != nil {
return err
}
defer file.Close()
if mimeType == "" {
mimeType = getMimeType(file)
}
state, err := u.InitiateMultipartUpload(keyName, mimeType)
if err != nil {
return err
}
fsize := getFileSize(file)
chunkCount := divideCeil(fsize, int64(state.BlkSize)) //向上取整
concurrentChan := make(chan error, jobs)
for i := 0; i != jobs; i++ {
concurrentChan <- nil
}
wg := &sync.WaitGroup{}
for i := 0; i != chunkCount; i++ {
uploadErr := <-concurrentChan //最初允许启动 10 个 goroutine,超出10个后,有分片返回才会开新的goroutine.
if uploadErr != nil {
err = uploadErr
break // 中间如果出现错误立即停止继续上传
}
wg.Add(1)
go func(pos int) {
defer wg.Done()
offset := int64(state.BlkSize * pos)
chunk := make([]byte, state.BlkSize)
bytesRead, _ := file.ReadAt(chunk, offset)
e := u.UploadPart(bytes.NewBuffer(chunk[:bytesRead]), state, pos)
concurrentChan <- e //跑完一个 goroutine 后,发信号表示可以开启新的 goroutine。
}(i)
}
wg.Wait() //等待所有任务返回
if err == nil { //再次检查剩余上传完的分片是否有错误
loopCheck:
for {
select {
case e := <-concurrentChan:
err = e
if err != nil {
break loopCheck
}
default:
break loopCheck
}
}
}
close(concurrentChan)
if err != nil {
u.AbortMultipartUpload(state)
return err
}
return u.FinishMultipartUpload(state)
}
//AbortMultipartUpload 取消分片上传,如果掉用 UploadPart 出现错误,可以调用本函数取消分片上传。
//state 参数是 InitiateMultipartUpload 返回的
func (u *UFileRequest) AbortMultipartUpload(state *MultipartState) error {
query := &url.Values{}
query.Add("uploadId", state.uploadID)
reqURL := u.genFileURL(state.keyName) + "?" + query.Encode()
req, err := http.NewRequest("DELETE", reqURL, nil)
if err != nil {
return err
}
authorization := u.Auth.Authorization("DELETE", u.BucketName, state.keyName, req.Header)
req.Header.Add("authorization", authorization)
return u.request(req)
}
//InitiateMultipartUpload 初始化分片上传,返回一个 state 用于后续的 UploadPart, FinishMultipartUpload, AbortMultipartUpload 的接口。
//
//keyName 表示传到 ufile 的文件名。
//
//mimeType 表示文件的 mimeType, 传空会报错,你可以使用 GetFileMimeType 方法检测文件的 mimeType。如果您上传的不是文件,您可以使用 http.DetectContentType https://golang.org/src/net/http/sniff.go?s=646:688#L11进行检测。
func (u *UFileRequest) InitiateMultipartUpload(keyName, mimeType string) (*MultipartState, error) {
reqURL := u.genFileURL(keyName) + "?uploads"
req, err := http.NewRequest("POST", reqURL, nil)
if err != nil {
return nil, err
}
// if mimeType == "" {
// return nil, fmt.Errorf("Mime Type 不能为空!!!")
// }
req.Header.Add("Content-Type", mimeType)
for k, v := range u.RequestHeader {
for i := 0; i < len(v); i++ {
req.Header.Add(k, v[i])
}
}
authorization := u.Auth.Authorization("POST", u.BucketName, keyName, req.Header)
req.Header.Add("authorization", authorization)
err = u.request(req)
if err != nil {
return nil, err
}
response := new(MultipartState)
err = json.Unmarshal(u.LastResponseBody, response)
if err != nil {
return nil, err
}
response.keyName = keyName
response.etags = make(map[int]string)
response.mimeType = mimeType
return response, err
}
//UploadPart 上传一个分片,buf 就是分片数据,buf 的数据块大小必须为 state.BlkSize,否则会报错。
//pardNumber 表示第几个分片,从 0 开始。例如一个文件按 state.BlkSize 分为 5 块,那么分片分别是 0,1,2,3,4。
//state 参数是 InitiateMultipartUpload 返回的
func (u *UFileRequest) UploadPart(buf *bytes.Buffer, state *MultipartState, partNumber int) error {
query := &url.Values{}
query.Add("uploadId", state.uploadID)
query.Add("partNumber", strconv.Itoa(partNumber))
reqURL := u.genFileURL(state.keyName) + "?" + query.Encode()
req, err := http.NewRequest("PUT", reqURL, buf)
if err != nil {
return err
}
if u.verifyUploadMD5 {
md5Str := fmt.Sprintf("%x", md5.Sum(buf.Bytes()))
req.Header.Add("Content-MD5", md5Str)
}
req.Header.Add("Content-Type", state.mimeType)
authorization := u.Auth.Authorization("PUT", u.BucketName, state.keyName, req.Header)
req.Header.Add("Authorization", authorization)
req.Header.Add("Content-Length", strconv.Itoa(buf.Len()))
resp, err := u.requestWithResp(req)
if err != nil {
return err
}
if !VerifyHTTPCode(resp.StatusCode) {
err = u.responseParse(resp)
if err != nil {
return err
}
return fmt.Errorf("Remote response code is %d - %s not 2xx call DumpResponse(true) show details",
resp.StatusCode, http.StatusText(resp.StatusCode))
}
defer resp.Body.Close()
etag := strings.Trim(resp.Header.Get("Etag"), "\"") //为保证线程安全,这里就不保留 lastResponse
if etag == "" {
etag = strings.Trim(resp.Header.Get("ETag"), "\"") //为保证线程安全,这里就不保留 lastResponse
}
state.mux.Lock()
state.etags[partNumber] = etag
state.mux.Unlock()
return nil
}
//UploadPartCopy 实现从一个已存在的Object中拷贝数据来上传一个Part。
func (u *UFileRequest) UploadPartCopy(state *MultipartState, partNumber int, sourceBucketName, sourceObject string, offset, size int64) error {
query := &url.Values{}
query.Add("uploadId", state.uploadID)
query.Add("partNumber", strconv.Itoa(partNumber))
reqURL := u.genFileURL(state.keyName) + "?" + query.Encode()
req, err := http.NewRequest("PUT", reqURL, nil)
if err != nil {
return err
}
req.Header.Add("X-Ufile-Copy-Source", fmt.Sprintf("/%s/%s", sourceBucketName, sourceObject))
req.Header.Add("X-Ufile-Copy-Source-Range", fmt.Sprintf("bytes=%d-%d", offset, offset+size-1))
req.Header.Add("Content-Type", state.mimeType)
authorization := u.Auth.Authorization("PUT", u.BucketName, state.keyName, req.Header)
req.Header.Add("Authorization", authorization)
resp, err := u.requestWithResp(req)
if err != nil {
return err
}
if !VerifyHTTPCode(resp.StatusCode) {
err = u.responseParse(resp)
if err != nil {
return err
}
return fmt.Errorf("Remote response code is %d - %s not 2xx call DumpResponse(true) show details",
resp.StatusCode, http.StatusText(resp.StatusCode))
}
defer resp.Body.Close()
etag := strings.Trim(resp.Header.Get("Etag"), "\"") //为保证线程安全,这里就不保留 lastResponse
if etag == "" {
etag = strings.Trim(resp.Header.Get("ETag"), "\"") //为保证线程安全,这里就不保留 lastResponse
}
state.mux.Lock()
state.etags[partNumber] = etag
state.mux.Unlock()
return nil
}
//FinishMultipartUpload 完成分片上传。分片上传必须要调用的接口。
//state 参数是 InitiateMultipartUpload 返回的
func (u *UFileRequest) FinishMultipartUpload(state *MultipartState) error {
query := &url.Values{}
query.Add("uploadId", state.uploadID)
reqURL := u.genFileURL(state.keyName) + "?" + query.Encode()
var etagsStr string
etagLen := len(state.etags)
for i := 0; i != etagLen; i++ {
etagsStr += state.etags[i]
if i != etagLen-1 {
etagsStr += ","
}
}
req, err := http.NewRequest("POST", reqURL, strings.NewReader(etagsStr))
if err != nil {
return err
}
req.Header.Add("Content-Type", state.mimeType)
authorization := u.Auth.Authorization("POST", u.BucketName, state.keyName, req.Header)
req.Header.Add("Authorization", authorization)
req.Header.Add("Content-Length", strconv.Itoa(len(etagsStr)))
return u.request(req)
}
func divideCeil(a, b int64) int {
div := float64(a) / float64(b)
c := math.Ceil(div)
return int(c)
}
//ListParts 获取已上传成功的分片列表
func (u *UFileRequest) ListParts(uploadId string, maxParts, partNumberMarker int) (list ListPartsResponse, err error) {
if maxParts == 0 {
maxParts = 100
}
query := &url.Values{}
query.Add("uploadId", uploadId)
query.Add("max-parts", strconv.Itoa(maxParts))
query.Add("part-number-marker", strconv.Itoa(partNumberMarker))
reqURL := u.genFileURL("") + "?muploadpart&" + query.Encode()
req, err := http.NewRequest("GET", reqURL, nil)
if err != nil {
return
}
authorization := u.Auth.Authorization("GET", u.BucketName, "", req.Header)
req.Header.Add("authorization", authorization)
err = u.request(req)
if err != nil {
return
}
err = json.Unmarshal(u.LastResponseBody, &list)
return
}