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

Spring Cloud Data Flow:多输入多输出示例应用需求问询

Spring Cloud Data Flow 多输入多输出示例应用

核心场景

这个示例实现一个订单处理处理器,同时监听订单消息和库存通知两个输入源,处理后分别输出已完成订单和库存预警两个消息流,完全适配SCDF的流编排能力。


1. 自定义处理器代码

创建一个Spring Cloud Stream处理器,通过@EnableBinding定义多输入输出通道:

import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.messaging.handler.annotation.SendTo;

// 定义多输入多输出通道的绑定接口
interface OrderProcessingChannels {
    String ORDER_INPUT = "order-input";
    String INVENTORY_INPUT = "inventory-input";
    String COMPLETED_ORDER_OUTPUT = "completed-order-output";
    String INVENTORY_ALERT_OUTPUT = "inventory-alert-output";
}

@EnableBinding(OrderProcessingChannels.class)
public class OrderProcessingProcessor {

    // 监听订单输入通道,处理后发送到已完成订单输出通道
    @StreamListener(OrderProcessingChannels.ORDER_INPUT)
    @SendTo(OrderProcessingChannels.COMPLETED_ORDER_OUTPUT)
    public String processOrder(String orderMessage) {
        // 模拟订单校验、用户关联等业务逻辑
        return "COMPLETED: " + orderMessage;
    }

    // 监听库存输入通道,处理后发送到库存预警输出通道
    @StreamListener(OrderProcessingChannels.INVENTORY_INPUT)
    @SendTo(OrderProcessingChannels.INVENTORY_ALERT_OUTPUT)
    public String processInventory(String inventoryMessage) {
        // 模拟库存阈值判断逻辑
        if (inventoryMessage.contains("low")) {
            return "ALERT: " + inventoryMessage;
        }
        return "OK: " + inventoryMessage;
    }
}

2. 依赖配置(pom.xml)

添加Spring Cloud Stream和基础依赖:

<dependencies>
    <dependency>
        <groupId>org.springframework.cloud</groupId>
        <artifactId>spring-cloud-starter-stream-rabbit</artifactId>
    </dependency>
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>
</dependencies>

3. SCDF流编排与部署

假设使用RabbitMQ作为消息中间件,按以下步骤操作:

3.1 注册自定义处理器到SCDF

app register --name order-processor --type processor --uri maven://com.your.package:order-processor:1.0.0

3.2 创建并部署两个流

分别绑定处理器的两个输入通道到不同的HTTP源,输出通道到log sink:

# 订单流:HTTP源(8081端口)→ 处理器order-input通道 → 已完成订单日志
stream create --name order-flow --definition "http --server.port=8081 | order-processor --spring.cloud.stream.bindings.order-input.destination=order-input | log --name=completed-order-log"

# 库存流:HTTP源(8082端口)→ 处理器inventory-input通道 → 库存预警日志
stream create --name inventory-flow --definition "http --server.port=8082 | order-processor --spring.cloud.stream.bindings.inventory-input.destination=inventory-input | log --name=inventory-alert-log"

# 部署两个流
stream deploy --name order-flow
stream deploy --name inventory-flow

4. 测试验证

  • 发送订单消息:

    curl -X POST http://localhost:8081 -d "Order #123: User ID 456"
    

    在SCDF的log sink中会看到输出:COMPLETED: Order #123: User ID 456

  • 发送库存预警消息:

    curl -X POST http://localhost:8082 -d "Product A: stock low"
    

    在SCDF的log sink中会看到输出:ALERT: Product A: stock low

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 06:16:07