【问题标题】:Kafka Stream and KGlobalTable Join issueKafka Stream 和 KGlobalTable Join 问题
【发布时间】:2019-06-07 06:21:57
【问题描述】:

我在使用 GlobalKTable 加入 KStream 时遇到问题,希望您能提供帮助。

给定两个Kafka主题orderscustomers

订单

"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"}}

我正在尝试以下方法:

  1. orders 主题构建 KStream
  2. customers 主题构建 GlobalKTable
  3. 构建另一个连接订单和客户的流(在客户表中查找 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


【解决方案1】:

让我们仔细看看你复制粘贴的内容:

customers话题中:

"100"   {"CustID":"100","CustName":"Customer1"}

您可以注意到键是一个字符串,并且此字符串包含双引号"100"。通常,字符串键的打印不带双引号。我宁愿期待看到:

 100    {"CustID":"100","CustName":"Customer1"}

换句话说,您的密钥的 Java 字符串表示是 ""100""(或 "\"100\""),而不是我们预期的 "100"

另一方面,orders 主题中的值是 Json {"ID":"1","Name":"Myorder1","CustID":"100"},属性 CustID 是字符串,这次用 Java 表示 "100"

当您加入 orderscustomers 时,您会尝试将订单 CustID 100 与客户键 "100" 匹配。由于 CustID 中缺少密钥中的双引号,这将失败。

【讨论】:

    【解决方案2】:

    @deepak 你可能需要实现你的 KTable

    builder.table(customers, Materialized.as(customerStore));
    

    然后流式传输订单并建立您的加入。

    【讨论】:

    • 我正在使用 .GlobalKTable customers = builder.globalTable(Customer",Consumed.with(Serdes.String(), customerSerde),Materialized.>as("my-state-store")) . 请建议以防我们需要不同的使用它。
    • 确保您确实需要一个 GlobalKTable,如果您要将 Orders 主题中使用的键切换为“CustID”,您可以改为使用 KTable(这会提供更好的性能)。如果两个主题具有相同的分区策略并且具有相同的键(在您的情况下,您可以使用“CustId”),那么您可以使用 Ktable 与 GlobalKTable。请参考 Matthias J. Sax 的回答,stackoverflow.com/questions/45975755/…
    猜你喜欢
    • 2023-02-23
    • 1970-01-01
    • 2019-10-05
    • 2021-06-07
    • 1970-01-01
    • 1970-01-01
    • 2023-01-30
    • 1970-01-01
    • 2021-06-09
    相关资源
    最近更新 更多