【发布时间】: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