You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spring Cloud Stream Kafka消费者WebClient调用返回null问题排查

Kafka消费者中WebClient调用无法更新对象的问题分析

我有两个基于Kafka的Spring Cloud Stream应用,分别作为生产者和消费者。消费者应用接收消息后,会调用数据库和WebClient获取信息并返回更新后的对象。但下方的消费者代码中,WebClient调用无法更新对象,对应字段返回null,而数据库调用正常,且被调用的服务本身运行正常。

@Service
public class BookService {

    @Autowired
    private BookRepository bookRepository;

    @Autowired
    private WebClient webClient;

    public Book updateBook(Book book) {
        book.setBookName(book.getBookName().toUpperCase());
        bookRepository.findById(book.getId()).ifPresent(book1 -> {
            book.setBookIsbn(book1.getBookIsbn()); // 此字段能正常设置
        });

        webClient
                .method(HttpMethod.GET)
                .uri("http://localhost:9898/rx/getBookDetails/" + book.getId())
                .contentType(MediaType.APPLICATION_JSON)
                .accept(MediaType.ALL)
                .retrieve()
                .onStatus(HttpStatus::isError, response -> Mono.error(new RuntimeException("Exception :: " + response)))
                .bodyToMono(Book.class)
                .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)))
                .onErrorResume(error -> Mono.just(new Book()))
                .map(book1 -> {
                    book.setAuthorId(book1.getAuthorId()); // authorId始终为null
                    return book;
                });

        System.out.println("Book processed successfully...");
        return book;
    }
}

核心问题

WebClient是基于Reactor的异步非阻塞工具,你当前的代码只是构建了请求的Mono对象,但没有触发订阅(subscribe),WebClient的请求根本没有实际执行——或者说在请求还没完成、map里的赋值操作还没执行时,就已经提前返回了book对象,所以authorId字段始终是null。

另外,bookRepository.findById是Spring Data JPA的同步阻塞方法,能正常完成赋值,但WebClient的异步逻辑不会等待其执行。


修复方案

方案1:改为异步非阻塞模式(推荐,符合Reactor生态)

把方法返回类型改为Mono<Book>,通过Reactor操作符串联所有逻辑:

@Service
public class BookService {

    @Autowired
    private BookRepository bookRepository;

    @Autowired
    private WebClient webClient;

    public Mono<Book> updateBook(Book book) {
        book.setBookName(book.getBookName().toUpperCase());
        
        // 处理数据库查询(同步方法包装为Mono)
        Mono<Book> dbProcessedBook = Mono.fromCallable(() -> {
            bookRepository.findById(book.getId()).ifPresent(book1 -> {
                book.setBookIsbn(book1.getBookIsbn());
            });
            return book;
        });

        // 串联WebClient请求逻辑
        return dbProcessedBook.flatMap(processedBook -> 
            webClient
                .get() // 简化GET请求写法
                .uri("http://localhost:9898/rx/getBookDetails/{id}", processedBook.getId()) // 用占位符避免字符串拼接
                .accept(MediaType.APPLICATION_JSON)
                .retrieve()
                .onStatus(HttpStatus::isError, response -> Mono.error(new RuntimeException("调用失败: " + response.statusCode())))
                .bodyToMono(Book.class)
                .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)))
                .onErrorResume(error -> Mono.just(new Book()))
                .map(remoteBook -> {
                    processedBook.setAuthorId(remoteBook.getAuthorId());
                    return processedBook;
                })
                .doOnSuccess(updatedBook -> System.out.println("Book processed successfully..."))
        );
    }
}

如果你的bookRepository支持Reactive(比如继承ReactiveCrudRepository),可以直接使用其返回的Mono,无需Mono.fromCallable包装,更贴合非阻塞规范。

方案2:强制同步阻塞(不推荐,仅用于兼容旧代码)

如果必须保持方法返回Book,可以在WebClient的Mono上调用block()触发执行并等待结果:

public Book updateBook(Book book) {
    book.setBookName(book.getBookName().toUpperCase());
    bookRepository.findById(book.getId()).ifPresent(book1 -> {
        book.setBookIsbn(book1.getBookIsbn());
    });

    // 调用block()等待异步请求完成
    Book remoteBook = webClient
            .get()
            .uri("http://localhost:9898/rx/getBookDetails/{id}", book.getId())
            .accept(MediaType.APPLICATION_JSON)
            .retrieve()
            .onStatus(HttpStatus::isError, response -> Mono.error(new RuntimeException("调用失败: " + response.statusCode())))
            .bodyToMono(Book.class)
            .retryWhen(Retry.backoff(3, Duration.ofSeconds(1)))
            .onErrorResume(error -> Mono.just(new Book()))
            .block(); // 阻塞等待结果

    if (remoteBook != null) {
        book.setAuthorId(remoteBook.getAuthorId());
    }

    System.out.println("Book processed successfully...");
    return book;
}

⚠️ 注意:在异步场景(比如Spring Cloud Stream消费者)中使用block()可能导致线程池阻塞,影响性能和扩展性,优先推荐方案1。

内容的提问来源于stack exchange,提问作者astar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 04:26:35