如何使用Kafka Streams DSL处理Transformer的多输出结果
Kafka DSL处理Transformer多输出的实现方法
核心实现思路
你当前在Transformer中通过To.child("xxx")向指定命名子节点转发消息的写法是正确的,只需要在调用transform方法时声明对应命名输出分支,即可拆分得到两个类型不同的独立KStream。
完整代码示例
1. 修正Transformer实现
由于所有输出都通过context.forward手动发送,不需要方法返回值输出,所以Transformer的返回值泛型设置为Void,transform方法直接返回null即可:
public class MyTransformer implements Transformer<String, Transaction, Void> { private ProcessorContext context; @Override public void init(ProcessorContext context) { this.context = context; } @Override public Void transform(String key, Transaction value) { if (/* 触发输出类型A的判断条件 */) { context.forward(key1, new A(...), To.child("first_child")); } else { context.forward(key2, new B(...), To.child("second_child")); } // 所有消息已手动转发,无需返回值 return null; } @Override public void close() { // 按需填写资源清理逻辑 } }
2. 拓扑构建代码
调用transform的重载方法,传入你定义的输出分支名称列表,得到分支结果数组后强转为对应类型的KStream即可:
someMethod(KStream<String, Transaction> transaction) { // 定义和Transformer中forward对应的输出分支名称,顺序和后续数组返回顺序一致 final String[] OUTPUT_BRANCHES = new String[]{"first_child", "second_child"}; // 调用带分支参数的transform重载,得到分支输出数组 final KStream<?, ?>[] outputBranches = transaction .transform(() -> new MyTransformer(...), Named.as("custom-transform"), OUTPUT_BRANCHES); // 按分支顺序强转为对应类型的KStream KStream<String, A> aStream = (KStream<String, A>) outputBranches[0]; KStream<String, B> bStream = (KStream<String, B>) outputBranches[1]; // 后续可分别对两个不同类型的流做独立处理 // aStream.xxx(); // bStream.xxx(); }
注意事项
- 输出分支数组的顺序和你定义
OUTPUT_BRANCHES数组的顺序严格对应,outputBranches[0]对应first_child的输出,outputBranches[1]对应second_child的输出 - 类型强转时要保证Transformer中forward的消息类型和强转目标类型一致,Kafka Streams运行时不会做泛型检查,类型不匹配会在后续处理阶段抛出ClassCastException
- 不要同时使用
context.forward和transform方法返回值输出消息,避免产生多余的无效消息
内容的提问来源于stack exchange,提问作者Fruit
相关产品推荐
相关产品推荐

