【问题标题】:Kafka Stream DSL non-key join current workaround explainedKafka Stream DSL 非密钥加入当前解决方法解释
【发布时间】:2019-11-18 05:29:15
【问题描述】:

我正在尝试理解以下提到的解决方法:

https://issues.apache.org/jira/browse/KAFKA-3705

如今,在 Kafka Streams DSL 中,KTable 连接仅基于键。如果 用户想通过 key a 加入一个 KTable A 和另一个 KTable B 通过 key b 但使用“外键”a,并假设它们是从两个主题中读取的 它们分别在 a 和 b 上分区,它们需要执行 以下模式:

tableB' = tableB.groupBy(/* select on field "a" */).agg(...); // now tableB' is partitioned on "a"

tableA.join(tableB', joiner);

我很难理解到底发生了什么。

这句话特别令人困惑:“如果用户想通过键 a 加入另一个 KTable B,但通过键 b 但使用“外键”a”。我也不明白上面的代码。

有人可以澄清一下这里发生了什么吗?

这里也提到了:

缩小流中KTables的语义和关系数据库中的表之间的差距。通常的做法是在对 RDBMS 中的表进行更改时将更改捕获到 Kafka 主题(JDBC-connect、Debezium、Maxwell)中。这些实体通常具有多个一对多的关系。通常,RDBMS 提供了很好的支持来通过连接解决这种关系。 Streams 在这里不足,解决方法 (group by - join -lateral view) 也没有得到很好的支持,并且不符合基于记录的处理的想法。 https://cwiki.apache.org/confluence/display/KAFKA/KIP-213+Support+non-key+joining+in+KTable

什么意思(分组-加入-横向视图)?我怀疑它与上面的代码有关,但又有点难以理解。任何人都可以对此有所了解吗?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    嗯,下面的代码是用非键连接连接两个 KTable 的伪代码:

    tableB' = tableB.groupBy(/* select on field "a" */).agg(...); // now tableB' is partitioned on "a"
    
    tableA.join(tableB', joiner);
    

    解释

    假设 tableA 有一个关键字段“a”。为了与 tableA 加入另一个 ktable,它应该是共同分区的。它应该有相同的键。因此,我们将使用字段“a

    重新设置 ktable tableB
    tableB' = tableB.groupBy(/* select on field "a" */).agg(...); // now tableB' is partitioned on "a"
    

    groupBy()selectKey()+ groupByKey() 操作的简写。

    groupBy(/* select on field "a" */) 将在字段 "a" 上重新键入 tableB 并按该键分组。因此,现在您有一个 KGroupedTable,其中包含字段“a”作为键。为了得到 KTable,你需要调用 .aggregate() 。这就是上面代码中发生的事情。

    附言.agg() 应重命名为 .aggregate()

    tableB' 准备就绪后,您可以使用以下代码加入 tableA

    tableA.join(tableB', joiner);
    

    这里的joiner指的是ValueJoiner的实现。

    示例

    // Java 8+ example, using lambda expressions
    KTable<String, String> joined = left.join(right, 
         /* Below line is ValueJoiner */
        (leftValue, rightValue) -> "left=" + leftValue + ", right=" + rightValue 
      );
    

    目前,这是 KTables 上的非键连接方式 您可以在文档中找到很好的解释:https://docs.confluent.io/current/streams/developer-guide/dsl-api.html#ktable-ktable-join

    【讨论】:

    • 这个可以直接用KSQL表达吗?
    • 是的,您也可以使用 KSQL。 docs.confluent.io/current/ksql/docs/developer-guide/…
    • 我在示例 KSQL 示例中没有看到的步骤是聚合。事实上,如果没有聚合,记录将被覆盖。我的意思是,据我了解,PartitionBy 只是重新输入密钥,但没有分组和聚合。至少示例中没有提到这一点。你能澄清一下吗?
    • 为什么不使用 GroupBy 和 CollectList 之类的东西,如果第一部分还没有重新设置密钥,则使用 PartitionBY。
    • 是的,你是对的 上面的代码案例在以下SELECT user, COLLECT_LIST(URL_EXTRACT_PATH(url)) AS CLICK_PATH FROM clicks GROUP BY user 中用于KSQL。但要注意COLLECT_LIST 仅支持单列。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-05-08
    • 1970-01-01
    • 2013-05-23
    • 2022-07-22
    • 1970-01-01
    相关资源
    最近更新 更多