【问题标题】:Select all fields as json string as new field in Flink SQL选择所有字段作为 json 字符串作为 Flink SQL 中的新字段
【发布时间】: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


    【解决方案1】:

    您可能必须实现自己的用户定义聚合函数。

    这就是我所做的,这里我假设 UDF 的输入看起来像

    to_json('col1', col1, 'col2', col2)

    public class RowToJson extends ScalarFunction {
        public String eval(@DataTypeHint(inputGroup = InputGroup.ANY) Object... row) throws Exception {
            if(row.length % 2 != 0) {
                throw new Exception("Wrong key/value pairs!");
            }
    
            String json = IntStream.range(0, row.length).filter(index -> index % 2 == 0).mapToObj(index -> {
                String name = row[index].toString();
                Object value = row[index+1];
                ... ...
            }).collect(Collectors.joining(",", "{", "}"));
            return json;
        }
    }
    

    如果您希望 udf 可用于分组依据,则必须从 AggregateFunction 扩展您的 udf 类

    public class RowsToJson extends AggregateFunction<String, List<String>>{
        @Override
        public String getValue(List<String> accumulator) {
            return accumulator.stream().collect(Collectors.joining(",", "[", "]"));
        }
    
        @Override
        public List<String> createAccumulator() {
            return new ArrayList<String>();
        }
    
        public void accumulate(List<String> acc, @DataTypeHint(inputGroup = InputGroup.ANY) Object... row) throws Exception {
            if(row.length % 2 != 0) {
                throw new Exception("Wrong key/value pairs!");
            }
            String json = IntStream.range(0, row.length).filter(index -> index % 2 == 0).mapToObj(index -> {
                String name = row[index].toString();
                Object value = row[index+1];
                ... ...
            }).collect(Collectors.joining(",", "{", "}"));
            acc.add(json);
        }
    
    }
    

    【讨论】:

    • 这并不能完全回答问题。您介意分享您的代码/指南吗?
    猜你喜欢
    • 2011-06-21
    • 2013-08-06
    • 2016-10-23
    • 1970-01-01
    • 1970-01-01
    • 2011-06-18
    • 1970-01-01
    • 1970-01-01
    • 2020-06-04
    相关资源
    最近更新 更多