Hello Gents, I am new at pyflink and I want to explore it, but I got stucked when I tried to create kafka table using flink connector and got below error Traceback (most recent call last):
File "C:\Users\Owner\PycharmProjects\pythonProject4\flink_app.py", line 63, in <module>
main()
File "C:\Users\Owner\PycharmProjects\pythonProject4\flink_app.py", line 58, in main
tbl_env.execute_sql("SELECT * FROM sales_usd LIMIT 10").print()
File "C:\Users\Owner\PycharmProjects\pythonProject4\.venv\lib\site-packages\pyflink\table\table_environment.py", line 804, in execute_sql
return TableResult(self._j_tenv.executeSql(stmt))
File "C:\Users\Owner\PycharmProjects\pythonProject4\.venv\lib\site-packages\py4j\java_gateway.py", line 1285, in call
return_value = get_return_value(
File "C:\Users\Owner\PycharmProjects\pythonProject4\.venv\lib\site-packages\pyflink\util\exceptions.py", line 162, in deco
raise java_exception
pyflink.util.exceptions.TableException: Failed to execute sql
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeQueryOperation(TableEnvironmentImpl.java:810)
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:1223)
at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeSql(TableEnvironmentImpl.java:728)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
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.lang.Thread.run(Thread.java:748) when I tried to do select query from the kafka table