【问题标题】:KSQL: How to cast JSON string to raw JSONKSQL:如何将 JSON 字符串转换为原始 JSON
【发布时间】:2022-08-03 03:03:19
【问题描述】:

我需要根据特定的 JSON 属性将消息从一个 Kafka 主题复制到另一个主题。也就是说,如果属性值为“A” - 复制消息,否则不复制。我试图找出使用 KSQL 的最简单方法。我的源消息都具有我的测试属性,但在其他方面具有非常不同和复杂的模式。有没有办法为此设置“无模式”?

源消息(示例):

{
    \"data\": {
        \"propertyToCheck\": \"value\",
        ... complex structure ...
    }
}

如果我在流中将“数据”定义为 VARCHAR,则可以使用 EXTRACTJSONFIELD 进一步检查该属性。

CREATE OR REPLACE STREAM Test1 (
    `data` VARCHAR
)
WITH (
    kafka_topic = \'Source_Topic\',
    value_format = \'JSON\'
);

然而,在这种情况下,我的 \"select\" 流将生成数据作为 JSON 字符串而不是原始 JSON(这是我想要的)。

CREATE OR REPLACE STREAM Test2 WITH (
    kafka_topic = \'Target_Topic\',
    value_format = \'JSON\'
)AS 
SELECT
  `data` AS `data`
FROM Test1
EMIT CHANGES;

任何想法如何使这项工作?

    标签: apache-kafka apache-kafka-streams ksqldb


    【解决方案1】:

    这是一种解决方法,但您可以按如下方式实现所需的行为:不要将消息模式定义为 VARCHAR,而是使用 BYTES 类型。然后将 FROM_BYTES 与 EXTRACTJSONFIELD 结合使用,从字节表示中读取您要过滤的属性。

    这是一个例子:

    这是一个源流,包含嵌套的 JSON 数据和一个示例数据行:

    CREATE STREAM test (data STRUCT<FOO VARCHAR, BAR VARCHAR>) with (kafka_topic='test', value_format='json', partitions=1);
    INSERT INTO test (data) VALUES (STRUCT(FOO := 'foo', BAR := 'bar'));
    

    现在,将数据表示为字节(使用 KAFKA 格式),而不是 JSON:

    CREATE STREAM test_bytes (data BYTES) WITH (kafka_topic='test', value_format='kafka');
    

    接下来,根据嵌套的 JSON 数据执行过滤:

    CREATE STREAM test_filtered_bytes WITH (kafka_topic='test_filtered') AS SELECT * FROM test_bytes WHERE extractjsonfield(from_bytes(data, 'utf8'), '$.DATA.FOO') = 'foo';
    

    新创建的主题“test_filtered”现在具有正确的 JSON 格式的数据,类似于源流“test”。我们可以通过以原始格式表示流并将其读回检查来进行验证:

    CREATE STREAM test_filtered (data STRUCT<FOO VARCHAR, BAR VARCHAR>) WITH (kafka_topic='test_filtered', value_format='json');
    SELECT * FROM test_filtered EMIT CHANGES;
    

    我验证了这些示例语句从最新的 ksqlDB 版本 (0.27.2) 开始对我有效。自从引入了 BYTES 类型和相关的内置函数以来,它们应该在所有 ksqlDB 版本上都可以正常工作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-10-10
      • 2021-07-21
      • 2023-03-08
      • 2018-04-04
      • 2018-06-28
      • 2011-03-23
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多