【问题标题】:route FROM and route TO with spring cloud stream and functions使用 spring cloud 流和函数的 route FROM 和 route TO
【发布时间】:2020-02-27 13:45:35
【问题描述】:

我对 Spring Cloud Stream 中的新路由功能有一些问题

我尝试实现一个简单的场景,我想发送一个带有header spring.cloud.function.definition = consume1 or consume2的消息

我希望应根据标头上发送的内容调用consume1或consume2,但方法是随机调用的。

我使用 rabbit 管理控制台将消息发送给交换消费者

我有以下日志:

2020-02-27 14:48:25.896  INFO 22132 --- [ consumer.app-1] com.example.demo.TestConsumer            : ==============>consume1 messge [[payload=ok, headers={amqp_receivedDeliveryMode=NON_PERSISTENT, amqp_receivedRoutingKey=#, amqp_receivedExchange=consumer, amqp_deliveryTag=1, deliveryAttempt=1, amqp_consumerQueue=consumer.app, amqp_redelivered=false, id=9a4dff25-88ef-4d76-93e2-c8719cda122d, spring.cloud.function.definition=consume1, amqp_consumerTag=amq.ctag-gGChFNCKIVd25yyR9H6-fQ, sourceData=(Body:'[B@3a92faa7(byte[2])' MessageProperties [headers={spring.cloud.function.definition=consume1}, contentLength=0, receivedDeliveryMode=NON_PERSISTENT, redelivered=false, receivedExchange=consumer, receivedRoutingKey=#, deliveryTag=1, consumerTag=amq.ctag-gGChFNCKIVd25yyR9H6-fQ, consumerQueue=consumer.app]), timestamp=1582811303347}]]
2020-02-27 14:48:25.984  INFO 22132 --- [nio-8080-exec-1] o.a.c.c.C.[Tomcat].[localhost].[/]       : Initializing Spring DispatcherServlet 'dispatcherServlet'
2020-02-27 14:48:25.984  INFO 22132 --- [nio-8080-exec-1] o.s.web.servlet.DispatcherServlet        : Initializing Servlet 'dispatcherServlet'
2020-02-27 14:48:25.991  INFO 22132 --- [nio-8080-exec-1] o.s.web.servlet.DispatcherServlet        : Completed initialization in 7 ms
2020-02-27 14:48:26.037  INFO 22132 --- [oundedElastic-1] o.s.i.monitor.IntegrationMBeanExporter   : Registering MessageChannel customer-1
2020-02-27 14:48:26.111  INFO 22132 --- [oundedElastic-1] o.s.c.s.m.DirectWithAttributesChannel    : Channel 'application.customer-1' has 1 subscriber(s).
2020-02-27 14:48:26.116  INFO 22132 --- [oundedElastic-1] o.s.a.r.c.CachingConnectionFactory       : Attempting to connect to: [localhost:5672]
2020-02-27 14:48:26.123  INFO 22132 --- [oundedElastic-1] o.s.a.r.c.CachingConnectionFactory       : Created new connection: rabbitConnectionFactory.publisher#32438e24:0/SimpleConnection@3e58666d [delegate=amqp://guest@127.0.0.1:5672/, localPort= 62514]
2020-02-27 14:48:26.139  INFO 22132 --- [-1.customer-1-1] o.s.i.h.s.MessagingMethodInvokerHelper   : Overriding default instance of MessageHandlerMethodFactory with provided one.
2020-02-27 14:48:26.140  INFO 22132 --- [-1.customer-1-1] com.example.demo.TestSink                : Data received customer-1...body
2020-02-27 14:49:14.185  INFO 22132 --- [ consumer.app-1] o.s.i.h.s.MessagingMethodInvokerHelper   : Overriding default instance of MessageHandlerMethodFactory with provided one.
2020-02-27 14:49:14.194  INFO 22132 --- [ consumer.app-1] com.example.demo.TestConsumer            : ==============>consume2 messge [[payload=ok, headers={amqp_receivedDeliveryMode=NON_PERSISTENT, amqp_receivedRoutingKey=#, amqp_receivedExchange=consumer, amqp_deliveryTag=1, deliveryAttempt=1, amqp_consumerQueue=consumer.app, amqp_redelivered=false, id=33581edb-2832-1c92-b765-a05794512b34, spring.cloud.function.definition=consume1, amqp_consumerTag=amq.ctag-RIp2nZdcG2a0hNQeImwtBw, sourceData=(Body:'[B@8159793(byte[2])' MessageProperties [headers={spring.cloud.function.definition=consume1}, contentLength=0, receivedDeliveryMode=NON_PERSISTENT, redelivered=false, receivedExchange=consumer, receivedRoutingKey=#, deliveryTag=1, consumerTag=amq.ctag-RIp2nZdcG2a0hNQeImwtBw, consumerQueue=consumer.app]), timestamp=1582811354186}]]
2020-02-27 14:49:14.203  INFO 22132 --- [oundedElastic-1] o.s.i.monitor.IntegrationMBeanExporter   : Registering MessageChannel customer-2
2020-02-27 14:49:14.213  INFO 22132 --- [oundedElastic-1] o.s.c.s.m.DirectWithAttributesChannel    : Channel 'application.customer-2' has 1 subscriber(s).
2020-02-27 14:49:14.216  INFO 22132 --- [-2.customer-2-1] o.s.i.h.s.MessagingMethodInvokerHelper   : Overriding default instance of MessageHandlerMethodFactory with provided one.
2020-02-27 14:49:14.216  INFO 22132 --- [-2.customer-2-1] com.example.demo.TestSink                : Data received customer-2...body

application.yml

spring:
  main:
    allow-bean-definition-overriding: true
spring.cloud.stream:
  function.definition: supplier;receive1;receive2;consume1;consume2
  function.routing:
    enabled: true

  bindings:
    consume1-in-0.destination: consumer
    consume1-in-0.group: app
    consume2-in-0.destination: consumer
    consume2-in-0.group: app
    receive1-in-0.destination: customer-1
    receive1-in-0.group: customer-1
    receive2-in-0.destination: customer-2
    receive2-in-0.group: customer-2

DemoApplication.java

import com.fasterxml.jackson.databind.ObjectMapper
import org.apache.commons.logging.Log
import org.apache.commons.logging.LogFactory
import org.springframework.boot.autoconfigure.SpringBootApplication
import org.springframework.boot.runApplication
import org.springframework.context.annotation.Bean
import org.springframework.http.HttpStatus
import org.springframework.messaging.Message
import org.springframework.messaging.support.MessageBuilder
import org.springframework.stereotype.Component
import org.springframework.web.bind.annotation.PathVariable
import org.springframework.web.bind.annotation.RequestMapping
import org.springframework.web.bind.annotation.RequestMethod.GET
import org.springframework.web.bind.annotation.ResponseStatus
import org.springframework.web.bind.annotation.RestController
import org.springframework.web.client.RestTemplate
import reactor.core.publisher.EmitterProcessor
import reactor.core.publisher.Flux
import java.util.function.Consumer
import java.util.function.Supplier


@SpringBootApplication
class DemoApplication

fun main(args: Array<String>) {
    runApplication<DemoApplication>(*args)
}

@RestController
class DynamicDestinationController(private val jsonMapper: ObjectMapper) {

    private val processor: EmitterProcessor<Message<String>> = EmitterProcessor.create<Message<String>>()

    @RequestMapping(path = ["/api/dest/{destName}"], method = [GET], consumes = ["*/*"])
    @ResponseStatus(HttpStatus.ACCEPTED)
    fun handleRequest(@PathVariable destName:String) {
        val message: Message<String> = MessageBuilder.withPayload("body")
                .setHeader("spring.cloud.stream.sendto.destination", destName).build()
        processor.onNext(message)
    }

    @Bean
    fun supplier(): Supplier<Flux<Message<String>>> {
        return Supplier { processor }
    }
}

const val destResourceUrl = "http://localhost:8080/api/dest"
@Component
class TestConsumer() {

    private val restTemplate: RestTemplate = RestTemplate()
    private val logger: Log = LogFactory.getLog(javaClass)

    @Bean
    fun consume1(): Consumer<Message<String>> = Consumer {
        logger.info("==============>consume1 messge [[payload=${it.payload}, headers=${it.headers}]]")
        restTemplate.getForEntity("$destResourceUrl/customer-1", String::class.java)
    }

    @Bean
    fun consume2(): Consumer<Message<String>> = Consumer {
        logger.info("==============>consume2 messge [[payload=${it.payload}, headers=${it.headers}]]")
        restTemplate.getForEntity("$destResourceUrl/customer-2", String::class.java)
    }
}


@Component
class TestSink {
    private val logger: Log = LogFactory.getLog(javaClass)
    @Bean
    fun receive1(): Consumer<String> = Consumer {
        logger.info("Data received customer-1..." + it);
    }

    @Bean
    fun receive2(): Consumer<String> = Consumer {
        logger.info("Data received customer-2..." + it);
    }
}

知道如何修复通往消费者的路线吗?

提前致谢。

demo-repo

【问题讨论】:

  • 为什么在一个应用程序中拥有源、处理器和接收器?我的意思是它只是为了测试而打算将其分解?
  • 这只是一个测试,但在一个真实的项目中,我需要在同一个应用程序中拥有多个功能,我正在研究迁移到新的功能范式。在单个应用程序中的当前项目上,我有多个绑定输入和输出,输入条件...
  • 我问的原因是它目前不适用于您的处理方式。 EmitterProcessor 绑定到单个供应商,因此数据总是在那里等等。 . .好消息是,很快我们将合并一个允许您这样做的问题。基本上,一旦this 被合并,我会在这里用完整的例子跟进。
  • 好的,谢谢你的回答。
  • 我不确定我的问题是 EmitterProcessor,你能看看这个简化的应用程序github.com/iguissouma/cloud-stream-functions-sample 如果你能确认它是否是同样的问题,我在队列 rabbit ha_consumer_queue 中发布了一条消息一个标头 spring.cloud.function.definition=simpleConsumer1 但执行的方法并不总是 simpleConsumer1?提前致谢

标签: spring-boot kotlin spring-cloud-stream spring-cloud-function


【解决方案1】:

其实我有点迷茫,所以我们一步一步来。这是使用sendto 功能的功能(仿照您的)应用程序,允许您将消息发送到特定(现有和/或动态解析的)目的地。

(在 java 中,但您可以将其改写为 Kotlin)

@Controller
public class WebSourceApplication {

    public static void main(String[] args) {
        SpringApplication.run(WebSourceApplication.class,
                "--spring.cloud.function.definition=supplier;consA;consB",
                "--spring.cloud.stream.bindings.consA-in-0.destination=consumerA",
                "--spring.cloud.stream.bindings.consA-in-0.group=consumerA-grp",
                "--spring.cloud.stream.bindings.consB-in-0.destination=consumerB",
                "--spring.cloud.stream.bindings.consB-in-0.group=consumerB-grp"
                );
    }

    EmitterProcessor<Message<String>> processor = EmitterProcessor.create();

    @RequestMapping(path = "/api/dest/{destName}", consumes = "*/*")
    @ResponseStatus(HttpStatus.ACCEPTED)
    public void delegateToSupplier(@RequestBody String body, @PathVariable String destName) {
        Message<String>  message = MessageBuilder.withPayload(body)
            .setHeader("spring.cloud.stream.sendto.destination", destName)
            .build();
        processor.onNext(message);
    }

    @Bean
    public Supplier<Flux<Message<String>>> supplier() {
        return () -> processor;
    }

    @Bean
    public Consumer<String> consA() {
        return v -> {
            System.out.println("Consuming from consA:  " + v);
        };
    }

    @Bean
    public Consumer<String> consB() {
        return v -> {
            System.out.println("Consuming from consB:  " + v);
        };
    }
}

当我卷曲它时,我会得到适当的消费者一致的调用:

curl -H "Content-Type: application/json" -X POST -d "Hello Spring Cloud Stream" http://localhost:8080/api/dest/consumerA
log: Consuming from consA:  Hello Spring Cloud Stream
. . .

curl -H "Content-Type: application/json" -X POST -d "Hello Spring Cloud Stream" http://localhost:8080/api/dest/consumerB
log: Consuming from consB:  Hello Spring Cloud Stream

注意:没有启用路由属性。该功能的主要目的是始终调用一个函数functionRouter,并让它代表您调用其他函数。它是 spring-cloud-function 的一个特性,这意味着它在 spring-cloud-srteam 和通道/目的地等之外工作。

这不是你想要完成的吗?根据 HTTP 请求中的某些 oath 变量将消息发送到不同的目的地?

【讨论】:

  • 发送到工作正常,我想要实现的就像旧样式的一个输入,带有两个侦听器和条件,但我不知道到消费者的路由是如何工作的
  • 所以我添加了TestConsumer 两个函数consume1 和cosume2,另一个微服务如何调用它们?使用旧样式,我向该输入发送一条消息,并执行条件验证的方法。如何实现相同的行为?
  • 您的问题是您在一个应用程序中测试所有内容,这违反了框架的设计并产生了某些冲突。例如,您不能真正拥有 spring.cloud.function.definitionspring.cloud.stream.function.routing.enabled 属性。所以让我发布一个不同的例子。但这就是为什么我最初问为什么你在一个应用程序中拥有所有这些?
【解决方案2】:

这是一个不同的微服务示例,它接收路由功能,然后路由到不同的功能

public class FunctionRoutingApplication {

    public static void main(String[] args) {
        SpringApplication.run(FunctionRoutingApplication.class,
                "--spring.cloud.stream.function.routing.enabled=true"
                );
    }

    @Bean
    public Consumer<String> consA() {
        return v -> {
            System.out.println("Consuming from consA:  " + v);
        };
    }

    @Bean
    public Consumer<String> consB() {
        return v -> {
            System.out.println("Consuming from consB:  " + v);
        };
    }
}

差不多就是这样。转到您的代理并将数据发送到 functionRouter-in-0 交换,同时提供 spring.cloud.function.definition=consA/consB 标头,您将看到一致的调用。

我还是错过了什么吗?

【讨论】:

  • functionRouter-in-0交换应该是自动创建还是手动创建?它可以被其他微服务共享吗?我最初的要求是创建一个微服务来安排将消息发送到其他微服务(这是动态目的地的需要),微服务可以安排/重新安排/取消安排消息的发送,所以我想要一个输入通道并基于方法被调用的路由表达式
  • 名称functionRouterconstant,如果不是overriden-in-0 是它的默认绑定名称和目标。与任何其他目的地一样,目的地(交换、队列、主题等)要么被创建(如果不存在),要么我们使用现有的。
  • 另外,这可能会有所帮助。函数的输出可以通过spring.cloud.stream.sendto.destination 属性路由到目的地。另一方面,到达目的地的输入可以通过 RoutingFunction 根据此处cloud.spring.io/spring-cloud-static/spring-cloud-stream/… 中描述的规则路由到特定函数。
  • 感谢您花时间解释这一切。
  • 如果我收到两种不同类型的消息,它们具有不同的标头字段和值,该怎么办?例如我的消费服务收到以下消息(在“订单”队列和“发货”队列上):Order1: { Header: {catalog:groceries} } , Order2: { Header: {catalog:tools} }Shipment1: { Header: {region:Europe} }, Shipment2: { Header: {region:America} } 我将为 2 个队列提供 2 个绑定,但我无法将每条消息路由到他们自己的消费者函数,因为我只能为headers['catalog']headers['region'] 全局设置路由条件。
猜你喜欢
  • 2022-12-01
  • 2020-02-01
  • 1970-01-01
  • 2022-11-20
  • 2019-06-14
  • 2023-02-06
  • 2015-05-19
  • 2020-01-27
  • 2017-12-25
相关资源
最近更新 更多