【发布时间】:2019-06-07 06:21:57
【问题描述】:
我在使用 GlobalKTable 加入 KStream 时遇到问题,希望您能提供帮助。
给定两个Kafka主题orders和customers:
订单
"1" {"ID":"1","Name":"Myorder1","CustID":"100"}
"2" {"ID":"2","Name":"MyOrder2","CustID":"200"}
客户
"100" {"CustID":"100","CustName":"Customer1"}
"200" {"CustID":"200","CustName":"Customer2"}
要求是用客户名称丰富订单流
"1" {"ID":"1","Name":"Myorder1","CustID":"100","CustName":"Customer1"}
"2" {"ID":"2","Name":"MyOrder2","CustID":"200","CustName":"Customer2"}}
我正在尝试以下方法:
- 从
orders主题构建 KStream - 从
customers主题构建 GlobalKTable - 构建另一个连接订单和客户的流(在客户表中查找 Order.CustID)
KStream<String, EnrichedOrder> enrichedstreams = orders.join(
customers,
new KeyValueMapper<String, Order, String>() {
@Override
public String apply(String key, Order value) {
return value.CustID;
}
},
new ValueJoiner<Order,Customer, EnrichedOrder>() {
@Override
public EnrichedOrder apply(Order order, Customer customer) {
EnrichedOrder eorder = new EnrichedOrder();
eorder.CustID = order.CustID;
eorder.CustName = customer.CustName;
eorder.ID = order.ID;
eorder.Name = order.Name;
return eorder;
}
}
);
但它没有给出任何结果,也没有抛出任何异常。
当使用 leftJoin 时,我得到了客户的 NullPointer 异常。
如果您遇到类似问题,请告诉我并提出解决方法。
【问题讨论】:
-
确保在流处理开始时
GlobalKTable已完全填充。可能是订单已经在处理,而客户表仍在填充。为避免这种情况,请启动流应用程序,然后才生成新的订单事件。此外,您可能需要在多次运行测试时重置偏移量。 -
有什么方法可以检查 GlobalKtable 是否已完全填充?就像 kstream 中的 foreach 循环一样。
-
@deepak,您能否提供创建 GlobalKTable 的代码?
-
@deepak,确保您的消息(订单和客户)有密钥。顺便问一下,你真的需要 GlobalKTable(而不是 KTable)吗?
-
@user152468 -- 您最初对“加载”表的担忧不应该适用 -- 在启动时,
GlobalKTable在任何处理开始之前被引导到主题的末尾(这是一个区别到KTables,提供时间同步的加入/处理,而GlobalKTables 不是时间同步的)。
标签: java apache-kafka apache-kafka-streams