【发布时间】:2022-10-04 21:29:42
【问题描述】:
Kafka 中的INPUT_STREAM 是使用下面的 ksql 语句创建的:
CREATE STREAM INPUT_STREAM (year STRUCT<month STRUCT<day STRUCT<hour INTEGER, minute INTEGER>>>) WITH (KAFKA_TOPIC = 'INPUT_TOPIC', VALUE_FORMAT = 'JSON');
它定义了四个级别的嵌套 json 模式,其中包含字段 year、month 和 day 和 hour 和 minute,如下所示:
{
"year": {
"month": {
"day": {
"hour": string,
"minute": string
}
}
}
}
我想创建第二个OUTPUT_STREAM,它将读取来自INPUT_STREAM 的消息并将其字段名称重新映射到一些自定义名称。我想获取hour 和minute 值并将它们放在one 和two 字段下方的嵌套json 中,如下所示:
{
"one": {
"two": {
"hour": string,
"minute": string
}
}
}
我继续将 ksql 语句放在一起创建OUTPUT_STREAM
CREATE STREAM OUTPUT_STREAM WITH (KAFKA_TOPIC='OUTPUT_TOPIC', REPLICAS=3) AS SELECT YEAR->MONTH->DAY->HOUR ONE->TWO->HOUR FROM INPUT_STREAM EMIT CHANGES;
该语句失败并出现错误。此语句中是否存在语法错误?是否可以像我在这里一样指定目标字段名称
...AS SELECT YEAR->MONTH->DAY->HOUR ONE->TWO->HOUR FROM... ?
我尝试使用STRUCT 而不是ONE->TWO->HOUR:
CREATE STREAM OUTPUT_STREAM WITH (KAFKA_TOPIC='OUTPUT_TOPIC', REPLICAS=3) AS SELECT YEAR->MONTH->DAY->HOUR ONE STRUCT<TWO STRUCT<HOUR VARCHAR>> FROM INPUT_STREAM EMIT CHANGES;
它也出错并且不起作用
【问题讨论】:
标签: apache-kafka confluent-platform confluent-schema-registry ksqldb