【发布时间】:2019-06-25 06:41:35
【问题描述】:
在我的基于 DSL 的转换中,我有一个流--> 分支,我希望将分支输出重定向到多个主题。
当前的branch.to() 方法只接受String。
stream.branch 是否有任何简单的选项,我可以将结果路由到多个主题。对于消费者,我可以通过提供一个字符串数组作为主题来订阅多个主题。
如果特定谓词满足查询,我的问题需要我采取多项操作。
我尝试使用stream.branch[index].to(string),但这不足以满足我的要求。我正在寻找类似stream.branch[index].to(string array of topics) 或stream.branch[index].to(string) 的东西。
我希望 branch.to 方法具有多个主题,或者是否有其他方法可以通过流实现相同的目标?
添加示例代码。删除实际变量名称。
我的谓词
Predicate <String, MyDomainObject> Predicate1 = new Predicate<String, MyDomainObject>() {
@Override
public boolean test(String key, MyDomainObject domObj) {
boolean result = false;
if condition on domObj
return result;
}
};
Predicate <String, MyDomainObject> Predicate2 = new Predicate<String, MyDomainObject>() {
@Override
public boolean test(String key, MyDomainObject domObj) {
boolean result = false;
if condition on domObj
return result;
}
};
KStream <String, MyDomainObject>[] branches= myStream.branch(
Predicate1, Predicate2
);
// here I need your suggestions.
// this is my current implementation
branches[0].to(singleTopic),
Produced.with(Serdes.String(), Serdes.serdeFrom(inSer, deSer)));
// I want to send notification to multiple topics. something like below
branches[0].to(topicList),
Produced.with(Serdes.String(), Serdes.serdeFrom(inSer, deSer)));
【问题讨论】:
-
欢迎来到 SO。如果您包含一些您在其上下文中尝试过的代码,将会很有帮助。我知道您在段落中包含了一些内容,但是如果您在代码块中看到您的整个函数,那就太好了。
-
我已经用我正在使用的代码更新了我的原始帖子。
标签: java apache-kafka apache-kafka-streams