最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

RxJava加Retrofit文件分段上傳實(shí)現(xiàn)詳解

 更新時間:2023年01月03日 09:17:05   作者:Chavin  
這篇文章主要為大家介紹了RxJava加Retrofit文件分段上傳實(shí)現(xiàn)詳解,有需要的朋友可以借鑒參考下,希望能夠有所幫助,祝大家多多進(jìn)步,早日升職加薪

前言

本文基于 RxJava 和 Retrofit 庫,設(shè)計(jì)并實(shí)現(xiàn)了一種用于大文件分塊上傳的工具,并對其進(jìn)行了全面的拆解分析。拋磚引玉,對同樣有處理文件分塊上傳訴求的讀者,可能會起到一定的啟發(fā)作用。

文章主體由四部分構(gòu)成:

  • 首先分析問題,問題拆解為:多線程分段讀取文件、構(gòu)建和發(fā)出文件片段上傳請求
  • 基于 JDK 隨機(jī)讀取文件的類庫,設(shè)計(jì)本地多線程分段讀取文件的單元
  • 基于 Retrofit 設(shè)計(jì)由文件片段構(gòu)建上傳的網(wǎng)絡(luò)請求
  • 從上述設(shè)計(jì)演變而來的完整代碼實(shí)現(xiàn)

  另外,在文章提供的完整代碼中,還附了一段由 PHP 編寫,用來接收多線程分段數(shù)據(jù)的服務(wù)端接口實(shí)現(xiàn),其中處理了因客戶端都線程上傳片段,導(dǎo)致服務(wù)端接收的文件片段無序,故需在適當(dāng)時機(jī)合并分塊構(gòu)成目標(biāo)文件。

受限于筆者的開發(fā)經(jīng)驗(yàn)與理論理解,文章的思路和代碼難免可能有偏頗,對于有改進(jìn)和優(yōu)化的部分,歡迎大家討論區(qū)提出。

問題拆解

要完成文件分段上傳到服務(wù)端,第一步是分段讀取本地文件。通常分段是為了多線程同時執(zhí)行上傳,提高設(shè)備計(jì)算和網(wǎng)絡(luò)資源利用率,減少上傳時間優(yōu)化體驗(yàn),這樣即需要一個支持多線程的文件分段讀取工具。由于文件可能超過設(shè)備內(nèi)存大小,在讀取這類超大文件時需要控制最大讀取量防止內(nèi)存溢出。此時文件已從磁盤數(shù)據(jù)轉(zhuǎn)換為內(nèi)存中的字節(jié)數(shù)據(jù),只需要將這些內(nèi)存數(shù)據(jù)傳給服務(wù)端即可。這樣問題被分成 3 個子問題:

  • 分段讀取文件到內(nèi)存中
  • 控制多線程數(shù)量
  • 將文件片段傳給服務(wù)端

問題 1 很好解決,利用 Java 的 RandomAccessFile 可對文件的隨機(jī)讀取的特性,即可按需讀取文件片段到內(nèi)存中。

問題 2 相對復(fù)雜一點(diǎn),但如果有閱讀過 JDK 中線程池源碼的讀者,就會發(fā)現(xiàn)這個問題的和控制線程池中線程數(shù)量其實(shí)是類似的。

問題 3 就不復(fù)雜了,Retrofit 基于 OKhttp ,OkHttp是很容易基于字節(jié)數(shù)組構(gòu)建 multipart/form-data 請求的。

分塊并發(fā)讀取文件

根據(jù)上述對問題 1、2 的拆解,可將讀取抽象為一個文件讀取器,構(gòu)建時傳入文件對象和分段大小以及最大并發(fā)數(shù),以及分段數(shù)據(jù)的回調(diào)。當(dāng)外部啟動讀取時將根據(jù)文件大小和配置的分段大小構(gòu)建若干個 Task 用于讀取對應(yīng)片段的數(shù)據(jù)。

public BlockReader(@NotNull File file, @NotNull BlockCallback callback, int poolSize, int blockSize) {
    mFile = file;
    mCallback = callback;
    mPoolSize = poolSize;
    mBlockSize = blockSize;
}
public void start(@Nullable BlockFilter filter) {
    Observable.empty().observeOn(Schedulers.computation()).doOnComplete(() -> {
        long length = mFile.length();
        for (long offset = 0; offset < length; offset += mBlockSize) {
            if (null != filter && filter.ignore(offset)) {
                continue;
            }
            mQueue.offer(new ReadTask(offset));
        }
        for (int i = 0; i < Math.min(mPoolSize, mQueue.size()); i++) {
            Observable.empty().observeOn(Schedulers.io()).doOnComplete(this::schedule).subscribe();
        }
    }).subscribe();
}

