-
Notifications
You must be signed in to change notification settings - Fork 1.8k
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
[FLINK-35599][cdc-connector][jdbc-mysql]Flink cdc pipeline sink jdbc mysql #3433
base: master
Are you sure you want to change the base?
[FLINK-35599][cdc-connector][jdbc-mysql]Flink cdc pipeline sink jdbc mysql #3433
Conversation
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Good Job! And I left some comment. @kissycn
public class MySqlDataSinkFactory implements DataSinkFactory { | ||
private static final Logger LOG = LoggerFactory.getLogger(MySqlDataSinkFactory.class); | ||
|
||
public static final String IDENTIFIER = "mysql-writer"; |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think it should be changed to mysql to be consistent with other sink pipeline connectors.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Will do, thanks!
String jdbcUrl = MysqlPooledDataSinkFactory.INSTANCE.getJdbcUrl(builder.build()); | ||
builder.connUrl(jdbcUrl); | ||
// print configs | ||
Map<String, String> map = config.toMap(); |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
In my experience, when PASSWORD is empty, a NullPointerException will be thrown.
However, MySQL allows the password to be empty.
+ " jdbc.properties.allowPublicKeyRetrieval: true\n" | ||
+ "\n" | ||
+ "sink:\n" | ||
+ " type: mysql-writer\n" |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I think mysql is better.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Will do, thanks!
|
||
@Override | ||
public void addColumn(TableId tableId, AddColumnEvent addedColumns) | ||
throws TableNotExistException, CatalogException { |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Exception that is never thrown: TableNotExistException.
Close: #FLINK-35599