【发布时间】: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