【问题标题】:Pipeline is not ingesting data into memsql table using procedure管道未使用过程将数据摄取到 memsql 表中
【发布时间】:2019-02-06 07:01:58
【问题描述】:

我将 json(20 键值对) 推送到 kafka 并能够使用它 - 对其进行测试以验证数据是否成功推送到 kafka。

以下脚本正在创建管道,但未将数据加载到 memsql 表中。我是否需要为 JSON 数据类型修改我的创建管道脚本。

CREATE OR REPLACE PIPELINE omnitracs_gps_evt_pipeline
AS LOAD DATA KAFKA '192.168.188.110:9092/ib_Omnitracs' 
INTO procedure INGEST_OMNITRACS_EVT_PROC;

DELIMITER //
CREATE OR REPLACE PROCEDURE INGEST_OMNITRACS_EVT_PROC(batch query(evt_json json))
AS
BEGIN
    INSERT INTO TEST(id, name) 
      SELECT evt_json::ignition,evt_json::positiontype
      FROM batch;
      ECHO SELECT 'HELLO';
END
//
DELIMITER ; 

TEST PIPELINE omnitracs_gps_evt_pipeline LIMIT 5;
START PIPELINE omnitracs_gps_evt_pipeline FOREGROUND LIMIT 5 BATCHES;

任何人都可以帮助它应该是什么。

【问题讨论】:

    标签: apache-kafka singlestore


    【解决方案1】:

    您可能应该修改 CREATE PIPELINE 的 AS LOAD DATA 子句以执行本机 JSON 加载,如下所述:https://docs.memsql.com/sql-reference/v6.7/load-data/#json-load-data

    有两个原因:

    • 所写的管道将期望来自 kafka 的输入位于 TSV 中 带有 1 个字段的格式。 TSV 是默认格式,它推断出预期的字段数 从参数到目标存储过程。实际上,输入 JSON 记录很有可能会成功解析,但我不会依赖这个。

    • 使用原生 JSON 管道的 subvalue_mapping 子句来执行 提取和插入 ::ignition 和 ::positiontype, 完全跳过存储过程的开销。此外,所写的管道将 实例化内存中的临时 JSON 数据结构,这是相对的 很贵。

    我建议如下:

    CREATE OR REPLACE PIPELINE omnitracs_gps_evt_pipeline
    AS LOAD DATA KAFKA '192.168.188.110:9092/ib_Omnitracs' 
    INTO TABLE TEST
    FORMAT JSON
    ( 
      id <- ignition_event,
      name <- position_type
    );
    

    【讨论】:

    • 我用上面提到的方法修改了我的管道脚本。但它仍然无法正常工作。 0 行受影响。 >测试管道omnitracs_gps_evt_pipeline LIMIT 5; >START PIPELINE omnitracs_gps_evt_pipeline FOREGROUND LIMIT 5 BATCHES;
    • 你能发给我们SELECT * FROM information_schema.PIPELINES_CURSORS吗?
    • DATABASE_NAME : MEMSQLDB PIPELINE_NAME : omnitracs_gps_evt_pipeline SOURCE_TYPE : KAFKA SOURCE_PARTITION_ID : 0 EARLIEST_OFFSET : 0 LATEST_OFFSET : 1320 CURSOR_OFFSET : 0 SUCCESSFUL_CURSOR_OFFSET : NULL UPDATED_UNIX_TIMESTAMP : 1549542203.604487 EXTRA_FIELDS : NULL
    【解决方案2】:

    管道的存储过程中不允许使用 ECHO SELECT。当您运行 START PIPELINE ... FOREGROUND 或在 CREATE PIPELINE 时间(如果已定义该过程)时,您应该得到一个错误提示。

    【讨论】:

      【解决方案3】:

      在 kafka 从生产者中删除 ProducerConfig.TRANSACTIONAL_ID_CONFIG 配置后,管道现在正在工作。

      CREATE PIPELINE FEB13_PIPELINE_2
      AS LOAD DATA KAFKA '192.168.188.110:9092/FEB13_PROC' 
      INTO procedure INGEST_EVT_PROC;
      
      DELIMITER //
      CREATE OR REPLACE PROCEDURE INGEST_EVT_PROC(batch query(evt_json json))
      AS
      BEGIN
          INSERT INTO TEST_FEB13(ID, NAME) 
            SELECT evt_json::ID,evt_json::NAME
            FROM batch;
      END
      //
      DELIMITER ;
      

      现在只是一个小疑问,即使双引号也被添加到表格列中。如何逃脱它。 JSON发送到kafka:“{'ID':1,'NAME':\'a\'}”

      【讨论】:

      • JSON_EXTRACT_ 会做需要的
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2022-08-24
      • 1970-01-01
      • 2015-10-08
      • 2022-01-13
      • 1970-01-01
      • 1970-01-01
      • 2017-09-28
      相关资源
      最近更新 更多