【问题标题】:Auto-incrementing column in ksqlDB?ksqlDB中的自动递增列?
【发布时间】:2020-06-18 00:10:15
【问题描述】:

我目前正在使用this process(见下文)在 ksqlDB 中生成一个自动递增列。但现在我想知道这种方法是否存在竞争条件或其他同步问题。这是在 ksqlDB 中生成自动递增列的好方法吗?如果没有,有没有更好的方法?

假设您想将一个 ksqlDB 流中的值插入到另一个 ksqlDB 流中,同时在 目标流。

首先,创建两个流:

CREATE STREAM dest (ROWKEY INT KEY, i INT, x INT) WITH (kafka_topic='test_dest', value_format='json', partitions=1);
CREATE STREAM src (x INT) WITH (kafka_topic='test_src', value_format='json', partitions=1);

接下来,创建一个包含目标流最大值的物化视图。

CREATE TABLE dest_maxi AS SELECT MAX(i) AS i FROM dest GROUP BY 1;

我们需要能够将源流加入物化视图。为此,我们将创建另一个中间流 使用始终设置为 1 的虚拟 one 列,这就是我们对物化视图进行分组的原因:

CREATE STREAM src_one AS SELECT x, 1 AS one FROM src;
INSERT INTO dest SELECT COALESCE(dest_maxi.i,0)+1 AS i, src_one.x AS x FROM src_one LEFT JOIN dest_maxi ON src_one.one = dest_maxi.ROWKEY PARTITION BY COALESCE(dest_maxi.i,0)+1 EMIT CHANGES;

现在您可以将值插入流 src 并观察它们以自动递增的 ID 出现在流 dest 中。

【问题讨论】:

    标签: ksqldb


    【解决方案1】:

    我认为你的方法行不通。 kslqDB 不保证跨两个不同查询处理记录的顺序。在您的情况下,这意味着没有订购保证

    CREATE TABLE dest_maxi AS <query>;
    

    将在之前运行和更新dest_maxi

    INSERT INTO dest <query>;
    

    运行。因此,我认为您会遇到问题。

    您似乎正在尝试获取一系列数字,例如

    1234
    24746
    24848
    4947
    34
    

    并添加一个自动递增的 id 列,以便结果如下所示:

    1, 1234
    2, 24746
    3, 24848
    4, 4947
    5, 34
    

    这样的东西应该会给你你想要的:

    -- source stream of numbers:
    CREATE STREAM src (
         x INT
      ) WITH (
        kafka_topic='test_src', 
        value_format='json'
      );
    
    -- intermediate 'table' of numbers and current count:
    CREATE TABLE with_counter 
       WITH (partitions = 1) AS
       SELECT
          1 as k,
          LATEST_BY_OFFSET(x) as x,
          COUNT(1) AS id
       FROM src
       GROUP BY 1
    
    
    -- if you need this back as a stream in ksqlDB you can run:
    CREATE STREAM dest (
         x INT,
         id BIGINT
       ) WITH (
         kafka_topic='WITH_COUNTER',
         value_format='json'
       );
    

    UDAF 计算每个键的值,因此我们按一个常数分组,确保所有输入行都集中到一个键(和分区 - 所以这不能很好地扩展!)。

    我们使用COUNT 来计算看到的行数,因此它的输出会自动递增,我们使用LATEST_BY_OFFSETx 的当前值抓取到我们的表中。

    with_counter 表的更改日志将包含您想要的输出,只有 1 的常量键:

    1 -> 1, 1234
    1 -> 2, 24746
    1 -> 3, 24848
    1 -> 4, 4947
    1 -> 5, 34
    

    我们将其作为dest 流重新导入ksqlDB。您可以正常使用。如果你想要一个没有密钥的主题,你可以运行:

    CREATE STREAM without_key AS SELECT * FROM dest;
    

    【讨论】:

    • 谢谢安德鲁!这太棒了,而且很有意义。
    猜你喜欢
    • 2016-03-07
    • 2012-04-10
    • 2014-04-15
    • 1970-01-01
    • 1970-01-01
    • 2013-03-03
    • 2020-01-15
    • 2014-10-01
    • 1970-01-01
    相关资源
    最近更新 更多