Skip to content

Commit

Permalink
Reverted camel-amqp version.
Browse files Browse the repository at this point in the history
  • Loading branch information
hekonsek committed Apr 27, 2016
1 parent 40f86a1 commit 9b2abab
Show file tree
Hide file tree
Showing 3 changed files with 9 additions and 3 deletions.
5 changes: 5 additions & 0 deletions bom/pom.xml
Expand Up @@ -514,6 +514,11 @@
<!-- Camel Spark is available starting from Camel 2.17 -->
<version>2.17-SNAPSHOT</version>
</dependency>
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-amqp</artifactId>
<version>2.16.2</version>
</dependency>
<dependency>
<groupId>org.apache.commons</groupId>
<artifactId>commons-exec</artifactId>
Expand Down
Expand Up @@ -20,7 +20,7 @@
import static io.rhiot.cloudplatform.runtime.spring.RhiotConstants.CAMEL_SPRINGBOOT_TYPE_CONVERSION;
import static io.rhiot.cloudplatform.runtime.spring.RhiotConstants.META_INF_RHIOT_BANNER_TXT;
import static java.util.Arrays.asList;
import static org.apache.camel.component.amqp.AMQPComponent.amqpComponent;
import static org.apache.camel.component.amqp.AMQPComponent.amqp10Component;

import io.rhiot.cloudplatform.encoding.spi.PayloadEncoding;
import io.rhiot.cloudplatform.connector.IoTConnector;
Expand Down Expand Up @@ -90,7 +90,7 @@ public static void main(String[] args) throws InterruptedException {
AMQPComponent amqp(@Value("${AMQP_SERVICE_HOST:localhost}") String amqpBrokerUrl,
@Value("${AMQP_SERVICE_PORT:5672}") int amqpBrokerPort) throws MalformedURLException {
LOG.debug("About to create AMQP component {}:{}", amqpBrokerUrl, amqpBrokerPort);
return amqpComponent("amqp://" + amqpBrokerUrl + ":" + amqpBrokerPort);
return amqp10Component("amqp://" + amqpBrokerUrl + ":" + amqpBrokerPort);
}

@Bean
Expand Down
Expand Up @@ -19,6 +19,7 @@
import io.rhiot.cloudplatform.encoding.spi.PayloadEncoding;
import org.apache.camel.Message;
import org.apache.camel.builder.RouteBuilder;
import org.apache.qpid.amqp_1_0.jms.impl.QueueImpl;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand Down Expand Up @@ -62,7 +63,7 @@ public void configure() throws Exception {

from(fromChannel).process(exchange -> {
Message message = exchange.getIn();
String channel = message.getHeader("JMSDestination", String.class);
String channel = message.getHeader("JMSDestination", QueueImpl.class).getQueueName();
byte[] incomingPayload = message.getBody(byte[].class);
OperationBinding operationBinding = operationBinding(payloadEncoding, channel, incomingPayload, message.getHeaders(), getContext().getRegistry());
exchange.setProperty(TARGET_PROPERTY, "bean:" + operationBinding.service() + "?method=" + operationBinding.operation() + "&multiParameterArray=true");
Expand Down

0 comments on commit 9b2abab

Please sign in to comment.