【问题标题】:KSQLDB: Select Fields As ArrayKSQLDB:选择字段作为数组
【发布时间】:2021-03-24 21:58:40
【问题描述】:

我正在尝试通过 KsqlDb 将数据从一个主题(读取主题)转换为另一个主题(写入主题)。

这是已经产生到阅读主题的数据

{
  "orderNumber": "01235656",
  "deliveryBarcode": "733998877",
  "requestId": "1616516663000",
  "status": "APPROVED_BY_SUPERVISOR"
}

我编写了这些 ksqldb 查询:

-- The general stream to read the topic is like this:
CREATE STREAM GENERAL_STREAM (
    deliveryBarcode VARCHAR,
    orderNumber VARCHAR,
    requestId VARCHAR,
    status VARCHAR
) WITH (
    kafka_topic = 'read-topic',
    value_format = 'json'
);


-- This is the stream to redirect the filtered data throgh 'write-topic'
CREATE STREAM REDIRECTION_STREAM
WITH (
    partitions = 6,
    replicas = 3,
    kafka_topic = 'write-topic',
    value_format = 'json'
) AS
SELECT
       AS_VALUE(requestId) `requestId`,
       ARRAY<STRUCT<
        deliveryBarcode 
        orderNumber
       >> `packages`
FROM EARTH_DELIVERY_COURIER_PUDOPACKAGESTATUSUPDATED_0
WHERE (payload -> status = 'APPROVED_BY_SUPERVISOR')
EMIT CHANGES;

但由于这部分原因,我的查询不起作用:

ARRAY<STRUCT<
        deliveryBarcode 
        orderNumber
       >> `packages`

我对 write-topic 的预期数据是这样的

{
  "requestId": "1616516663000"
  "packages":[
    {
      "ordernumber":"01235656",
      "barcodenumber":"733998877"
    }
  ]
}

我应该如何修改这些查询,以便能够按预期以数组格式生成“包”字段?

【问题讨论】:

    标签: apache-kafka confluent-platform ksqldb


    【解决方案1】:

    您可以使用 array() 和 as_map() 函数来生成您期望的输出。

    这是我用来解决问题的 CSAS:

    CREATE STREAM REDIRECTION_STREAM 
    WITH (kafka_topic='write_topic', value_format='json') 
    AS SELECT 
      AS_VALUE(requestId) requestId, 
      ARRAY[
        AS_MAP(
          ARRAY['deliverybarcode', 'ordernumber'], 
          ARRAY[deliverybarcode, ordernumber]
        )
      ] packages 
    FROM GENERAL_STREAM;
    

    这是上述主题的输出

    print 'write_topic' from BEGINNING;
    Key format: ¯\_(ツ)_/¯ - no data processed
    Value format: JSON or KAFKA_STRING
    rowtime: 2021/03/24 13:24:44.557 Z, key: <null>, value: {"REQUESTID":"1616516663000","PACKAGES":[{"ordernumber":"01235656","deliverybarcode":"733998877"}]}, partition: 0
    ^CTopic 
    

    【讨论】:

    • 出色的解决方案塞尔吉奥。有用。谢谢!
    猜你喜欢
    • 2016-10-01
    • 1970-01-01
    • 2016-10-23
    • 2011-06-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-07-15
    • 2011-07-03
    相关资源
    最近更新 更多