【发布时间】:2021-04-29 22:18:57
【问题描述】:
我正在使用 Flink Table API。我有一个表定义,我想选择所有字段并将它们转换为新字段中的 JSON 字符串。
我的表有三个字段; a: String, b: Int, c: Timestamp.
如果我这样做
INSERT INTO kinesis
SELECT a, b, c from my_table
kinesis 流有 json 记录;
{
"a" : value,
"b": value,
"c": value
}
但是,我想要类似于 Spark 的功能;
INSERT INTO kinesis
SELECT "constant_value" as my source, to_json(struct(*)) as playload from my_table
所以,预期的结果是;
{
"my_source": "constant_value",
"payload": "json string from the first example that has a,b,c"
}
我在 Flink 中看不到任何 to_json 或 struct() 函数。可以实现吗?
【问题讨论】:
标签: apache-flink flink-streaming flink-sql