【问题标题】:Mule File Inbound Flow : Control Number of threadsMule 文件入站流程:控制线程数
【发布时间】:2014-06-18 20:41:30
【问题描述】:

我想控制文件入站和消息处理器中的线程数。假设如果我的输入目录中有 5 个文件,那么我应该能够一次处理 2 个文件。一旦处理了这些文件(文件内容由消息处理器处理),那么只有它应该拾取其他文件。我曾尝试在流级别使用同步处理策略,但它只处理一个文件,我想要多个线程,但每个线程将直接从接收文件处理消息以发送响应。我尝试了大卫建议的方法,但它也不起作用。一次只拾取一个文件。

<flow name="fileInboundTestFlow2" doc:name="fileInboundTestFlow2" processingStrategy="synchronous">
    <poll frequency="1000">
        <component class="FilePollerComponent" doc:name="File Poller"></component>
    </poll>
    <collection-splitter />
    <request-reply >
            <vm:outbound-endpoint path="out"/>
            <vm:inbound-endpoint path="response">
                               <collection-aggregator/>
            </vm:inbound-endpoint>                      
    </request-reply>
    <file:outbound-endpoint path="E:/fileTest/processed" />
</flow>

public class FilePollerComponent implements Callable{

private String pollDir="E://fileTest" ;

private int numberOfFiles = 3;

public String getPollDir()
{
    return pollDir;
}

public void setPollDir(String pollDir)
{
    this.pollDir = pollDir;
}



public int getNumberOfFiles()
{
    return numberOfFiles;
}

public void setNumberOfFiles(int numberOfFiles)
{
    this.numberOfFiles = numberOfFiles;
}

@Override
public Object onCall(MuleEventContext eventContext) throws Exception
{
    File f = new File(pollDir);
    List<File> filesToReturn = new ArrayList<File>(numberOfFiles);
    if(f.isDirectory())
    {
        File[] files = f.listFiles();
        int i = 0;
        for(File file : files)
        {
            if(file.isFile())
                filesToReturn.add(file);
            if(i==numberOfFiles)
                break ;
            i++;
        }
    }
    else
    {
        throw new Exception("Invalid Directory");
    }
    return filesToReturn;
}}

【问题讨论】:

    标签: mule


    【解决方案1】:

    文件入站端点是一个轮询器,因此它使用一个线程。如果您使流程同步,您将搭载这个单线程,因此一次处理一个文件。

    您需要创建一个只允许 2 个线程的流处理策略。以允许 500 个线程为例:http://www.mulesoft.org/documentation/display/current/Flow+Processing+Strategies#FlowProcessingStrategies-Fine-TuningaQueued-AsynchronousProcessingStrategy

    编辑:上述提案不满足此要求:

    如果我配置了 3 个线程,那么将从输入目录中选择三个文件,并且在处理完这些文件之前,不应选择其他文件

    确实,上述提议总是并行处理 3 个文件,而不是按 3 个为一组进行处理。

    所以我提出了这种替代方法:

    • 将流处理策略配置为同步
    • 使用poll 元素作为源
    • 在其中放置一个自定义组件,该组件从可配置目录中选择 3 个不同的文件。无需锁定任何东西,因为流的同步策略可以防止重新进入。返回java.io.Files 中的java.util.List
    • 在其后添加collection-splitter
    • 在入站端点中添加带有aggregatorrequest-reply 以实现fork-join 模式(http://blogs.mulesoft.org/aggregation-with-mule-fork-and-join-pattern/)。文件的处理将在另一个流程中进行,而轮询流程将阻塞,直到处理完所有 3 个文件。

    【讨论】:

    • 我已经尝试过了,但它正在无限循环中。它继续获取文件上的锁。我想要的是,单线程应该接收、处理和发送消息。因此,如果我配置 3 个线程,则将从输入目录中选择三个文件,并且在处理这些文件之前,不应选择其他文件。请帮忙
    • 无限循环听起来像一个错误。我已经通过答案进行了审查。
    • 如果我使用同步流。即使文件列表有三个,它仍然一次拾取一个文件。
    • 您错过了vm:inbound-endpoint 中的聚合器来阻止流执行,直到 3 个文件的处理完成。
    • 仍然只拾取一个文件,它没有按需要拾取 3 个文件。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2012-09-21
    • 1970-01-01
    • 1970-01-01
    • 2014-08-26
    • 2011-05-07
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多