【问题标题】:Invalid SQL identifier - org.apache.flink.sql.parser.impl.ParseException: Encountered "TABLE" at line 2, column 16无效的 SQL 标识符 - org.apache.flink.sql.parser.impl.ParseException:在第 2 行第 16 列遇到“TABLE”
【发布时间】:2021-11-21 13:54:58
【问题描述】:

我正在尝试运行一个 PyFlink 作业,该作业从源 Kafka 主题中获取数据并将其放入 hdfs。有一个奇怪的与 SQL 相关的错误不断出现。这是来自 Apache-Flink (PyFlink) Table API Sink 中的 SQL 语句:

SQL:

sql_statement_sink = """
            CREATE TABLE avro_sink (
                timeTime STRING,
                correlationId STRING,
                spanId STRING,
                appName STRING,
                messageType STRING,
                message STRING,
                tag STRING,
                journey as SPLIT_INDEX(tag, '_', 2)
            ) PARTITIONED BY (
                journey,
                appName,
                messageType
            ) WITH (
                'connector' = 'filesystem',
                'partition.default-name' = 'others',
                'format" = 'avro',
                'path' = 'file:///Users/ahmedawny/PycharmProjects/ms_log_consumer/output'
            ) 
        """

完全错误:

WARNING: An illegal reflective access operation has occurred
WARNING: Illegal reflective access by org.apache.flink.api.java.ClosureCleaner (file:/Users/ahmedawny/PycharmProjects/%20ms_log_consumer/venv/lib/python3.8/site-packages/pyflink/lib/flink-dist_2.11-1.14.0.jar) to field java.util.Properties.serialVersionUID
WARNING: Please consider reporting this to the maintainers of org.apache.flink.api.java.ClosureCleaner
WARNING: Use --illegal-access=warn to enable warnings of further illegal reflective access operations
WARNING: All illegal access operations will be denied in a future release
Traceback (most recent call last):
  File "log_consumer.py", line 96, in <module>
    main(**vars(args))
  File "log_consumer.py", line 77, in main
    statement_set.add_insert(avro_sink, table_known_tag)
  File "/Users/ahmedawny/PycharmProjects/ ms_log_consumer/venv/lib/python3.8/site-packages/pyflink/table/statement_set.py", line 116, in add_insert
    self._j_statement_set.addInsert(target_path_or_descriptor, table._j_table, overwrite)
  File "/Users/ahmedawny/PycharmProjects/ ms_log_consumer/venv/lib/python3.8/site-packages/py4j/java_gateway.py", line 1285, in __call__
    return_value = get_return_value(
  File "/Users/ahmedawny/PycharmProjects/ ms_log_consumer/venv/lib/python3.8/site-packages/pyflink/util/exceptions.py", line 146, in deco
    return f(*a, **kw)
  File "/Users/ahmedawny/PycharmProjects/ ms_log_consumer/venv/lib/python3.8/site-packages/py4j/protocol.py", line 326, in get_return_value
    raise Py4JJavaError(
py4j.protocol.Py4JJavaError: An error occurred while calling o697.addInsert.
: org.apache.flink.table.api.SqlParserException: Invalid SQL identifier 
            CREATE TABLE avro_sink (
                timeTime STRING,
                correlationId STRING,
                spanId STRING,
                appName STRING,
                messageType STRING,
                message STRING,
                tag STRING,
                journey as SPLIT_INDEX(tag, '_', 2)
            ) PARTITIONED BY (
                journey,
                appName,
                messageType
            ) WITH (
                'connector' = 'filesystem',
                'partition.default-name' = 'others',
                'format" = 'avro',
                'path' = 'file:///Users/ahmedawny/PycharmProjects/ms_log_consumer/output'
            ) 
        .
        at org.apache.flink.table.planner.parse.CalciteParser.parseIdentifier(CalciteParser.java:96)
        at org.apache.flink.table.planner.delegation.ParserImpl.parseIdentifier(ParserImpl.java:109)
        at org.apache.flink.table.api.internal.StatementSetImpl.addInsert(StatementSetImpl.java:76)
        at org.apache.flink.table.api.bridge.java.internal.StreamStatementSetImpl.addInsert(StreamStatementSetImpl.java:48)
        at org.apache.flink.table.api.bridge.java.internal.StreamStatementSetImpl.addInsert(StreamStatementSetImpl.java:28)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
        at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
        at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
        at java.base/java.lang.reflect.Method.invoke(Method.java:566)
        at org.apache.flink.api.python.shaded.py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
        at org.apache.flink.api.python.shaded.py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
        at org.apache.flink.api.python.shaded.py4j.Gateway.invoke(Gateway.java:282)
        at org.apache.flink.api.python.shaded.py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
        at org.apache.flink.api.python.shaded.py4j.commands.CallCommand.execute(CallCommand.java:79)
        at org.apache.flink.api.python.shaded.py4j.GatewayConnection.run(GatewayConnection.java:238)
        at java.base/java.lang.Thread.run(Thread.java:834)
Caused by: org.apache.flink.sql.parser.impl.ParseException: Encountered "TABLE" at line 2, column 20.
Was expecting one of:
    <EOF> 
    "." ...
    
        at org.apache.flink.sql.parser.impl.FlinkSqlParserImpl.generateParseException(FlinkSqlParserImpl.java:40981)
        at org.apache.flink.sql.parser.impl.FlinkSqlParserImpl.jj_consume_token(FlinkSqlParserImpl.java:40792)
        at org.apache.flink.sql.parser.impl.FlinkSqlParserImpl.TableApiIdentifier(FlinkSqlParserImpl.java:6316)
        at org.apache.flink.table.planner.parse.CalciteParser.parseIdentifier(CalciteParser.java:87)
        ... 15 more

提前致谢。

添加更多句子,因为 StackOverflow 不允许使用“大部分代码”发布。 添加更多句子作为 StackOverflow 不允许发布“大部分代码”。添加更多句子作为 StackOverflow 不允许发布“大部分代码”。添加更多句子作为 StackOverflow 不允许发布“大部分代码”。添加更多作为 StackOverflow 的句子不允许发布“大部分代码”。添加更多句子作为 StackOverflow 不允许发布“大部分代码”。添加更多句子作为 StackOverflow 不允许发布“大部分代码”。

【问题讨论】:

    标签: sql python-3.x apache-flink flink-sql pyflink


    【解决方案1】:

    journey 字段存在语法错误,请将其更改为journey String。向接收器插入数据时使用SPLIT_INDEX 函数。

    【讨论】:

    • 我这样做了,但事实并非如此。现在它在“STRING”之后的“as”上给出错误。它不喜欢它
    猜你喜欢
    • 2019-09-18
    • 1970-01-01
    • 1970-01-01
    • 2023-03-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-05-02
    • 1970-01-01
    相关资源
    最近更新 更多