【问题标题】:KSQL table get old and new valueKSQL 表获取新旧值
【发布时间】:2019-03-03 07:51:32
【问题描述】:

在 KSQL 中是否可以从表中流出新旧值?我们想使用一个表作为值的存储,当一个更改流出一个“反转”值,它是前一个,以某种方式标记,以及新值,以便我们可以处理下游系统中的增量?

【问题讨论】:

    标签: apache-kafka ksqldb


    【解决方案1】:

    Kafka 表通常用于存储最新值。因此,例如说表中存在键为“123”的流,并且主题上出现了具有相同键“123”但列值不同的新流,这将覆盖(更新)表中的现有值。

    所以在 Table 上做这件事可能不是一个好主意。

    您的用例对我来说并不清楚,但我的建议是您需要在流源中使用某种机制或使用时间戳来处理增量提要。

    【讨论】:

    • 我们可以在主题上添加一个自定义处理器来执行此操作,我只是认为它非常适合 KSQL
    • 正确,您也可以尝试自定义 Kafka 连接器从各种来源提取数据,或者您可以编写自己的连接器来提取增量提要。
    【解决方案2】:

    是的,这是可能的。确实需要一些杂耍。

    创建表以保持上次状态

    create table v1_mux_connection_ping_ta
    as 
    select 
      assetid, 
      LATEST_BY_OFFSET(pingable) pingable
    from v1_mux_connection_ping_st_parse 
    group by assetid;
    

    问题是它也不会发出任何变化。一种解决方案是将表格转换为流。

    CREATE STREAM v1_mux_connection_ping_ta_s 
      (assetId VARCHAR KEY, pingable VARCHAR) 
    WITH (kafka_topic='V1_MUX_CONNECTION_PING_TA', value_format='JSON');
    

    只得到改变的值

    create table d_opt_details as 
    select 
      s.assetId,
      LATEST_BY_OFFSET(s.pingable) new,
      LATEST_BY_OFFSET(s.pingable, 2)[1] old
    from v1_mux_connection_ping_ta_s s
    group by
      s.assetId;
    
    create table opt_details as 
    select 
      s.assetId, s.new as pingable
    from d_opt_details s
    where new != old;
    

    【讨论】:

      猜你喜欢
      • 2019-10-26
      • 2014-03-21
      • 2017-03-16
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多