DTStack / DTStack/chunjun

local模式测试sink使用 file-x ,source使用 stream-x报错,这是为什么呢

Open
#683 2 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
4.1k
Forks
1.7k
PR merge metrics
No merged PRs in 30d

Description

测试数据:
![image](https://user-images.githubusercontent.com/26226118/162566121-52dfd6f6-e0fe-493c-a0ba-f08a9670de83.png)

sql文件:
![image](https://user-images.githubusercontent.com/26226118/162566145-c39f7c71-d76d-4bb4-82dd-e135036de62b.png)

执行命令:
bin/flinkx -mode local -jobType sql -job /data/flinkx/job.sql -flinkxDistDir flinkx-dist

提示报错:
----------sql start---------
1>
2>
3>
4> INSERT INTO SINK
5> SELECT *
6> FROM FILE_SOURCE

----------sql end---------

Could not find any factory for identifier 'file-x' that implements 'org.apache.flink.table.factories.DynamicTableSourceFactory' in the classpath.

Available factory identifiers are:

datagen
filesystem

Unable to create a source for reading table 'default_catalog.default_database.FILE_SOURCE'.

Table options are:

'connector'='file-x'
'format'='csv'
'path'='/data/flinkx/text.csv'
at com.dtstack.flinkx.Main.exeSqlJob(Main.java:149)
at com.dtstack.flinkx.Main.main(Main.java:108)
at com.dtstack.flinkx.client.local.LocalClusterClientHelper.submit(LocalClusterClientHelper.java:35)
at com.dtstack.flinkx.client.Launcher.main(Launcher.java:126)
Caused by: com.dtstack.flinkx.throwable.DtSqlParserException:
----------sql start---------
1>
2>
3>
4> INSERT INTO SINK
5> SELECT *
6> FROM FILE_SOURCE

----------sql end---------

Could not find any factory for identifier 'file-x' that implements 'org.apache.flink.table.factories.DynamicTableSourceFactory' in the classpath.

Available factory identifiers are:

datagen
filesystem

Unable to create a source for reading table 'default_catalog.default_database.FILE_SOURCE'.

Table options are:

'connector'='file-x'
'format'='csv'
'path'='/data/flinkx/text.csv'
at com.dtstack.flinkx.sql.parser.SqlParser.lambda$parseSql$1(SqlParser.java:71)
at java.util.stream.ForEachOps$ForEachOp$OfRef.accept(ForEachOps.java:184)
at java.util.stream.ReferencePipeline$2$1.accept(ReferencePipeline.java:175)
at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1382)
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:481)
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:471)
at java.util.stream.ForEachOps$ForEachOp.evaluateSequential(ForEachOps.java:151)
at java.util.stream.ForEachOps$ForEachOp$OfRef.evaluateSequential(ForEachOps.java:174)
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234)
at java.util.stream.ReferencePipeline.forEach(ReferencePipeline.java:418)
at com.dtstack.flinkx.sql.parser.SqlParser.parseSql(SqlParser.java:65)
at com.dtstack.flinkx.Main.exeSqlJob(Main.java:140)
... 3 more
Caused by: org.apache.flink.table.api.ValidationException: Could not find any factory for identifier 'file-x' that implements 'org.apache.flink.table.factories.DynamicTableSourceFactory' in the classpath.

Available factory identifiers are:

datagen
filesystem

Unable to create a source for reading table 'default_catalog.default_database.FILE_SOURCE'.

Table options are:

'connector'='file-x'
'format'='csv'
'path'='/data/flinkx/text.csv'
at org.apache.flink.table.factories.FactoryUtil.createTableSource(FactoryUtil.java:176)
at org.apache.flink.table.planner.plan.schema.CatalogSourceTable.createDynamicTableSource(CatalogSourceTable.java:254)
at org.apache.flink.table.planner.plan.schema.CatalogSourceTable.toRel(CatalogSourceTable.java:100)
at org.apache.calcite.sql2rel.SqlToRelConverter.toRel(SqlToRelConverter.java:3585)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertIdentifier(SqlToRelConverter.java:2507)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertFrom(SqlToRelConverter.java:2144)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertFrom(SqlToRelConverter.java:2093)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertFrom(SqlToRelConverter.java:2050)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertSelectImpl(SqlToRelConverter.java:663)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertSelect(SqlToRelConverter.java:644)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertQueryRecursive(SqlToRelConverter.java:3438)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertQuery(SqlToRelConverter.java:570)
at org.apache.flink.table.planner.calcite.FlinkPlannerImpl.org$apache$flink$table$planner$calcite$FlinkPlannerImpl$$rel(FlinkPlannerImpl.scala:165)
at org.apache.flink.table.planner.calcite.FlinkPlannerImpl.rel(FlinkPlannerImpl.scala:157)
at org.apache.flink.table.planner.operations.SqlToOperationConverter.toQueryOperation(SqlToOperationConverter.java:902)
at org.apache.flink.table.planner.operations.SqlToOperationConverter.convertSqlQuery(SqlToOperationConverter.java:871)
at org.apache.flink.table.planner.operations.SqlToOperationConverter.convert(SqlToOperationConverter.java:250)
at org.apache.flink.table.planner.operations.SqlToOperationConverter.convertSqlInsert(SqlToOperationConverter.java:564)
at org.apache.flink.table.planner.operations.SqlToOperationConverter.convert(SqlToOperationConverter.java:248)
at org.apache.flink.table.planner.delegation.ParserImpl.parse(ParserImpl.java:77)
at org.apache.flink.table.api.internal.StatementSetImpl.addInsertSql(StatementSetImpl.java:50)
at com.dtstack.flinkx.sql.parser.InsertStmtParser.execStmt(InsertStmtParser.java:47)
at com.dtstack.flinkx.sql.parser.AbstractStmtParser.handleStmt(AbstractStmtParser.java:50)
at com.dtstack.flinkx.sql.parser.AbstractStmtParser.handleStmt(AbstractStmtParser.java:52)
at com.dtstack.flinkx.sql.parser.AbstractStmtParser.handleStmt(AbstractStmtParser.java:52)
at com.dtstack.flinkx.sql.parser.SqlParser.lambda$parseSql$1(SqlParser.java:68)
... 14 more
Caused by: java.lang.RuntimeException: Could not find any factory for identifier 'file-x' that implements 'org.apache.flink.table.factories.DynamicTableSourceFactory' in the classpath.

Available factory identifiers are:

datagen
filesystem
at org.apache.flink.table.factories.FactoryUtil.discoverFactory(FactoryUtil.java:302)
at org.apache.flink.table.factories.FactoryUtil.getDynamicTableFactory(FactoryUtil.java:419)
at org.apache.flink.table.factories.FactoryUtil.createTableSource(FactoryUtil.java:172)
... 39 more

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.