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

Spring Cloud Kafka Binder消息异常:仅processOrder执行问题排查

问题描述

我尝试通过REST端点调用processOrder()生成消息,期望将processOrder()的结果传递给processShipping()和processPayment()。但每当调用REST端点http://localhost:8080/processOrder时,仅processOrder()被执行,其余方法未触发。相关代码、配置及依赖如下:

Java函数代码

package com.example.kafkademo.functions;

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

import java.util.function.Consumer;
import java.util.function.Function;

@Configuration
public class MessageFunctions {

    @Bean
    public Function<String, String> processOrder(){
        return orderId -> {
            System.out.println("processOrder: " + orderId);
            System.out.println(orderId);
            return orderId + " : " + System.currentTimeMillis();
        };
    }

    @Bean
    public Consumer<String> processShipping(){
        return orderId -> {
            System.out.println("processShipping: " + orderId);
            System.out.println(orderId);
        };
    }

    @Bean
    public Consumer<String> processPayment(){
        return orderId -> {
            System.out.println("processPayment: " + orderId);
            System.out.println(orderId);
        };
    }
}

application.yml配置

spring:
  application:
      name: kafka-demo

  cloud:
    function:
      definition: processOrder;processPayment;processShipping
    stream:
      bindings:
        processOrder-out-0:
          destination: order_topic
        processPayment-in-0:
          destination: order_topic
        processShipping-in-0:
          destination: order_topic

  kafka:
    listener:
      port: 9094
    bootstrap-servers:
      - localhost:9094

Gradle依赖配置

dependencies {
    implementation 'org.springframework.boot:spring-boot-starter-web'
    implementation 'org.springframework.cloud:spring-cloud-stream'
    implementation 'org.springframework.cloud:spring-cloud-stream-binder-kafka'
    implementation 'org.springframework.kafka:spring-kafka'
    implementation 'org.springframework.cloud:spring-cloud-starter-function-web' 
    testImplementation 'org.springframework.boot:spring-boot-starter-test'
    testImplementation 'org.springframework.cloud:spring-cloud-stream-test-binder'
    testImplementation 'org.springframework.kafka:spring-kafka-test'
    testRuntimeOnly 'org.junit.platform:junit-platform-launcher'
}
解决方案

问题根源

直接通过Web端点调用processOrder时,Spring Cloud Function默认只会同步执行函数并返回结果,不会自动将输出发送到Spring Cloud Stream绑定的Kafka主题,导致下游的processPayment和processShipping没有消息可消费。

修复步骤

1. 修正配置并启用函数桥接

更新application.yml,删除无效的spring.kafka.listener.port配置(该配置不存在,Kafka监听器端口由Broker决定),并开启Stream的函数桥接功能,让Web调用的函数输出自动流转到绑定的Kafka主题:

spring:
  application:
      name: kafka-demo

  cloud:
    function:
      definition: processOrder;processPayment;processShipping
    stream:
      function:
        bridge:
          enabled: true # 开启Web函数到Stream绑定的桥接
      bindings:
        processOrder-out-0:
          destination: order_topic
        processPayment-in-0:
          destination: order_topic
        processShipping-in-0:
          destination: order_topic

  kafka:
    bootstrap-servers:
      - localhost:9094

2. 确保Kafka主题可用

  • 手动创建order_topic:执行Kafka命令kafka-topics.sh --create --topic order_topic --bootstrap-server localhost:9094
  • 或者在Kafka的server.properties中开启自动创建主题:auto.create.topics.enable=true

3. 测试验证

重启应用后,向http://localhost:8080/processOrder发送POST请求,请求体传入字符串(例如"ORD123")。此时:

  1. processOrder执行并返回结果
  2. 结果会被自动发送到order_topic
  3. processPayment和processShipping会消费order_topic中的消息并执行

4. 版本兼容性检查

确保Spring Boot与Spring Cloud版本匹配,避免因版本差异导致绑定失效(例如Spring Boot 2.7.x对应Spring Cloud 2021.0.x)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 16:09:53