【问题标题】:How can i make a camel route thread safe?我怎样才能使骆驼路线线程安全?
【发布时间】:2017-04-03 20:46:15
【问题描述】:

我有一条骆驼路线,它消耗来自 apache activeMQ 的任务。 当 ActiveMQ 只有一个消费者时,一切正常。

但是当消费者增加到两个(或更多)时,应用程序的行为就会不恰当。

这是我的路线:

<routeContext id="myRoute"  xmlns="http://camel.apache.org/schema/spring">
    <route errorHandlerRef="myErrorHandler" id="myErrorRoute">
        <from uri="activemq:queue:{{my.queue}}" />
        <log loggingLevel="DEBUG" message="Message received from my queue : ${body}"></log>
        <multicast>
            <pipeline>
                <log loggingLevel="DEBUG" message="Adding to redis : ${body}"></log>
                <to uri="spring-redis://localhost:6379?serializer=#stringSerializer" />
            </pipeline>
            <pipeline>
                <transform>
                    <method ref="insertBean" method="myBatchInsertion"></method>
                </transform>
                <choice>
                    <when>
                        <simple> ${body.size()} == ${properties:my.batch.size}</simple>
                        <log message="Going to insert my batch in database" />
                        <to uri="mybatis:batchInsert?statementType=InsertList"></to>
          <log message="Inserted in my table : ${in.header.CamelMyBatisResult}"></log>
                        <choice>
                            <when>
                                <simple>${properties:my.write.file} == true</simple>
                                <bean beanType="com.***.***.processors.InsertToFile"
                                    method="processMy(${exchange}" />
                                <log message="Going to write to file : ${in.header.CamelFileName}" />
                                <to uri="file://?fileExist=Append&amp;bufferSize=32768"></to>
                            </when>
                        </choice>
                    </when>
                </choice>
            </pipeline>
        </multicast>
    </route>
</routeContext>

下面是豆子:

public class InsertBeanImpl {
  public  List<Out> myOutList = new CopyOnWriteArrayList<Out>();

  public List<Out> myBatchInsertion(Exchange exchange) {
        if (myOutList.size() >= myBatchSize) {
            Logger.sysLog(LogValues.info,this.getClass().getName(),"Reached max PayLoad size : "+myOutList.size() + " , going to clear batch");
            myOutList.clear();
        }
        Out myOut = exchange.getIn().getBody(Out.class);
        Logger.sysLog(LogValues.APP_INFO, this.getClass().getName(), myOut.getMasterId()+" | "+"Adding to batch masterId : "+myOut.getMasterId());
        synchronized(myOut){
            myOutList.add(myOut);
        }
        Logger.sysLog(LogValues.info, this.getClass().getName(), "Count of batch : "+myOutList.size());
        return myOutList;
    }
}



public class SaveToFile {
    static String currentFileName = null;
    static int iSub = 0;
    String path;
    String absolutePath;

    @Autowired
    private Utility utility;


    public void processMy(Exchange exchange) {
        getFileName(exchange, currentFileNameSub, iSub);
    }


    public void getFileName(Exchange exchange, String outFile, int i) {
        exchange.getIn().setBody(getFromJson(exchange));
        path = (String) exchange.getIn().getHeader("path");
        Calendar date = null;
        date = new GregorianCalendar();
        NumberFormat format = NumberFormat.getIntegerInstance();
        format.setMinimumIntegerDigits(2);
        String pathSuffix = "/" + date.get(Calendar.YEAR) + "/"
                + format.format((date.get(Calendar.MONTH) + 1)) + "/"
                    + format.format(date.get(Calendar.DAY_OF_MONTH));
        String fileName = new SimpleDateFormat("yyyy-MM-dd").format(new Date());
        double megabytes = 100 * 1024 * 1024;
        if (outFile != null) {
            if (!fileName.equals(outFile.split("_")[0])) {
                outFile = null;
                i = 0;
            }
        }
        if (outFile == null) {
            outFile = fileName + "_" + i;
        }
        while (new File(path + "/" + pathSuffix + "/" + outFile).length() >= megabytes) {
            outFile = fileName + "_" + (++i);
        }
        absolutePath = path + "/" + pathSuffix + "/" + outFile;
        exchange.getIn().setHeader("CamelFileName", absolutePath);
    }

    public String getFromJson(Exchange exchange) {
        synchronized(exchange){
            List<Out> body = exchange.getIn().getBody(CopyOnWriteArrayList.class);
            Logger.sysLog(LogValues.info, this.getClass().getName(), "body > "+body.size());
            String text = "";
            for (int i = 0; i < body.size(); i++) {
                Out msg = body.get(i);
                text = text.concat(utility.convertObjectToJsonStr(msg) + "\n");
            }
            return text;
        }

    }
}

由于处理器不同步且不是线程安全的,因此在多个消费者的情况下路由不会按预期工作。

谁能告诉我如何使我的路由线程安全或同步?

我尝试使处理器同步,但没有帮助。还有其他方法吗?

【问题讨论】:

  • 我认为您需要准确解释这里的问题所在。如果你的处理器不是线程安全的,那么问题就不是 Camel 的错。
  • @AdamHawkes 请查看更新后的问题。那是我的问题,我怎样才能使它线程安全。我一点也不怪骆驼:P
  • 确定哪些 bean/组件存在问题会很有帮助,然后我们可以深入了解如何解决这些特定的代码。
  • 通过添加 com.***.***.processors.InsertToFile 和 "insertBean" bean 的源代码来更新您的问题。我怀疑其中任何一个都不是线程安全的。
  • @TadayoshiSato 请查看更新后的问题。

标签: java thread-safety apache-camel activemq


【解决方案1】:

让你的处理器线程安全,不要过多地使用 Synchronize,或者根本不使用。

为此,您根本不能在其中包含可更改的实例变量。

只有公共属性可以存在,例如某些设置对所有线程都有效,并且在任何方法执行期间都不会更改。那么就不需要使用影响性能的Synchronize机制了。

Camel 中的其他所有内容都是线程安全的,如果处理器实现不是线程安全的,则无法告诉 Camel 是线程安全的。

【讨论】:

    【解决方案2】:

    显然InsertBeanImpl 和SaveToFile 都不是线程安全的。一般来说,Camel 路由中使用的 bean 应该是无状态的,即它们不应该有可变字段。

    对于InsertBeanImpl,看起来您真正想做的是将多条消息聚合为一条。对于这样的用例,我会考虑改用 Camel 聚合器 [1],这样您可以更轻松地实现线程安全的解决方案。

    [1]http://camel.apache.org/aggregator.html

    对于SaveToFile,我认为没有理由将path 和absolutePath 作为一个字段。将它们移动到 getFileName 方法中的局部变量中。

    【讨论】:

      猜你喜欢
      • 2012-04-15
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-09-05
      • 1970-01-01
      • 2011-05-24
      相关资源
      最近更新 更多