MartijnVisser commented on code in PR #245:
URL:
https://github.com/apache/flink-connector-jdbc/pull/245#discussion_r4154103641
##########
flink-connector-jdbc-core/src/main/java/org/apache/flink/connector/jdbc/internal/JdbcOutputFormat.java:
##########
@@ -256,10 +271,12 @@ public Connection getConnection() {
@Override
public LineageVertex getLineageVertex() {
+ String query = lineageQuery;
Review Comment:
Intentional, and it's now the table name itself. Flink only reads lineage
before `open()`, and the planner still uses the catalog name (FLINK-39935), so
I dropped the release note.
##########
flink-connector-jdbc-core/src/main/java/org/apache/flink/connector/jdbc/internal/JdbcOutputFormat.java:
##########
@@ -90,9 +92,22 @@ public JdbcOutputFormat(
@Nonnull JdbcConnectionProvider connectionProvider,
@Nonnull JdbcExecutionOptions executionOptions,
@Nonnull StatementExecutorFactory<JdbcExec>
statementExecutorFactory) {
+ this(connectionProvider, executionOptions, statementExecutorFactory,
null);
Review Comment:
The DataStream `JdbcSink` has its own `getLineageVertex()`, which parses the
query statement directly. `JdbcWriter` only creates this format on the task
manager, where nobody reads lineage.
##########
flink-connector-jdbc-core/src/main/java/org/apache/flink/connector/jdbc/core/table/sink/JdbcOutputFormatBuilder.java:
##########
@@ -86,19 +88,24 @@ public JdbcOutputFormatBuilder setFieldDataTypes(DataType[]
fieldDataTypes) {
Arrays.stream(fieldDataTypes)
.map(DataType::getLogicalType)
.toArray(LogicalType[]::new);
+ final String sql =
+ dmlOptions
+ .getDialect()
+ .getInsertIntoStatement(
+ dmlOptions.getTableName(),
dmlOptions.getFieldNames());
+ // Lineage is read at plan time, before open(). The plain INSERT names
the target table in
+ // every dialect; the upsert forms do not all parse.
+ final String lineageQuery =
+ FieldNamedPreparedStatementImpl.parseNamedStatement(sql, new
HashMap<>());
Review Comment:
Turns out no parsing is needed here: the builder already has the table name,
so it passes that directly and no longer touches the Impl class.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]