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
相关产品推荐
相关产品推荐

