【问题标题】:Is there a KSQL statement to update values in table?是否有用于更新表中值的 KSQL 语句?
【发布时间】:2023-03-03 10:52:01
【问题描述】:

我在主题中有一个数据流,应该被视为 ksql 表(只有给定键的最后一个值很重要),这些数据是关于其他主题中某些数据特定字段的更新。 KSQLDB 中是否有任何方法可以处理更新其他流/表/主题中的值的流?目标主题的实体有 20 个字段,但我的包含 update 的流更新了 3 个字段,所以我只想更新这 3 个字段,其他 17 个字段在目标主题中应保持不变(视为表)。

【问题讨论】:

    标签: apache-kafka ksqldb


    【解决方案1】:

    您可以使用 JOIN STATEMENT 稍作调整来解决您的问题,按照示例,将创建一个包含 5 个字段的表,但只需要从另一个表更新字段技能和级别。

    1.从源主题创建表:

    CREATE TABLE TBL_EMPLOYEE( `employee_id` VARCHAR, `name` varchar, `lastName` varchar, `age` INT, `skill` VARCHAR, `level` VARCHAR ) WITH ( KAFKA_TOPIC = 'employee-topic-input', PARTITIONS = 3, VALUE_FORMAT = 'JSON', KEY = '`employee_id`');
    

    2.创建表来处理所需的更新(它可以是流或表,由另一个查询产生)

    CREATE TABLE TBL_EMPLOYEE_DESIRED_UPDATES (`employee_id` VARCHAR, `skill` VARCHAR, `level` VARCHAR) WITH( KAFKA_TOPIC = 'employee-desired-updates-topic', PARTITIONS = 3, VALUE_FORMAT ='JSON', KEY = '`employee_id`');
    

    3.创建最终表以更新必填字段,左连接允许第一个表上的所有元素。如果第二个表没有任何更新,技能和等级字段将是相同的。

    SET 'auto.offset.reset' = 'earliest';
    CREATE TABLE TBL_EMPLOYEE_FINAL AS 
        SELECT 
            EMP.`employee_id` AS `employee_id`, 
            EMP.`name` AS `name`, 
            EMP.`lastName` AS `lastName`, 
            IFNULL(UPD.`skill`, EMP.`skill`) as `skill`, 
            IFNULL(UPD.`level`, EMP.`level`) as `level` 
        FROM TBL_EMPLOYEE AS EMP 
        LEFT JOIN TBL_EMPLOYEE_DESIRED_UPDATES  UPD ON EMP.ROWKEY = UPD.ROWKEY EMIT CHANGES;
    

    例子:

    INSERT INTO TBL_EMPLOYEE (`employee_id`, `name`, `lastName`, `age`, `skill`, `level`) VALUES ('117', 'John', 'Constantine', 30, 'java', 'jr');
    INSERT INTO TBL_EMPLOYEE (`employee_id`, `name`, `lastName`, `age`, `skill`, `level`) VALUES ('118', 'Anthony', 'Stark', 40, 'AWS', 'architect');
    INSERT INTO TBL_EMPLOYEE (`employee_id`, `name`, `lastName`, `age`, `skill`, `level`) VALUES ('119', 'Clark', 'Kent', 35, 'python', 'senior');
    ksql> SELECT * FROM TBL_EMPLOYEE_FINAL EMIT CHANGES;
    +-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+
    |ROWTIME                        |ROWKEY                         |employee_id                    |name                           |lastName                       |skill                          |level                          |
    +-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+
    |1611440363833                  |119                            |119                            |Clark                          |Kent                           |python                         |senior                         |
    |1611440361284                  |117                            |117                            |John                           |Constantine                    |java                           |jr                             |
    |1611440361408                  |118                            |118                            |Anthony                        |Stark                          |AWS                            |architect                      |
    

    第二步是发送更新

    INSERT INTO TBL_EMPLOYEE_DESIRED_UPDATES  (`employee_id`, `skill`, `level` ) VALUES ('118', 'mongo', 'senior');
    

    结果

    ksql> SELECT  * from TBL_EMPLOYEE_FINAL EMIT CHANGES;
    +-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+
    |ROWTIME                        |ROWKEY                         |employee_id                    |name                           |lastName                       |skill                          |level                          |
    +-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+-------------------------------+
    |1611440363833                  |119                            |119                            |Clark                          |Kent                           |python                         |senior                         |
    |1611440361284                  |117                            |117                            |John                           |Constantine                    |java                           |jr                             |
    |1611440361408                  |118                            |118                            |Anthony                        |Stark                          |AWS                            |architect                      |
    |1611440585726                  |118                            |118                            |Anthony                        |Stark                          |mongo                          |senior                         |
    

    您必须将表中的最新元素视为具有两次修改的新元素。另一个是表更改日志的一部分。记录是不可变的。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2010-09-24
      • 1970-01-01
      • 1970-01-01
      • 2014-02-28
      • 1970-01-01
      • 1970-01-01
      • 2018-10-15
      • 2018-08-03
      相关资源
      最近更新 更多