【问题标题】:Springboot v2.0.0.M6 WebClient making multiple duplicate HTTP POST callsSpringboot v2.0.0.M6 WebClient 进行多次重复的 HTTP POST 调用
【发布时间】:2019-04-07 08:53:03
【问题描述】:

我使用的是 spring-boot 版本 2.0.0.M6。 我需要从spring-boot应用程序说APP1到另一个应用程序(播放框架)说APP2进行异步HTTP调用。 因此,如果我需要从 APP1 到 APP2 进行 20 次不同的异步调用,APP2 会收到 20 个请求,其中很少有重复请求,这意味着这些重复请求替换了几个不同的请求。 预期:

api/v1/call/1
api/v1/call/2
api/v1/call/3
api/v1/call/4

实际:

api/v1/call/1
api/v1/call/2
api/v1/call/4
api/v1/call/4

我正在使用 spring 响应式 WebClient。

下面是build.gradle中的spring boot版本

buildscript {
ext {
    springBootVersion = '2.0.0.M6'
    //springBootVersion = '2.0.0.BUILD-SNAPSHOT'
}
repositories {
    mavenCentral()
    maven { url "https://repo.spring.io/snapshot" }
    maven { url "https://repo.spring.io/milestone" }
    maven {url "https://plugins.gradle.org/m2/"}
}
dependencies {
    classpath("org.springframework.boot:spring-boot-gradle-plugin:${springBootVersion}")
    classpath("se.transmode.gradle:gradle-docker:1.2")


}
}

我的 WebClient 初始化 sn-p

private WebClient webClient = WebClient.builder()
        .clientConnector(new ReactorClientHttpConnector((HttpClientOptions.Builder builder) -> builder.disablePool()))
        .build();

我的 POST 方法

public <T> Mono<JsonNode> postClient(String url, T postData) {
    return Mono.subscriberContext().flatMap(ctx -> {
        String cookieString = ctx.getOrDefault(Constants.SubscriberContextConstnats.COOKIES, StringUtils.EMPTY);
        URI uri = URI.create(url);
        return webClient.post().uri(uri).body(BodyInserters.fromObject(postData)).header(HttpHeaders.COOKIE, cookieString)
          .exchange().flatMap(clientResponse ->
          {
              return clientResponse.bodyToMono(JsonNode.class);
          })
         .onErrorMap(err -> new TurtleException(err.getMessage(), err))
         .doOnSuccess(jsonData -> {
         });
    });
}

调用此 postClient 方法的代码

private void getResultByKey(PremiumRequestHandler request, String key, BrokerConfig brokerConfig) {

    /* Live calls for the insurers */
    LOG.info("[PREMIUM SERVICE] LIVE CALLLLL MADE FOR: " + key + " AND REQUEST ID: " + request.getRequestId());

    String uri = brokerConfig.getHostUrl() + verticalResolver.determineResultUrl(request.getVertical()) + key;
    LOG.info("[PREMIUM SERVICE] LIVE CALL WITH URI : " + uri + " FOR REQUEST ID: " + request.getRequestId());
    Mono<PremiumResponse> premiumResponse = reactiveWebClient.postClient(uri, request.getPremiumRequest())
            .map(json -> PlatformUtils.mapToClass(json, PremiumResponse.class));

    premiumResponse.subscribe(resp -> {
        resp.getPremiumResults().forEach(result -> {
            LOG.info("Key " + result.getKey());

            repository.getResultRepoRawType(request.getVertical())
                    .save(result).subscribe();
            saveResult.subscriberContext(ctx -> {

                MultiBrokerMongoDBFactory.setDatabaseNameForCurrentThread(brokerConfig.getBroker());
                return ctx;
            }).subscribe();
        });
    }, error -> {
        LOG.info("[PREMIUM SERVICE] ERROR RECIEVED FOR " + key + " AND REQUEST ID" + request.getRequestId() + " > " + error.getMessage());
    });

}

已将日志放在客户端代码的端点,此时无法看到多个请求。

这可能是 WebClient 中的一个错误,其中 URI 在多线程环境中被交换。

尝试改变 WebClient,但 URI 仍然被交换

请帮忙。

Git 仓库添加 github.com/praveenk007/ps-demo

【问题讨论】:

  • 我不认为你的代码 sn-p 编译,因为它指的是未知变量(如saveResult)。一般来说,应该首先尝试使用最新的 GA 版本(2.1.0),因为这个 Milestone 版本已经过时了。此外,使代码 sn-p 更简单并减少噪音(通过数据库调用和其他东西)会有所帮助。
  • 尝试将 springboot 版本升级到 2.1.0 但在运行时给出 noMethod 错误。需要检查文档。同时会在 GitHub 上做一个单独的项目来复制这个问题。
  • 如果您尝试复制此问题,您可以使用 postClient 代码并通过传递动态 url 和一些静态 post 对象从循环中调用它。循环运行 20-25 次。
  • 能否请您分享您的代码或回购中的示例?因为我没有与此相关的问题。
  • @JonathanJohx 这是一个独立代码的仓库。但无法在此处复制错误。 github.com/praveenk007/ps-demo

