【问题标题】:Inconsistent output from multithreaded FTP InputStreams来自多线程 FTP InputStreams 的输出不一致
【发布时间】:2015-08-21 09:42:10
【问题描述】:

我正在尝试创建一个 java 程序,将某些资产文件从 FTP 服务器下载到本地文件。因为我的(免费)FTP 服务器不支持超过几兆字节的文件大小,所以我决定在上传文件时拆分文件,并在程序下载文件时重新组合它们。这行得通,但速度很慢,因为对于每个文件,它必须获取InputStream,这需要一些时间。

我使用的 FTP 服务器可以下载文件而无需实际登录服务器,因此我使用此代码获取 InputStream

private static final InputStream getInputStream(String file) throws IOException {
    return new URL("http://site.website.com/path/" + file).openStream();
}

要获取资产文件一部分的InputStream,我正在使用以下代码:

public static InputStream getAssetInputStream(String asset, int num) throws IOException, FTPException {
    try {
        return getInputStream("assets/" + asset + "_" + num + ".raf");
    } catch (Exception e) {
        // error handling
    }
}

因为getAssetInputStreams(String, int) 方法需要一些时间来运行(特别是如果文件大小超过一兆字节),我决定将实际下载文件的代码设为多线程。这就是我的问题所在。

final Map<Integer, Boolean> done = new HashMap<Integer, Boolean>();
final Map<Integer, byte[]> parts = new HashMap<Integer, byte[]>();

for (int i = 0; i < numParts; i++) {
    final int part = i;
    done.put(part, false);

    new Thread(new Runnable() {
        @Override
        public void run() {
            try {
                InputStream is = FTP.getAssetInputStream(asset, part);
                ByteArrayOutputStream baos = new ByteArrayOutputStream();

                byte[] buf = new byte[DOWNLOAD_BUFFER_SIZE];
                int len = 0;

                while ((len = is.read(buf)) > 0) {
                    baos.write(buf, 0, len);
                    curDownload.addAndGet(len);
                    totAssets.addAndGet(len);
                }

                parts.put(part, baos.toByteArray());
                done.put(part, true);
            } catch (IOException e) {
                // error handling
            } catch (FTPException e) {
                // error handling
            }
        }
    }, "Download-" + asset + "-" + i).start();
}

while (done.values().contains(false)) {
    try {
        Thread.sleep(100);
    } catch(InterruptedException e) {
        e.printStackTrace();
    }
}

File assetFile = new File(dir, "assets/" + asset + ".raf");
assetFile.createNewFile();
FileOutputStream fos = new FileOutputStream(assetFile);

for (int i = 0; i < numParts; i++) {
    fos.write(parts.get(i));
}

fos.close();

此代码有效,但并非总是如此。当我在台式计算机上运行它时,它几乎总是有效。不是 100% 的时间,但通常它工作得很好。在我的笔记本电脑上,它的互联网连接要差得多,它几乎从不工作。结果是一个不完整的文件。有时,它会下载文件的 50%。有时,它会下载 90% 的文件,但每次都不同。

现在,如果我将 .start() 替换为 .run(),则代码 100% 可以正常工作,即使在我的笔记本电脑上也是如此。然而,它非常慢,所以我宁愿不使用.run()

有没有办法可以更改我的代码以使其在多线程中工作?任何帮助将不胜感激。

