MultiUploader.java
16.3 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
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
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
package com.cnlive.demo.upload;
import android.content.Context;
import android.os.Handler;
import android.os.Message;
import android.util.Log;
import com.cnlive.lib.upload.listener.CNAbortMultipartUploadResponseListener;
import com.cnlive.lib.upload.listener.CNCompleteMultipartUploadResponseListener;
import com.cnlive.lib.upload.listener.CNInitiateMultipartUploadResponseListener;
import com.cnlive.lib.upload.listener.CNListPartsResponseListener;
import com.cnlive.lib.upload.listener.CNUploadPartResponseListener;
import com.cnlive.lib.upload.model.CNPart;
import com.cnlive.lib.upload.model.CNPartETag;
import com.cnlive.lib.upload.model.UploadError;
import com.cnlive.lib.upload.model.acl.CNCannedAccessControlList;
import com.cnlive.lib.upload.model.result.CNCompleteMultipartUploadResult;
import com.cnlive.lib.upload.model.result.CNInitiateMultipartUploadResult;
import com.cnlive.lib.upload.model.result.CNListPartsResult;
import com.cnlive.lib.upload.services.CNAbortMultipartUploadRequest;
import com.cnlive.lib.upload.services.CNCompleteMultipartUploadRequest;
import com.cnlive.lib.upload.services.CNInitiateMultipartUploadRequest;
import com.cnlive.lib.upload.services.CNUploadPartRequest;
import com.cnlive.lib.upload.upload.CNUpload;
import org.apache.http.Header;
import java.io.File;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
import java.util.Iterator;
import java.util.List;
import java.util.concurrent.atomic.AtomicInteger;
/**
* Created by Lynn on 2017/8/7.
* 分块上传的更高层的封装,可以修改partSize等变量,log输出格式等。也可以在initiateMultipartUpload这一步之后加入保存uploadid的逻辑(最好存在你自己后台)。
* 当然也可参考MultiUploader.java自行实现你的接口。
*/
public class MultiUploader {
private static final String TAG = "MultiUploader";
private String bucketName;
private String key;
private File file;
private int concurrentNo = 2;
private long partSize = 5 * 1024 * 1024;
private String uploadId;
private AtomicInteger cur;
private CNUpload cnUpload;
private static final int INIT_DONE = 0;
private static final int PARATS_DONE = 1;
private static final int COMPLETE_DONE = 2;
private static final int GET_UPLOADED_DONE = 3;
private static final int LIMIT_DONE = 4;
private boolean pause = false;
private boolean isAbortMultipartUpload = false;
List<CNPartETag> doneParts = Collections.synchronizedList(new ArrayList<CNPartETag>());
private List<Integer> leftParts = Collections.synchronizedList(new ArrayList<Integer>());
private CNCompleteMultipartUploadResponseListener multiListener;
private final MyHandler mHandler = new MyHandler();
private void create(CNUpload cnUpload, File file, String uploadId, long partSize) {
this.cnUpload = cnUpload;
this.file = file;
this.uploadId = uploadId;
this.partSize = partSize;
}
public MultiUploader(CNUpload cnUpload, File file, String uploadId, long partSize) {
this.create(cnUpload, file, uploadId, partSize);
}
public MultiUploader(CNUpload cnUpload, File file, String uploadId) {
this.create(cnUpload, file, uploadId, this.partSize);
}
public MultiUploader(CNUpload cnUpload, File file) {
this.create(cnUpload, file, null, this.partSize);
}
public MultiUploader(CNUpload cnUpload, File file, long partSize) {
this.create(cnUpload, file, null, partSize);
}
/**
* 并发传的块数,默认就是2。一般不用改。根据你自测的性能状况来调整。
*
* @param no
*/
public void setConcurrentNo(int no) {
this.concurrentNo = no;
}
public void setListener(CNCompleteMultipartUploadResponseListener listener) {
this.multiListener = listener;
}
public String getKey() {
return this.key;
}
public String getUploadId() {
return this.uploadId;
}
/**
* 初始化分块上传
*
* @param name 文件名
* @param des 描述
* @param userId 用户id
* @param filePath 文件路径
* @param file 图片
* @param isHorizontalScreen 横竖屏
* @param callback 回到地址
* @return
*/
public boolean upload(String name, String des, String userId, String filePath, File file, int isHorizontalScreen, String callback) {
// pause = false;
if (this.uploadId != null) {
return false;
} else {
isAbortMultipartUpload = false;
final CNInitiateMultipartUploadRequest request = new CNInitiateMultipartUploadRequest(bucketName, key);
request.setCannedACL(CNCannedAccessControlList.PublicRead);
request.setContentType("video/mp4");
cnUpload.initiateMultipartUpload(request, userId, filePath, name, des, file, isHorizontalScreen, callback, new CNInitiateMultipartUploadResponseListener() {
@Override
public void onFailure(int i, UploadError uploadError, Header[] headers, String s, Throwable throwable) {
Log.w(TAG, "init multiupload fail, statusCode = " + i, throwable);
}
@Override
public void onSuccess(int i, Header[] headers, CNInitiateMultipartUploadResult cnInitiateMultipartUploadResult) {
uploadId = cnInitiateMultipartUploadResult.getUploadId();
Log.d(TAG, "init multiupload success, uploadId = " + uploadId + ", key = " + key + " : " + cnInitiateMultipartUploadResult.getBucket() + " : " + cnInitiateMultipartUploadResult.getKey());
bucketName = cnInitiateMultipartUploadResult.getBucket();
key = cnInitiateMultipartUploadResult.getKey();
mHandler.sendEmptyMessage(INIT_DONE);
}
});
return true;
}
}
private List<CNPartETag> convertPart(List<CNPart> list) {
ArrayList res = new ArrayList();
Iterator iterator = list.iterator();
while (iterator.hasNext()) {
CNPart part = (CNPart) iterator.next();
res.add(new CNPartETag(part.getPartNumber(), part.getETag()));
}
return res;
}
/**
* 罗列出已经上传的块
*/
public void getUploadedParts() {
final ArrayList res = new ArrayList();
CNListPartsResponseListener listener = new CNListPartsResponseListener() {
@Override
public void onFailure(int i, UploadError uploadError, Header[] headers, String s, Throwable throwable) {
Log.w(TAG, "list parts, statusCode = " + i, throwable);
}
@Override
public void onSuccess(int i, Header[] headers, CNListPartsResult cnListPartsResult) {
res.addAll(convertPart(cnListPartsResult.getCNParts()));
Log.e(TAG, "getUploadedParts = " + res.size());
if (pause) {
doneParts = res;
pause = false;
}
if (!cnListPartsResult.isTruncated()) {
mHandler.sendMessage(Message.obtain(mHandler, GET_UPLOADED_DONE, res));
} else {
Log.e(TAG, "File size too large. You may not use phone to upload");
}
}
};
cnUpload.listParts(bucketName, key, uploadId, listener);
}
/**
* 获取到还没有上传的块
*
* @param uploadedParts 已经上传的块
* @return
*/
public List<Integer> getLeftParts(List<CNPartETag> uploadedParts) {
ArrayList res = new ArrayList();
long start = 0L;
int partNumber = 1;
HashSet set = new HashSet();
Iterator iterator = uploadedParts.iterator();
while (iterator.hasNext()) {
CNPartETag cnPartETag = (CNPartETag) iterator.next();
set.add(Integer.valueOf(cnPartETag.getPartNumber()));
}
while (start < file.length()) {
if (!set.contains(Integer.valueOf(partNumber))) {
res.add(Integer.valueOf(partNumber));
}
++partNumber;
start += partSize;
}
return res;
}
/**
* 继续上传
*/
public void reUpload() {
if (this.uploadId == null) {
Log.i(TAG, "no upload id, cannot reupload");
} else {
getUploadedParts();
}
}
private CNPartETag cnPartETag = new CNPartETag();
/**
* 开始分块上传
*
* @param startNo
* @param N
*/
public void doWithLimit(int startNo, final int N) {
for (int i = startNo; i < startNo + concurrentNo && i < leftParts.size(); ++i) {
// Log.e(TAG, isPause() + " ");
// if (isPause()) {
// break;
// } else {
final int partNumber = ((Integer) this.leftParts.get(i)).intValue();
long offset = partNumber * this.partSize - partSize;
CNUploadPartRequest request = new CNUploadPartRequest(bucketName, key, uploadId, file, offset, partNumber, Math.min(file.length() - offset, partSize));
request.setCannedACL(CNCannedAccessControlList.PublicReadWrite);
request.setContentType("video/mp4");
cnUpload.uploadPart(request, new MyUploadPartResponseHandler(key, partNumber, uploadId) {
@Override
public void onSuccess(int statusCode, Header[] responseHeaders, CNPartETag result, String key, int partNo, String uploadId) {
cnPartETag = result;
result.setPartNumber(partNo);
if (!isAbortMultipartUpload)
doneParts.add(result);
cur.incrementAndGet();
Log.e(TAG, "doneParts = " + doneParts.size() + " ,N = " + N);
if (doneParts.size() >= N) {
mHandler.sendEmptyMessage(PARATS_DONE);
} else if (cur.get() % concurrentNo == 0) {
mHandler.sendMessage(Message.obtain(mHandler, LIMIT_DONE, cur.get(), N));
}
}
@Override
public void onFailure(int statusCode, UploadError error, Header[] responseHeaders, String response, Throwable throwable, String key, int partNo, String uploadId) {
Log.w(TAG, "upload part fail, uploadId = " + uploadId + " ,key = " + key + " ,partNo = " + partNo, throwable);
if (multiListener != null) {
multiListener.onFailure(statusCode, error, responseHeaders, response, throwable);
}
}
@Override
public void onTaskProcess(double progress, String key, int partNo, String uploadId) {
if (progress > 99) {
Log.d(TAG, "progress : " + progress + " ,key = " + key + " ,partNo = " + partNo);
}
}
});
// }
}
}
/**
* 取消分块上传
*/
public void cancelUpload() {
CNAbortMultipartUploadRequest request = new CNAbortMultipartUploadRequest(bucketName, key, uploadId);
cnUpload.abortMultipartUpload(request, new CNAbortMultipartUploadResponseListener() {
@Override
public void onFailure(int i, UploadError uploadError, Header[] headers, String s, Throwable throwable) {
Log.e(TAG, "abortMultipart fail , stateCode = " + i, throwable);
}
@Override
public void onSuccess(int i, Header[] headers) {
Log.e(TAG, "abortMultipart success");
isAbortMultipartUpload = true;
if (uploadId != null) {
uploadId = null;
}
if (doneParts != null) doneParts.clear();
}
});
}
public void reUpload(List<CNPartETag> uploadedParts) {
leftParts = getLeftParts(uploadedParts);
Log.e(TAG, "reupload uploadedParts = " + uploadedParts.size() + " , leftParts = " + leftParts.size());
int N = leftParts.size() + uploadedParts.size();
cur = new AtomicInteger(0);
// doneParts.addAll(uploadedParts);
if (leftParts.isEmpty()) {
mHandler.sendEmptyMessage(PARATS_DONE);
}
Log.e(TAG, "reupload uploadedsize = " + doneParts.size() + " ,N = " + N);
doWithLimit(0, N);
}
private void uploadParts() {
reUpload(new ArrayList<CNPartETag>());
}
/**
* 组装上传的块
*/
private void completeUpload() {
if (multiListener == null) {
multiListener = new CNCompleteMultipartUploadResponseListener() {
@Override
public void onFailure(int i, UploadError uploadError, Header[] headers, String s, Throwable throwable) {
Log.w(TAG, "complete upload fail, statusCode = " + i, throwable);
}
@Override
public void onSuccess(int i, Header[] headers, CNCompleteMultipartUploadResult completeMultipartUploadResult, String videoId) {
Log.i(TAG, "complete upload, key = " + key + ", bucketName = " + bucketName + " , uploadId = " + uploadId + "\n" + completeMultipartUploadResult.toString());
}
};
}
CNCompleteMultipartUploadRequest request = new CNCompleteMultipartUploadRequest(bucketName, key);
request.setContentType("video/mp4");
request.setCNPartETags(doneParts);
request.setUploadId(uploadId);
cnUpload.completeMultipartUpload(request, multiListener);
}
class MyHandler extends Handler {
MyHandler() {
}
@Override
public void handleMessage(Message msg) {
switch (msg.what) {
case INIT_DONE:
uploadParts();
break;
case PARATS_DONE:
completeUpload();
break;
case COMPLETE_DONE:
default:
break;
case GET_UPLOADED_DONE:
List res = (List) msg.obj;
reUpload(res);
break;
case LIMIT_DONE:
int no = msg.arg1;
int N = msg.arg2;
doWithLimit(no, N);
break;
}
}
}
abstract class MyUploadPartResponseHandler implements CNUploadPartResponseListener {
private String key;
private int partNo;
private String uploadId;
public abstract void onSuccess(int statusCode, Header[] responseHeaders, CNPartETag result, String key, int partNo, String uploadId);
public abstract void onFailure(int statusCode, UploadError error, Header[] responseHeaders, String response, Throwable throwable, String key, int partNo, String uploadId);
public abstract void onTaskProcess(double progress, String key, int partNo, String uploadId);
public MyUploadPartResponseHandler(String key, int partNo, String uploadId) {
this.key = key;
this.partNo = partNo;
this.uploadId = uploadId;
}
@Override
public void onTaskProgress(double v) {
this.onTaskProcess(v, this.key, this.partNo, this.uploadId);
}
@Override
public void onSuccess(int i, Header[] headers, CNPartETag cnPartETag) {
this.onSuccess(i, headers, cnPartETag, this.key, this.partNo, this.uploadId);
}
@Override
public void onFailure(int i, UploadError uploadError, Header[] headers, String s, Throwable throwable) {
this.onFailure(i, uploadError, headers, s, throwable, this.key, this.partNo, this.uploadId);
}
}
// public boolean isPause() {
// return pause;
// }
// public void setPause(boolean pause) {
// this.pause = pause;
// }
/**
* 此暂停方法,会立即暂停,继续上传时,当前块重新上传
* 也可以通过不上传下一块的方法暂停上传
*
* @param context
*/
public void pause(Context context) {
cnUpload.pause(context);
pause = true;
}
}