架构

生产者模板

师成师成· 更新于 2026-09-28· 阅读 14 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

ProducerTemplate

ProducerTemplate 接口允许你以多种不同的方式向端点发送消息交换(message exchange),从而方便地在 Java 代码中使用 Camel 的 Endpoint 实例。

如果你只是想向同一个端点发送大量消息,可以为它配置一个默认端点;也可以把 Endpoint 或 uri 作为第一个参数传入。

你不应为每次消息调用都创建一个 ProducerTemplate;应该在启动时创建一个实例并一直持有它。另外,在使用完 ProducerTemplate 之后,应调用 stop() 方法来关闭它所使用的全部资源。

sendBody() 方法可以让你轻松地向端点发送任意对象,示例如下:

仅使用 Java:通过 ProducerTemplate 发送消息

ProducerTemplate template = exchange.getContext().createProducerTemplate();

// send to default endpoint
template.sendBody("<hello>world!</hello>");

// send to a specific queue
template.sendBody("activemq:MyQueue", "<hello>world!</hello>");

// send with a body and header
template.sendBodyAndHeader("activemq:MyQueue",
   "<hello>world!</hello>",
   "CustomerRating", "Gold");

你还可以提供一个 Exchange 或 Processor 来自定义交换。

Send 与 Request 方法

ProducerTemplate 支持消息交换模式(MEP),用于控制所使用的消息传递方式:

换句话说,ProducerTemplate 上所有以 sendXXX 开头的方法用于 InOnly 消息传递,所有以 requestXXX 开头的方法用于 InOut 消息传递。

下面看一个调用端点以获取响应(InOut)的示例:

纯 Java 方式:使用 requestBody 进行请求-应答消息传递

Object response = template.requestBody("<hello/>");

// you can type convert the response to what you want such as String
String ret = template.requestBody("<hello/>", String.class);

// or specify the endpoint uri in the method
String ret = template.requestBody("cxf:bean:HelloWorldService", "<hello/>", String.class);

流式接口

FluentProducerTemplate 在常规的 ProducerTemplate 之上提供了流式语法。

以下是一些示例:

设置消息头与消息体

这是使用流式构建器设置消息头和消息体最常见的写法:

仅 Java:使用 FluentProducerTemplate 设置消息头与消息体

Integer result = FluentProducerTemplate.on(context)
    .withHeader("key-1", "value-1")
    .withHeader("key-2", "value-2")
    .withBody("Hello")
    .to("direct:inout")
    .request(Integer.class);

使用处理器

这里我们使用 Processor 来准备要发送的消息。

仅限 Java:使用 FluentProducerTemplate 配合 Processor

Integer result = FluentProducerTemplate.on(context)
    .withProcessor(exchange -> exchange.getIn().setBody("Hello World"))
    .to("direct:exception")
    .request(Integer.class);

使用模板定制器进行高级配置

这种用法很少见,但 TemplateCustomizer 可用于高级场景,以控制 FluentProducerTemplate 的各个方面,例如配置使用自定义线程池:

仅 Java:在 FluentProducerTemplate 中使用模板定制器

Object result = FluentProducerTemplate.on(context)
    .withTemplateCustomizer(
        template -> {
            template.setExecutorService(myExecutor);
            template.setMaximumCacheSize(10);
        }
    )
    .withBody("the body")
    .to("direct:start")
    .request();

如何获取处理 Exchange 时抛出的异常?

你已将一个 Exchange 发送给 Camel,但它在处理过程中因抛出异常而失败。如何获取这个异常?

如果你使用的是 CamelTemplate(或 CamelProducer),通常会使用 sendBody/requestBody 方法,这些方法只返回 Exchange 的响应正文。因此,如果处理过程中抛出了异常,Camel 并不会重新抛出该异常。要解决这个问题,你可以使用普通的 send/request 方法,这些方法接受一个 Exchange 对象并返回一个 Exchange 对象。

通过返回的 Exchange,你可以检测它是否失败,并获取导致失败的异常。下面的代码示例展示了这一用法:

纯 Java:从返回的 Exchange 中获取抛出的异常

@Test
public void testOk() {
    int result = (Integer) template.sendBody("direct:input", ExchangePattern.InOut, "Hello London");
    assertEquals(1, result);
}

@Test
public void testFailure() {
    // must create an exchange to get the result as an exchange where we can get the caused exception
    Exchange exchange = getMandatoryEndpoint("direct:input").createExchange(ExchangePattern.InOut);
    exchange.getIn().setBody("Hello Paris");

    Exchange out = template.send("direct:input", exchange);
    assertTrue("Should be failed", out.isFailed());
    assertTrue("Should be IllegalArgumentException", out.getException() instanceof IllegalArgumentException);
    assertEquals("Forced exception", out.getException().getMessage());
}

protected RouteBuilder createRouteBuilder() throws Exception {
    return new RouteBuilder() {
        public void configure() throws Exception {
            from("direct:input").bean(new ExceptionBean());
        }
    };
}

public static class ExceptionBean {
    public int doSomething(String request) throws Exception {
        if (request.equals("Hello London")) {
            return 1;
        } else {
            throw new IllegalArgumentException("Forced exception");
        }
    }
}

配置默认缓存大小

你可以全局配置 ProducerTemplate 和 ConsumerTemplate 的默认缓存大小,这两个模板将由 CamelContext 创建或通过依赖注入获得。

可以通过在 CamelContext 上设置全局选项来完成,如以下 Java 代码所示:

仅限 Java:通过 CamelContext 全局选项配置缓存大小

getCamelContext().getGlobalOptions().put(Exchange.MAXIMUM_CACHE_POOL_SIZE, "50");

或者在 application.properties 中:

camel.main.producerTemplateCacheSize = 50

默认的最大缓存大小为 1000。

重新使用 ProducerTemplate

在使用 ProducerTemplate 时,理想情况下应在 Camel 应用程序的整个生命周期内重复使用该模板。

一个常见的错误做法是在 Processor 或 bean 方法调用中创建新的 ProducerTemplate。

你不应该为每条消息的调用都创建一个 ProducerTemplate;正确的做法是在启动时创建一个实例并一直持有它。

使用完 ProducerTemplate 后,你应该调用 stop() 方法来关闭它所使用的所有资源。

如果不这样做,Camel 应用程序可能会创建越来越多的资源(线程等)。

更好的做法是显式地在启动时创建一个 ProducerTemplate,并将其注入到你的 Processor 或 bean 中,以便重复使用该模板。

另请参阅

参见 ConsumerTemplate

评论

登录后参与评论

正在加载评论…