【问题讨论】:

  • "此代码有效,但并非总是如此。"你这是什么意思?你收到一些例外吗?哪些组件失败并详细说明失败的原因(必要时提供堆栈跟踪)。此外,您的链接 (http://site.website.com/path/) 指向 HTTP,并且您在访问 doneparts 时存在竞争条件。
  • 好点,我忘了说。它不会抛出异常,只是不会下载整个文件,只是下载其中的一部分。有时大约是文件的一半,有时可能是文件的 90%,有时是完整的文件。
  • “不抛出异常”是指它不抛出未捕获的异常还是根本不抛出任何异常(包括IOExceptionFTPException)。另外,为什么要为done 使用 HashMap?为什么不只是一个整数来计算完成了多少件。另外我们能看到检测文件是否完成的代码吗?
  • “不抛出异常”我的意思是它实际上不会抛出任何异常。程序正常终止,就像刚刚完成下载一样。我将编辑我的原始帖子以显示检测完整文件的代码。

标签: java multithreading download ftp inputstream


【解决方案1】:

首先,更换您的 FTP 服务器,有很多免费的 FTP 服务器支持任意文件大小并提供附加功能,但我离题了...

您的代码似乎有许多不相关的问题,这些问题都可能导致您看到的行为,解决方法如下:

  1. 从多个线程的未受保护/非同步访问中访问 doneparts 映射时存在竞争条件。这可能会导致数据损坏和线程之间这些变量的同步丢失,甚至可能导致done.values().contains(false) 返回true,即使它实际上不是。

  2. 您正在高频率地反复呼叫done.values().contains()。虽然 javadoc 没有明确说明,但哈希映射可能会以 O(n) 方式遍历每个值,以检查给定映射是否包含值。再加上其他线程正在修改地图这一事实,您将获得未定义的行为。根据values()javadoc:

    如果在对集合进行迭代时修改了映射(通过迭代器自己的删除操作除外),则迭代的结果是不确定的。

  3. 您以某种方式调用 new URL("http://site.website.com/path/" + file).openStream();,但声明您正在使用 FTP。链接中的http:// 定义了协议openStream() 尝试打开,http:// 不是ftp://。不确定这是一个错字还是你的意思是 HTTP(或者你有一个 HTTP 服务器提供相同的文件)。

  4. 鉴于并非所有部分都“完成”(基于您的忙等待循环设计),任何引发任何类型 Exception 的线程都会导致代码失败。当然,您可能会被修改一些其他逻辑来防止这种情况发生,但否则这是代码的潜在问题。

  5. 您没有关闭任何已打开的流。这可能意味着底层套接字本身也保持打开状态。这不仅构成资源泄漏,如果服务器本身有某种最大同时连接数限制,您只会导致新连接失败,因为旧的、已完成的传输没有关闭。

基于上述问题,我建议将下载逻辑移动到 Callable 任务中,并通过ExecutorService 运行它们,如下所示:

LinkedList<Callable<byte[]>> tasksToExecute = new LinkedList<>();

// Populate tasks to run
for(int i = 0; i < numParts; i++){
    final int part = i;

    // Lambda to 
    tasksToExecute.add(() -> {
        InputStream is = null;

        try{
            is = FTP.getAssetInputStream(asset, part);
            ByteArrayOutputStream baos = new ByteArrayOutputStream();

            byte[] buf = new byte[DOWNLOAD_BUFFER_SIZE];
            int len = 0;

            while((len = is.read(buf)) > 0){
                baos.write(buf, 0, len);
                curDownload.addAndGet(len);
                totAssets.addAndGet(len);
            }

            return baos.toByteArray();
        }catch(IOException e){
            // handle exception
        }catch(FTPException e){
            // handle exception
        }finally{
            if(is != null){
                try{
                    is.close();
                }catch(IOException ignored){}
            }
        }

        return null;
    });
}

// Retrieve an ExecutorService instance, note the use of work stealing pool is Java 8 only
// This can be substituted for newFixedThreadPool(nThreads) for Java < 8 as well for tight control over number of simultaneous links
ExecutorService executor = Executors.newWorkStealingPool(4);

// Tells the executor to execute all the tasks and give us the results
List<Future<byte[]>> resultFutures = executor.invokeAll(tasksToExecute);

// Populates the file
File assetFile = new File(dir, "assets/" + asset + ".raf");
assetFile.createNewFile();

try(FileOutputStream fos = new FileOutputStream(assetFile)){
    // Iterate through the futures, writing them to file in order
    for(Future<byte[]> result : resultFutures){
        byte[] partData = result.get();

        if(partData == null){
            // exception occured during downloading this part, handle appropriately
        }else{
            fos.write(partData);
        }
    }
}catch(IOException ex(){
    // handle exception
}

使用执行器服务,您可以进一步优化您的多线程方案,因为输出文件将在片段可用时立即开始写入(按顺序),并且线程本身被重用以节省线程创建成本。

如前所述,可能存在过多同时链接导致服务器拒绝连接的情况(或者更危险的是,编写 EOF 让您认为该部分已下载)。在这种情况下,可以通过newFixedThreadPool(nThreads) 调整工作线程的数量,以确保在任何给定时间,只有nThreads 的下载量可以同时发生。

【讨论】:

    猜你喜欢
    • 2011-01-10
    • 2011-07-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-01-18
    • 1970-01-01
    相关资源
    最近更新 更多