使用 Apache Flink SQL 从 Kafka 消息中获取嵌套字段

bas*_*721 8 apache-flink flink-sql pyflink

我正在尝试使用 Apache Flink 1.11 创建一个源表,我可以在其中访问 JSON 消息中的嵌套属性。我可以从根属性中提取值,但我不确定如何访问嵌套对象。

文档建议它应该是一种MAP类型,但是当我设置它时,出现以下错误

: java.lang.UnsupportedOperationException: class org.apache.calcite.sql.SqlIdentifier: MAP
Run Code Online (Sandbox Code Playgroud)

这是我的 SQL

: java.lang.UnsupportedOperationException: class org.apache.calcite.sql.SqlIdentifier: MAP
Run Code Online (Sandbox Code Playgroud)

我的 JSON 看起来像这样:

        CREATE TABLE input(
            id VARCHAR,
            title VARCHAR,
            properties MAP
        ) WITH (
            'connector' = 'kafka-0.11',
            'topic' = 'my-topic',
            'properties.bootstrap.servers' = 'localhost:9092',
            'properties.group.id' = 'python-test',
            'format' = 'json'
        )
Run Code Online (Sandbox Code Playgroud)

mor*_*aes 8

您可以使用它ROW来提取 JSON 消息中的嵌套字段。您的 DDL 语句将类似于:

CREATE TABLE input(
             id VARCHAR,
             title VARCHAR,
             properties ROW(`foo` VARCHAR)
        ) WITH (
            'connector' = 'kafka-0.11',
            'topic' = 'my-topic',
            'properties.bootstrap.servers' = 'localhost:9092',
            'properties.group.id' = 'python-test',
            'format' = 'json'
        );
Run Code Online (Sandbox Code Playgroud)


ayk*_*dem 7

[2022年更新]

在 Apache Flink 1.13 版本中,没有系统内置 JSON 函数。它们是在 1.14 版本中引入的。检查这个

如果您使用的版本<1.14,请参阅下面的解决方案。

如何使用嵌套 JSON 输入创建表?

JSON 输入示例:

{
    "id": "message-1",
    "title": "Some Title",
    "properties": {
        "foo": "bar",
        "nested_foo":{
            "prop1" : "value1",
            "prop2" : "value2"
        }
    }
}
Run Code Online (Sandbox Code Playgroud)

创建语句

CREATE TABLE input(
                id VARCHAR,
                title VARCHAR,
                properties ROW(`foo` VARCHAR, `nested_foo` ROW(`prop1` VARCHAR, `prop2` VARCHAR))
        ) WITH (
            'connector' = 'kafka-0.11',
            'topic' = 'my-topic',
            'properties.bootstrap.servers' = 'localhost:9092',
            'properties.group.id' = 'python-test',
            'format' = 'json'
        );
Run Code Online (Sandbox Code Playgroud)

如何选择嵌套列?

SELECT properties.foo, properties.nested_foo.prop1 FROM input;
Run Code Online (Sandbox Code Playgroud)

请注意,如果您输出结果

SELECT properties FROM input
Run Code Online (Sandbox Code Playgroud)

您会看到行格式的结果。该专栏的内容properties将是

+I[bar, +I[prop1,prop2]]
Run Code Online (Sandbox Code Playgroud)