多線程調(diào)度部分,可通過加鎖和記錄狀態(tài)變量統(tǒng)計(jì)當(dāng)前正運(yùn)行的線程數(shù),則可控制字節(jié)數(shù)組數(shù),這樣就相當(dāng)于控制住了最大內(nèi)存占用。

private void schedule() {
    if (mRunning.get() >= mPoolSize) {
        return;
    }
    ReadTask task;
    synchronized (mQueue) {
        if (mRunning.get() >= mPoolSize) {
            return;
        }
        task = mQueue.poll();
        if (null != task) {
            mRunning.incrementAndGet();
        }
    }
    if (null != task) {
        task.run();
    }
}

最后是文件隨機(jī)讀取,直接調(diào)用 RandomAccessFile 的 API 即可:

private class ReadTask implements Action {
    @Override
    public void run() {
        try (RandomAccessFile raf = new RandomAccessFile(mFile, RAF_MODE);
                ByteArrayOutputStream out = new ByteArrayOutputStream(mBlockSize)) {
            raf.seek(mOffset);
            byte[] buf = new byte[DEF_BLOCK_SIZE];
            long cnt = 0;
            for (int bytes = raf.read(buf); bytes != -1 && cnt < mBlockSize; bytes = raf.read(buf)) {
                out.write(buf, 0, bytes);
                cnt += bytes;
            }
            out.flush();
            mCallback.onFinished(mOffset, out.toByteArray());
        } catch (IOException e) {
            mCallback.onFinished(mOffset, null);
        } finally {
            mRunning.decrementAndGet();
            schedule();
        }
    }
}

文件片段上傳

上傳部分則使用 Retrofit 提供的注解和 OKHttp 的類庫構(gòu)建請求。但值得一提的是需要在磁盤IO線程同步完成網(wǎng)絡(luò)IO,這樣可以避免網(wǎng)絡(luò)IO速度落后磁盤IO太多而導(dǎo)致任務(wù)堆積造成內(nèi)存溢出。

public interface BlockUploader {
    @POST("test/upload.php")
    @Multipart
    Single<Response<ResponseBody>> upload(@Header("filename") String filename,
                                          @Header("total") long total,
                                          @Header("offset") long offset,
                                          @Part List<MultipartBody.Part> body);
}
private static void syncUpload(String fileName, long fileLength, long offset, byte[] bytes) {
    RequestBody data = RequestBody.create(MediaType.parse("application/octet-stream"), bytes);
    MultipartBody body = new MultipartBody.Builder()
            .addFormDataPart("file", fileName, data)
            .setType(MultipartBody.FORM)
            .build();
    retrofit.create(BlockUploader.class).upload(fileName, fileLength, offset, body.parts()).subscribe(resp -> {
        if (resp.isSuccessful()) {
            System.out.println("? offset: " + offset + " upload succeed " + resp.code());
        } else {
            System.out.println("? offset: " + offset + " upload failed " + resp.code());
        }
    }, throwable -> {
        System.out.println("! offset: " + offset + " upload failed");
    });
}

完整代碼

為控制篇幅,完整代碼請移步 Github,服務(wù)端部分處理形如:

以上就是RxJava加Retrofit文件分段上傳示例的詳細(xì)內(nèi)容,更多關(guān)于RxJava Retrofit文件上傳的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

最新評論

宝山区| 惠来县| 贵定县| 溆浦县| 尼木县| 乌海市| 广州市| 贵港市| 墨江| 鄂尔多斯市| 曲水县| 桃江县| 平远县| 合水县| 延长县| 德钦县| 滦南县| 肃北| 石家庄市| 调兵山市| 科技| 故城县| 凌云县| 肇东市| 襄汾县| 巩留县| 大渡口区| 武冈市| 依安县| 巴马| 贵州省| 广水市| 呼玛县| 开封市| 新乐市| 鸡西市| 资源县| 长海县| 平武县| 萝北县| 嫩江县|