eskabetxe commented on code in PR #245:
URL: 
https://github.com/apache/flink-connector-jdbc/pull/245#discussion_r4130868167


##########
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:
   Since lineageQuery takes precedence over the executor's statement, 
post-open() lineage for upsert sinks on dialects with native upsert (e.g. 
Postgres ON CONFLICT) now reports the plain INSERT instead of the upsert 
statement. Same table name and more consistent plan-time vs runtime, just 
confirming it's intentional and worth a line in the PR description/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:
   Nit/out-of-scope note: the DataStream JdbcSink path still goes through this 
constructor, so its getLineageVertex() remains empty before open() — same 
plan-time bug for LineageVertexProvider consumers of DataStream jobs. Not a 
blocker here; suggest a follow-up JIRA to thread a lineage query through 
JdbcSink/JdbcSinkBuilder too.



##########
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:
   Pulling the parsed statement via a static on FieldNamedPreparedStatementImpl 
works, but a tiny helper like LineageUtils.normalizedQuery(sql) (or moving 
parseNamedStatement to the non-Impl utility) would keep the builder from 
depending on an Impl class. Either way fine to leave as is since the static 
already exists.



-- 
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]

Reply via email to