[Bug] [mysqlcdc] 运行mysqlcdc脚本报错
- Dominant language
- Java
- Stars
- 4.1k
- Forks
- 1.7k
- PR merge metrics
- No merged PRs in 30d
Description
### Search before asking
- [X] I had searched in the [issues](https://github.com/DTStack/chunjun/issues) and found no similar issues.
### What happened
flink 版本是 1.12.7,chunjun 是 chunjun-1.12_release 分支最新代码

Standalone 模式下运行 mysqlcdc 脚本报错
Could not find any factory for identifier 'mysql-cdc' that implements 'org.apache.flink.table.factories.DynamicTableSourceFactory' in the classpath
具体异常信息如下
```
Unable to create a source for reading table 'default_catalog.default_database.test1'.
Table options are:
'connector'='mysql-cdc'
'database-name'='dgy'
'hostname'='192.168.1.121'
'password'='******'
'port'='3306'
'table-name'='test1'
'username'='root'
at org.apache.flink.table.factories.FactoryUtil.createTableSource(FactoryUtil.java:177)
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.convertJoin(SqlToRelConverter.java:2864)
at org.apache.calcite.sql2rel.SqlToRelConverter.convertFrom(SqlToRelConverter.java:2162)
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.convertFrom(SqlToRelConverter.java:2169)
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.chunjun.sql.parser.InsertStmtParser.execStmt(InsertStmtParser.java:47)
at com.dtstack.chunjun.sql.parser.AbstractStmtParser.handleStmt(AbstractStmtParser.java:50)
at com.dtstack.chunjun.sql.parser.AbstractStmtParser.handleStmt(AbstractStmtParser.java:52)
at com.dtstack.chunjun.sql.parser.AbstractStmtParser.handleStmt(AbstractStmtParser.java:52)
at com.dtstack.chunjun.sql.parser.SqlParser.lambda$parseSql$1(SqlParser.java:69)
... 24 more
Caused by: java.lang.RuntimeException: Could not find any factory for identifier 'mysql-cdc' 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:303)
at org.apache.flink.table.factories.FactoryUtil.getDynamicTableFactory(FactoryUtil.java:420)
at org.apache.flink.table.factories.FactoryUtil.createTableSource(FactoryUtil.java:173)
... 58 more
```
### What you expected to happen
脚本正常运行
### How to reproduce
```
CREATE TABLE test1(
id INT,
name varchar,
p_id int,
updateTime timestamp,
primary key(id) not enforced
)WITH(
'connector'='mysql-cdc',
'hostname'='192.168.1.121',
'port'='3306',
'database-name'='dgy',
'table-name'='test1',
'username'='root',
'password'='Datxx2023'
);
CREATE TABLE test2(
id int,
type varchar,
updateTime timestamp,
primary key(id) not enforced
)WITH(
'connector'='mysql-cdc',
'hostname'='192.168.1.121',
'port'='3306',
'database-name'='dgy',
'table-name'='test2',
'username'='root',
'password'='Dxx0x'
);
CREATE TABLE testall(
id int,
name varchar,
type varchar,
updateTime1 timestamp,
updateTime2 timestamp,
primary key(id) not enforced
) WITH (
'connector' = 'mysql-x',
'url' = 'jdbc:mysql://192.168.1.121:3306/dgy',
'table-name' = 'testall',
'username' = 'root',
'password' = 'Dxxxx023',
'sink.buffer-flush.max-rows' = '1024', -- 批量写数据条数,默认:1024
'sink.buffer-flush.interval' = '3000', -- 批量写时间间隔,默认:3000毫秒
'sink.all-replace' = 'true', -- 解释如下(其他rdb数据库类似):默认:false。定义了PRIMARY KEY才有效,否则是追加语句
-- sink.all-replace = 'true' 生成如:INSERT INTO `result3`(`mid`, `mbb`, `sid`, `sbb`) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE `mid`=VALUES(`mid`), `mbb`=VALUES(`mbb`), `sid`=VALUES(`sid`), `sbb`=VALUES(`sbb`) 。会将所有的数据都替换。
-- sink.all-replace = 'false' 生成如:INSERT INTO `result3`(`mid`, `mbb`, `sid`, `sbb`) VALUES (?, ?, ?, ?) ON DUPLICATE KEY UPDATE `mid`=IFNULL(VALUES(`mid`),`mid`), `mbb`=IFNULL(VALUES(`mbb`),`mbb`), `sid`=IFNULL(VALUES(`sid`),`sid`), `sbb`=IFNULL(VALUES(`sbb`),`sbb`) 。如果新值为null,数据库中的旧值不为null,则不会覆盖。
'sink.parallelism' = '1' -- 写入结果的并行度,默认:null
);
CREATE TABLE testall_es
(
id int,
name varchar,
type varchar,
updateTime1 timestamp,
updateTime2 timestamp,
primary key(id) not enforced
)
WITH (
'connector' = 'elasticsearch7-x',
'hosts' = '192.168.1.227:9200',
'index' = 'testall',
'client.connect-timeout' = '10000'
);
insert
into
testall
select
id,
name,
type,
updateTime1,
updateTime2
from
( SELECT
ck.id,
ck.name,
py.type,
ck.updateTime as updateTime1,
py.updateTime as updateTime2
from
test1 ck
left join
test2 py
on ck.p_id = py.id ) tt;
insert
into
testall_es
select
id,
name,
type,
updateTime1,
updateTime2
from
( SELECT
ck.id,
ck.name,
py.type,
ck.updateTime as updateTime1,
py.updateTime as updateTime2
from
test1 ck
left join
test2 py
on ck.p_id = py.id ) tt;
```
```
./flink-1.12.7111/bin/start-cluster.sh
sh bin/chunjun-standalone.sh -job test.sql
```
### Anything else
_No response_
### Version
1.12_release
### Are you willing to submit PR?
- [X] Yes I am willing to submit a PR!
### Code of Conduct
- [X] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct)
Contributor guide
Assessment
This issue has not been assessed yet.