【问题标题】:Spring Cloud Stream content-based routing on payload负载上基于 Spring Cloud Stream 内容的路由
【发布时间】:2020-10-14 10:10:46
【问题描述】:

我们正在使用带有 Kafka 和 Avro(本机编码器/解码器)的 Spring Cloud Stream v2.2。我们正在尝试使用基于负载条件的基于内容的路由。据我了解,根据Spring Cloud Stream docs,基于内容的路由只能在标头上实现,因为有效负载在到达条件时还没有经过类型转换过程。因此,除非条件基于字节格式,否则它不会按预期工作。但是,我知道当在本机模式下使用 Avro 时,会跳过消息头并且不处理类型协商。所以我不确定基于内容的路由是否可以按预期在有效负载上工作。

@StreamListener(target = Channels.INPUT, condition =
      "payload.context['type']=='one' or"
          + " payload.context['type']=='two'")
  public void doStuff(TypeOneAndTwoData inputData){
...
channels.outputChannel().send(MessageBuilder.withPayload(inputData).build());
}


@StreamListener(target = Channels.INPUT, condition =
      "payload.context['type']=='three' or"
          + " payload.context['type']=='four'")
  public void doOtherStuff(TypeThreeAndFourData inputData){
...
channels.outputChannel().send(MessageBuilder.withPayload(inputData).build());
}

根据我已经到位的日志记录,我可以看到偶尔会触发doStuff,有时会触发doOtherStuff。但是,似乎它们中的大多数都没有被触发并且消息被跳过。根据输入数据,我确信context.type 只能有“一”、“二”、“三”和“四”这 4 个值,因此根据输入,不可能期望有其他值,但经常我可以在日志中看到以下条目:

Cannot find a @StreamListener matching for message with id: null

我有几个问题:

  • 在消息被反序列化为原生 Avro 格式的相应 POJO 类后,负载条件是否对消息起作用?
  • 为什么有时条件有效,有时无效? id: null 有什么意义吗?
  • 从线程的角度来看,基于内容的路由如何工作?当我们有两个具有不同条件的 StreamListner 或者它们单线程工作时,是否会运行多个线程?在这种情况下,如何至少一次保证消息传递得到管理?条件是否应该相互排斥?

【问题讨论】:

    标签: spring avro spring-kafka spring-cloud-stream spring-cloud-stream-binder-kafka


    【解决方案1】:
    • 是的;它(当前)受支持。

    • id:null 很不幸;对于 Kafka,message.headers['id'] 默认为空 - 无论如何它没有任何意义,所以这条日志消息没有多大用处。

    • 问题是没有一个条件与转换后的有效负载匹配

    您可以在DispatchingStreamListenerMessageHandler.handleRequestMessage() 中设置断点以找出问题所在。

    编辑

    我没有回答你的第三个问题。

    调用是单线程的;不,一条消息可以匹配多个条件;只有在没有条件匹配时才会得到该日志。请参阅我上面引用的方法。

    如果多个条件匹配监听器抛出的任何异常,将停止处理(并调用重试/DLQ 处理等)。

    【讨论】:

    • 目前,您的意思是 2.2 版也支持它吗?我们的版本有点过时了……
    • 所有支持的(以及下一个 3.1)版本都支持它。但是@StreamListener 和朋友们deprecated in that version 支持更新的Function<?>al 模型;我认为该模型中没有 condition 的等价物,但我可能是错的。
    • 我没有回答你的第三个原始问题;已更新。
    • 谢谢。这是否意味着非原生 Avro 和其他格式(可能文档已过时)也支持有效载荷上的条件,或者不支持,只是因为我们使用带有原生编码/解码的 Avro?
    • 只有转换后的对象才能引用payload;但是,您可以在 JSON 有效负载的表达式中使用 JsonPath - Spring Integration(由 SCSt 使用)注册自定义 SpEL 函数#jsonPath。它适用于各种有效负载,包括byte[]docs.spring.io/spring-integration/docs/5.3.2.RELEASE/reference/…
    猜你喜欢
    • 2021-12-06
    • 1970-01-01
    • 2019-12-11
    • 2018-06-23
    • 2022-11-04
    • 2023-03-19
    • 2017-03-28
    • 2017-12-25
    • 1970-01-01
    相关资源
    最近更新 更多