标签: java spring-boot spring-webflux


【解决方案1】:

我碰巧遇到过类似的问题:

当并行 (webflux) 调用相同的服务(此处标记为 ExternalService)时,有时会发送相同的请求,而问题似乎确实存在于 Webclient 中。

解决方案是改变 Webclient 的创建方式。

之前:

这里ExternalCall是配置的客户端,它调用ExternalService。所以这里的并行执行是 internalCall 方法。请注意,我们将 WebClient.RequestBodySpec 类传递给 ExternalCall

@Bean
ExternalCall externalCall(WebClient.Builder webClientBuilder) {

    var exchangeStrategies = getExchangeStrategies();

    var endpoint = "v1/data";
    var timeout = 10000;
    var uri = "https://externalService.com/";

    var requestBodySpec = webClientBuilder.clone()
            .clientConnector(new ReactorClientHttpConnector(HttpClient.create()))
            .exchangeStrategies(exchangeStrategies)
            .build()
            .post()
            .uri(endpoint)
            .accept(TEXT_XML)
            .contentType(TEXT_XML);

    return new ExternalCall(requestBodySpec, uri, timeout);
}

然后在我的 ExternalCall 类中

private Mono<String> internalCall(String rq) {
    return requestBodySpec.bodyValue(rq)
            .retrieve()
            .bodyToMono(String.class)
            .timeout(timeout, Mono.error(() -> new TimeoutException(String.format("%s - timeout after %s seconds", "ExternalService", timeout.getSeconds()))));
}

之后:

我们将 WebClient 类传递给 ExternalCall

@Bean
ExternalCall externalCall(WebClient.Builder webClientBuilder) {

    var exchangeStrategies = getExchangeStrategies();

    var timeout = 10000;
    var uri = "https://externalService.com/";

    var webClient = webClientBuilder.clone()
            .clientConnector(new ReactorClientHttpConnector(HttpClient.create()))
            .exchangeStrategies(exchangeStrategies)
            .build();


    return new ExternalCall(webClient, uri, timeout);
}

现在我们在 ExternalCall 类中指定 RequestBodySpec:

private Mono<String> internalCall(String rq) {
    return webClient
            .post()
            .uri(endpoint)
            .accept(TEXT_XML)
            .contentType(TEXT_XML)
            .bodyValue(rq)
            .retrieve()
            .bodyToMono(String.class)
            .timeout(timeout, Mono.error(() -> new TimeoutException(String.format("%s - timeout after %s seconds", "ExternalService", timeout.getSeconds()))));
}

结论:显然,您创建 WebClient.RequestBodySpec 实例的那一刻很重要。希望对某人有所帮助

【讨论】:

    【解决方案2】:

    添加我的一些观察:

    webClient.get()webClient.post() 在每次调用 URI 时总是返回新的 DefaultRequestBodyUriSpec,我认为它看起来不像 URI 被交换。

    class DefaultWebClient implements WebClient {
    
    ..
        @Override
        public RequestHeadersUriSpec<?> get() {
            return methodInternal(HttpMethod.GET);
        }
        @Override
        public RequestBodyUriSpec post() {
            return methodInternal(HttpMethod.POST);
        }
    
    
    ..
    
    
        @Override
        public Mono<ClientResponse> exchange() {
            ClientRequest request = (this.inserter != null ?
                    initRequestBuilder().body(this.inserter).build() :
                    initRequestBuilder().build());
            return exchangeFunction.exchange(request).switchIfEmpty(NO_HTTP_CLIENT_RESPONSE_ERROR);
        }
    
        private ClientRequest.Builder initRequestBuilder() {
            URI uri = (this.uri != null ? this.uri : uriBuilderFactory.expand(""));
            return ClientRequest.create(this.httpMethod, uri)
                    .headers(headers -> headers.addAll(initHeaders()))
                    .cookies(cookies -> cookies.addAll(initCookies()))
                    .attributes(attributes -> attributes.putAll(this.attributes));
        }
    ..
    }
    

    methodInternal 方法如下所示

    @SuppressWarnings("unchecked")
    private RequestBodyUriSpec methodInternal(HttpMethod httpMethod) {
        return new DefaultRequestBodyUriSpec(httpMethod);
    }
    

    另外在发出实际请求的同时,还会创建新的ClientRequest

    类源

    https://github.com/spring-projects/spring-framework/blob/master/spring-webflux/src/main/java/org/springframework/web/reactive/function/client/DefaultWebClient.java

    【讨论】:

      猜你喜欢
      • 2013-09-03
      • 2021-08-19
      • 2019-02-14
      • 1970-01-01
      • 1970-01-01
      • 2018-06-11
      • 2017-09-16
      • 1970-01-01
      • 2018-03-14
      相关资源
      最近更新 更多