DTStack / DTStack/chunjun

[Bug] [mysqlcdc] 运行mysqlcdc脚本报错

Open
#1,752 5 comments 0 reactions 0 assignees View on GitHub
bug
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 分支最新代码

![image](https://github.com/DTStack/chunjun/assets/8185717/40fafc25-dfff-4bf0-9141-8bfba33fb203)

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

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.