diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java
index 229e383039..2decfbdd9f 100644
--- a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java
+++ b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/AbstractHttpRequestExecutingMessageHandler.java
@@ -61,7 +61,6 @@ import org.springframework.util.LinkedMultiValueMap;
import org.springframework.util.MultiValueMap;
import org.springframework.util.StringUtils;
import org.springframework.web.util.DefaultUriBuilderFactory;
-import org.springframework.web.util.UriComponentsBuilder;
/**
* Base class for http outbound adapter/gateway.
@@ -123,7 +122,7 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac
* {@link org.springframework.web.client.RestTemplate}. The default value is
* true.
* @param encodeUri true if the URI should be encoded.
- * @see UriComponentsBuilder
+ * @see org.springframework.web.util.UriComponentsBuilder
* @deprecated since 5.3 in favor of {@link #setEncodingMode}
*/
@Deprecated
@@ -295,6 +294,7 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac
}
@Override
+ @Nullable
protected Object handleRequestMessage(Message> requestMessage) {
HttpMethod httpMethod = determineHttpMethod(requestMessage);
if (this.extractPayloadExplicitlySet && logger.isWarnEnabled() && !shouldIncludeRequestBody(httpMethod)) {
@@ -320,34 +320,32 @@ public abstract class AbstractHttpRequestExecutingMessageHandler extends Abstrac
return exchange(uri, httpMethod, httpRequest, expectedResponseType, requestMessage, uriVariables);
}
+ @Nullable
protected abstract Object exchange(Object uri, HttpMethod httpMethod, HttpEntity> httpRequest,
Object expectedResponseType, Message> requestMessage, Map uriVariables);
protected Object getReply(ResponseEntity> httpResponse) {
- if (this.expectReply) {
- HttpHeaders httpHeaders = httpResponse.getHeaders();
- Map headers = this.headerMapper.toHeaders(httpHeaders);
- if (this.transferCookies) {
- doConvertSetCookie(headers);
- }
-
- AbstractIntegrationMessageBuilder> replyBuilder;
- MessageBuilderFactory messageBuilderFactory = getMessageBuilderFactory();
- if (httpResponse.hasBody()) {
- Object responseBody = httpResponse.getBody();
- replyBuilder = (responseBody instanceof Message>)
- ? messageBuilderFactory.fromMessage((Message>) responseBody)
- : messageBuilderFactory.withPayload(responseBody); // NOSONAR - hasBody()
-
- }
- else {
- replyBuilder = messageBuilderFactory.withPayload(httpResponse);
- }
- replyBuilder.setHeader(org.springframework.integration.http.HttpHeaders.STATUS_CODE,
- httpResponse.getStatusCode());
- return replyBuilder.copyHeaders(headers);
+ HttpHeaders httpHeaders = httpResponse.getHeaders();
+ Map headers = this.headerMapper.toHeaders(httpHeaders);
+ if (this.transferCookies) {
+ doConvertSetCookie(headers);
}
- return null;
+
+ AbstractIntegrationMessageBuilder> replyBuilder;
+ MessageBuilderFactory messageBuilderFactory = getMessageBuilderFactory();
+ if (httpResponse.hasBody()) {
+ Object responseBody = httpResponse.getBody();
+ replyBuilder = (responseBody instanceof Message>)
+ ? messageBuilderFactory.fromMessage((Message>) responseBody)
+ : messageBuilderFactory.withPayload(responseBody); // NOSONAR - hasBody()
+
+ }
+ else {
+ replyBuilder = messageBuilderFactory.withPayload(httpResponse);
+ }
+ replyBuilder.setHeader(org.springframework.integration.http.HttpHeaders.STATUS_CODE,
+ httpResponse.getStatusCode());
+ return replyBuilder.copyHeaders(headers);
}
/**
diff --git a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java
index d6f7fa40a0..bbc2bed4aa 100755
--- a/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java
+++ b/spring-integration-http/src/main/java/org/springframework/integration/http/outbound/HttpRequestExecutingMessageHandler.java
@@ -171,6 +171,7 @@ public class HttpRequestExecutingMessageHandler extends AbstractHttpRequestExecu
}
@Override
+ @Nullable
protected Object exchange(Object uri, HttpMethod httpMethod, HttpEntity> httpRequest,
Object expectedResponseType, Message> requestMessage, Map uriVariables) {
@@ -197,7 +198,13 @@ public class HttpRequestExecutingMessageHandler extends AbstractHttpRequestExecu
}
}
- return getReply(httpResponse);
+ if (isExpectReply()) {
+ return getReply(httpResponse);
+ }
+ else {
+ return null;
+ }
+
}
catch (RestClientException e) {
throw new MessageHandlingException(requestMessage,
diff --git a/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/outbound/WebFluxRequestExecutingMessageHandler.java b/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/outbound/WebFluxRequestExecutingMessageHandler.java
index 8fcd0402a6..525ad8565f 100644
--- a/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/outbound/WebFluxRequestExecutingMessageHandler.java
+++ b/spring-integration-webflux/src/main/java/org/springframework/integration/webflux/outbound/WebFluxRequestExecutingMessageHandler.java
@@ -204,9 +204,82 @@ public class WebFluxRequestExecutingMessageHandler extends AbstractHttpRequestEx
}
@Override
+ @Nullable
protected Object exchange(Object uri, HttpMethod httpMethod, HttpEntity> httpRequest,
Object expectedResponseType, Message> requestMessage, Map uriVariables) {
+ WebClient.RequestBodySpec requestSpec =
+ createRequestBodySpec(uri, httpMethod, httpRequest, requestMessage, uriVariables);
+
+ Mono responseMono = exchangeForResponseMono(requestSpec);
+
+ if (isExpectReply()) {
+ return createReplyFromResponse(expectedResponseType, responseMono);
+ }
+ else {
+ responseMono.subscribe(v -> { }, ex -> sendErrorMessage(requestMessage, ex));
+ return null;
+ }
+ }
+
+ private Object createReplyFromResponse(Object expectedResponseType, Mono responseMono) {
+ return responseMono
+ .flatMap(response -> {
+ ResponseEntity.BodyBuilder httpEntityBuilder =
+ ResponseEntity.status(response.statusCode())
+ .headers(response.headers().asHttpHeaders());
+
+ Mono> bodyMono;
+
+ if (expectedResponseType != null) {
+ if (this.replyPayloadToFlux) {
+ BodyExtractor extends Flux>, ReactiveHttpInputMessage> extractor;
+ if (expectedResponseType instanceof ParameterizedTypeReference>) {
+ extractor = BodyExtractors.toFlux(
+ (ParameterizedTypeReference>) expectedResponseType);
+ }
+ else {
+ extractor = BodyExtractors.toFlux((Class>) expectedResponseType);
+ }
+ Flux> flux = response.body(extractor);
+ bodyMono = Mono.just(flux);
+ }
+ else {
+ BodyExtractor extends Mono>, ReactiveHttpInputMessage> extractor;
+ if (expectedResponseType instanceof ParameterizedTypeReference>) {
+ extractor = BodyExtractors.toMono(
+ (ParameterizedTypeReference>) expectedResponseType);
+ }
+ else {
+ extractor = BodyExtractors.toMono((Class>) expectedResponseType);
+ }
+ bodyMono = response.body(extractor);
+ }
+ }
+ else if (this.bodyExtractor != null) {
+ Object body = response.body(this.bodyExtractor);
+ if (body instanceof Mono) {
+ bodyMono = (Mono>) body;
+ }
+ else {
+ bodyMono = Mono.just(body);
+ }
+ }
+ else {
+ bodyMono = Mono.empty();
+ }
+
+ return bodyMono
+ .map(httpEntityBuilder::body)
+ .defaultIfEmpty(httpEntityBuilder.build());
+ }
+ )
+ .map(this::getReply);
+ }
+
+ private WebClient.RequestBodySpec createRequestBodySpec(Object uri, HttpMethod httpMethod,
+ HttpEntity> httpRequest, Message> requestMessage, Map uriVariables) {
+
WebClient.RequestBodyUriSpec requestBodyUriSpec = this.webClient.method(httpMethod);
WebClient.RequestBodySpec requestSpec;
@@ -222,100 +295,42 @@ public class WebFluxRequestExecutingMessageHandler extends AbstractHttpRequestEx
if (inserter != null) {
requestSpec.body(inserter);
}
+ return requestSpec;
+ }
- Mono responseMono =
- requestSpec.exchange()
- .flatMap(response -> {
- HttpStatus httpStatus = response.statusCode();
- if (httpStatus.isError()) {
- return response.body(BodyExtractors.toDataBuffers())
- .reduce(DataBuffer::write)
- .map(dataBuffer -> {
- byte[] bytes = new byte[dataBuffer.readableByteCount()];
- dataBuffer.read(bytes);
- DataBufferUtils.release(dataBuffer);
- return bytes;
- })
- .defaultIfEmpty(new byte[0])
- .map(bodyBytes -> {
- throw new WebClientResponseException(
- "ClientResponse has erroneous status code: "
- + httpStatus.value() + " "
- + httpStatus.getReasonPhrase(),
- httpStatus.value(),
- httpStatus.getReasonPhrase(),
- response.headers().asHttpHeaders(),
- bodyBytes,
- response.headers().contentType()
- .map(MimeType::getCharset)
- .orElse(StandardCharsets.ISO_8859_1));
- }
- );
- }
- else {
- return Mono.just(response);
- }
- });
-
- if (isExpectReply()) {
- return responseMono
- .flatMap(response -> {
- ResponseEntity.BodyBuilder httpEntityBuilder =
- ResponseEntity.status(response.statusCode())
- .headers(response.headers().asHttpHeaders());
-
- Mono> bodyMono;
-
- if (expectedResponseType != null) {
- if (this.replyPayloadToFlux) {
- BodyExtractor extends Flux>, ReactiveHttpInputMessage> extractor;
- if (expectedResponseType instanceof ParameterizedTypeReference>) {
- extractor = BodyExtractors.toFlux(
- (ParameterizedTypeReference>) expectedResponseType);
+ private Mono exchangeForResponseMono(WebClient.RequestBodySpec requestSpec) {
+ return requestSpec.exchange()
+ .flatMap(response -> {
+ HttpStatus httpStatus = response.statusCode();
+ if (httpStatus.isError()) {
+ return response.body(BodyExtractors.toDataBuffers())
+ .reduce(DataBuffer::write)
+ .map(dataBuffer -> {
+ byte[] bytes = new byte[dataBuffer.readableByteCount()];
+ dataBuffer.read(bytes);
+ DataBufferUtils.release(dataBuffer);
+ return bytes;
+ })
+ .defaultIfEmpty(new byte[0])
+ .map(bodyBytes -> {
+ throw new WebClientResponseException(
+ "ClientResponse has erroneous status code: "
+ + httpStatus.value() + " "
+ + httpStatus.getReasonPhrase(),
+ httpStatus.value(),
+ httpStatus.getReasonPhrase(),
+ response.headers().asHttpHeaders(),
+ bodyBytes,
+ response.headers().contentType()
+ .map(MimeType::getCharset)
+ .orElse(StandardCharsets.ISO_8859_1));
}
- else {
- extractor = BodyExtractors.toFlux((Class>) expectedResponseType);
- }
- Flux> flux = response.body(extractor);
- bodyMono = Mono.just(flux);
- }
- else {
- BodyExtractor extends Mono>, ReactiveHttpInputMessage> extractor;
- if (expectedResponseType instanceof ParameterizedTypeReference>) {
- extractor = BodyExtractors.toMono(
- (ParameterizedTypeReference>) expectedResponseType);
- }
- else {
- extractor = BodyExtractors.toMono((Class>) expectedResponseType);
- }
- bodyMono = response.body(extractor);
- }
- }
- else if (this.bodyExtractor != null) {
- Object body = response.body(this.bodyExtractor);
- if (body instanceof Mono) {
- bodyMono = (Mono>) body;
- }
- else {
- bodyMono = Mono.just(body);
- }
- }
- else {
- bodyMono = Mono.empty();
- }
-
- return bodyMono
- .map(httpEntityBuilder::body)
- .defaultIfEmpty(httpEntityBuilder.build());
- }
- )
- .map(this::getReply);
- }
- else {
- responseMono.subscribe(v -> { }, ex -> sendErrorMessage(requestMessage, ex));
-
- return null;
- }
+ );
+ }
+ else {
+ return Mono.just(response);
+ }
+ });
}
@Nullable