【问题标题】:How Create KSQLdb Stream fields from nested JSON Object如何从嵌套的 JSON 对象创建 KSQLdb 流字段
【发布时间】:2020-09-29 23:38:39
【问题描述】:

我有一个主题,我以以下格式发送 json:

 {
  "schema": {
   "type": "string",
   "optional": true
  },
  "payload": “CustomerData{version='1', customerId=‘76813432’,      phone=‘76813432’}”
} 

我想创建一个带有 customerId 和 phone 的流,但我不确定如何根据嵌套的 json 对象定义流。 (已编辑)

CREATE  STREAM customer (
    payload.version VARCHAR,
    payload.customerId VARCHAR,
    payload.phone VARCHAR
  ) WITH (
    KAFKA_TOPIC='customers',
    VALUE_FORMAT='JSON'
  );    

会是这样吗?如何在定义流字段时取消引用嵌套对象?

实际上,上述内容不适用于它所说的字段定义:

Caused by: line 2:12: 
extraneous input '.' expecting {'EMIT', 'CHANGES',
'INTEGER', 'DATE', 'TIME', 'TIMESTAMP', 'INTERVAL', 'YEAR', 'MONTH', 'DAY',

【问题讨论】:

    标签: apache-kafka confluent-platform ksqldb


    【解决方案1】:

    应用函数extractjsonfield

    您可以使用一个名为 extractjsonfield 的 ksqlDB 函数。

    首先,您需要提取 schemapayload 字段:

    CREATE STREAM customer (
      schema VARCHAR,
      payload VARCHAR
    ) WITH (
        KAFKA_TOPIC='customers',
        VALUE_FORMAT='JSON'
    ); 
    

    然后你可以选择json中的嵌套字段:

    SELECT EXTRACTJSONFIELD(payload, '$.version') AS version FROM customer;
    

    但是,您的有效负载数据似乎没有有效的 JSON 格式


    应用 STRUCT 架构

    如果您的整个负载被编码为 JSON 字符串,这意味着您的数据如下所示:

    {
      "schema": {
       "type": "string",
       "optional": true
      },
      "payload": {
        "version"="1",
        "customerId"="76813432",
        "phone"="76813432"
      }
    } 
    

    您可以如下定义 STRUCT:

    CREATE STREAM customer (
      schema STRUCT<
        type VARCHAR,
        optional BOOLEAN>,
      payload STRUCT<
        version VARCHAR,
        customerId VARCHAR,
        phone VARCHAR>
    ) 
    WITH (
        KAFKA_TOPIC='customers',
        VALUE_FORMAT='JSON'
    );
    

    最后可以像这样引用单个字段:

    CREATE STREAM customer_analysis AS
    SELECT
      payload->version as VERSION,
      payload->customerId as CUSTOMER_ID,
      payload->phone as PHONE
    FROM customer
    EMIT CHANGES;
    

    【讨论】:

    • 好的,您能否说明 EXTRACTJSONFIELD(message, '$.payload.version') 将如何应用于流定义中
    • 我想看看如何同时使用 struct 和 EXTRACTJSONFIELD 方法。
    • 如果有效载荷再次有一个子对象,是否可以再次应用 STRUCT?
    • 所以为了澄清我需要知道例如消息是否是在 CREATE STREAM 命令期间的内置对象?以及在这种情况下如何设置流字段名称?
    • 另外请验证 STRING 和 VARCHAR 是否与您在答案中的定义相同?
    猜你喜欢
    • 2021-10-02
    • 1970-01-01
    • 2018-12-29
    • 1970-01-01
    • 2017-02-14
    • 2017-12-22
    • 2014-05-15
    • 1970-01-01
    • 2021-09-26
    相关资源
    最近更新 更多