【问题标题】:Spring cloud stream routing on payload有效负载上的 Spring 云流路由
【发布时间】:2022-01-12 03:34:24
【问题描述】:

我想为我的微服务使用 Spring Cloud Stream 来处理来自 kafka 的事件。

我从一个可以容纳多个 JSON 有效负载的主题中读取数据(我有一个主题,因为它到达的所有消息都来自同一个主题)。 根据不同的payload,我有不同的云功能要处理。

如何根据负载中的属性将传入事件路由到特定函数?

假设我有可以具有以下属性的 JSON 消息:

{
  "type":"A"
  "content": xyz
}

所以输入消息可以有一个属性A或B

假设我想在类型为 A 时调用一些 bean 函数,在类型为 B 时调用另一个 bean 函数

【问题讨论】:

  • 你试过写Function<KStream的方法,然后调用KStream.branch的方法吗? docs.spring.io/spring-cloud-stream-binder-kafka/docs/3.1.0.M1/…
  • 我不需要写发送到不同的主题,我想根据其有效负载中的属性调用函数(处理来自 kafka 的事件)
  • 那么您是否尝试过使用.filter() 对该属性进行检查,然后使用.map() 函数将这些过滤后的事件传递给您要运行的函数?
  • 我明白了...谢谢。但是假设我有 3 或 4 个可能的值,我该如何根据它处理或路由?
  • 我假设像这样的KStream kstream1 = input.filter((k, v) -> v.prop == 1).map(...); KStream kstream2 = input.filter(k, v), v.prop == 2).map(...); ...“路由”意味着您有新主题的传出记录,这意味着您使用.branch,并且您可以拥有这些主题的消费者而是调用您的其他处理方法

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


【解决方案1】:

从问题中不清楚您是使用基于消息通道的 Kafka binder 还是 Kafka Streams binder。上面的 cmets 暗示了对 KStream 的一些引用。假设您使用的是基于消息通道的 Kafka binder,您可以选择使用 Spring Cloud Stream 中的消息路由功能。文档的这一部分解释了基本用法:https://docs.spring.io/spring-cloud-stream/docs/3.2.1/reference/html/spring-cloud-stream.html#_event_routing

您可以提供 routing-expression,这是一个 SpEL 表达式,用于传递正确的属性值。

如果您想要超越 SpEL 表达式所表达的高级路由功能,您还可以实现自定义 MessageRoutingCallback。有关详细信息,请参阅此示例应用程序:https://github.com/spring-cloud/spring-cloud-stream-samples/tree/main/routing-samples/message-routing-callback

【讨论】:

  • 如何将路由功能绑定到特定的主题名称?还是所有主题的所有流量都经过functionRouter?
  • 如果您希望路由功能从特定主题消费,请执行以下操作:spring.cloud.stream.bindings.functionRouter-in-0.destination: <your-topic-name>
  • 谢谢。顺便说一句,我按照您的示例进行操作,似乎在实现 MessageRoutingCallback 时,我在 Message 中将字节 [] 作为有效负载。知道为什么吗?
  • 不知道为什么,可能是默认的序列化程序。与示例应用程序相比,它在相应的函数中转换为正确的类型。您可能需要调试它。如果您可以提供MCRE,那么我们也许可以查看并排除故障。
猜你喜欢
  • 2021-12-02
  • 2019-01-31
  • 1970-01-01
  • 2019-07-26
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多