【问题标题】:KSQL left join giving 'null' result even when data is present即使存在数据,KSQL 左连接也会给出“空”结果
【发布时间】:2021-06-13 18:41:50
【问题描述】:

我正在学习 K-SQL/KSQL-DB,目前正在探索联接。以下是我遇到的问题。

我有 1 个流 'DRIVERSTREAMREPARTITIONEDKEYED' 和一个表 'COUNTRIES',下面是它们的描述。

ksql> describe DRIVERSTREAMREPARTITIONEDKEYED;
Name: DRIVERSTREAMREPARTITIONEDKEYED
 Field       | Type
--------------------------------------
 COUNTRYCODE | VARCHAR(STRING)  (key)
 NAME        | VARCHAR(STRING)
 RATING      | DOUBLE
--------------------------------------

ksql> describe countries;

Name                 : COUNTRIES
 Field       | Type
----------------------------------------------
 COUNTRYCODE | VARCHAR(STRING)  (primary key)
 COUNTRYNAME | VARCHAR(STRING)
----------------------------------------------

这是他们拥有的样本数据,

ksql> select * from DRIVERSTREAMREPARTITIONEDKEYED emit changes;
+---------------------------------------------+---------------------------------------------+---------------------------------------------+
|COUNTRYCODE                                  |NAME                                         |RATING                                       |
+---------------------------------------------+---------------------------------------------+---------------------------------------------+
|SGP                                          |Suresh                                       |3.5                                          |
|IND                                          |Mahesh                                       |2.4                                          |

ksql> select * from countries emit changes;
+---------------------------------------------------------------------+---------------------------------------------------------------------+
|COUNTRYCODE                                                          |COUNTRYNAME                                                          |
+---------------------------------------------------------------------+---------------------------------------------------------------------+
|IND                                                                  |INDIA                                                                |
|SGP                                                                  |SINGAPORE                                                            |

我正在尝试对它们进行“左外”连接,流位于左侧,但下面是我得到的输出,

select d.name,d.rating,c.COUNTRYNAME from DRIVERSTREAMREPARTITIONEDKEYED d left join countries c on d.COUNTRYCODE=c.COUNTRYCODE emit changes;
+---------------------------------------------+---------------------------------------------+---------------------------------------------+
|NAME                                         |RATING                                       |COUNTRYNAME                                  |
+---------------------------------------------+---------------------------------------------+---------------------------------------------+
|Suresh                                       |3.5                                          |null                                         |
|Mahesh                                       |2.4                                          |null                                         |

在理想情况下,我应该将“COUNTRYNAME”列中的数据作为流和数据中的“COUNTRYCODE”列具有匹配的数据。

我尝试了很多搜索但无济于事。 我正在使用“融合平台:6.1.1”

【问题讨论】:

    标签: apache-kafka left-join confluent-platform ksqldb


    【解决方案1】:

    为了让连接起作用,我们有责任验证被连接的两个实体的键是否位于同一个分区中,KsqlDB 无法验证两个连接输入的分区策略是否相同。

    在我的情况下,我的“驱动程序”主题有 2 个分区,我在其上创建了一个流“DriversStream”,该流也有 2 个分区,但我想加入的表“国家”只有 1 个分区,因此,我“重新键入”了“DriversStream”并创建了问题中显示的另一个流“DRIVERSTREAMREPARTITIONEDKEYED”。

    但是表和流的数据不在同一个分区,所以连接失败。

    我用 1 个分区“DRIVERINFO”创建了另一个主题。

     kafka-topics --bootstrap-server localhost:9092 --create --partitions 1 --replication-factor 1 --topic DRIVERINFO
    

    然后在其上创建了一个流“DRIVERINFOSTREAM”。

     CREATE STREAM DRIVERINFOSTREAM (NAME STRING, RATING DOUBLE, COUNTRYCODE STRING) WITH (KAFKA_TOPIC='DRIVERINFO', VALUE_FORMAT='JSON');
    
    

    终于将它加入了 'COUNTRIES' 表,该表终于奏效了。

    ksql> select d.name,d.rating,c.COUNTRYNAME from DRIVERINFOSTREAM d left join countries c on d.COUNTRYCODE=c.COUNTRYCODE EMIT CHANGES;
    +-------------------------------------------+-------------------------------------------+-------------------------------------------+
    |NAME                                       |RATING                                     |COUNTRYNAME                                |
    +-------------------------------------------+-------------------------------------------+-------------------------------------------+
    |Suresh                                     |2.4                                        |SINGAPORE                                  |
    |Mahesh                                     |3.6                                        |INDIA                                      |
    
    
    

    详情请参考以下链接,

    KSQL join

    Partitioning data for Joins

    【讨论】:

      猜你喜欢
      • 2019-11-11
      • 1970-01-01
      • 2018-04-25
      • 2011-02-23
      • 1970-01-01
      • 2019-02-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多