yidawang-shopify commented on code in PR #17668:
URL: https://github.com/apache/iceberg/pull/17668#discussion_r4078181010


##########
flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/IcebergTableSink.java:
##########
@@ -124,49 +130,139 @@ public IcebergTableSink(
     this.useDynamicSink = true;
   }
 
-  @SuppressWarnings("deprecation")
   @Override
   public SinkRuntimeProvider getSinkRuntimeProvider(Context context) {
     Preconditions.checkState(
         !overwrite || context.isBounded(),
         "Unbounded data stream doesn't support overwrite operation.");
 
+    if (canProvideSinkV2()) {
+      IcebergSink sink = buildIcebergSink();
+      Integer parallelism = sink.writeParallelism();
+      return parallelism != null ? SinkV2Provider.of(sink, parallelism) : 
SinkV2Provider.of(sink);
+    }
+
     return (DataStreamSinkProvider)
         (providerContext, dataStream) -> {
           if (useDynamicSink) {
             return createDynamicIcebergSink(dataStream);
           }
 
-          ResolvedSchema physicalColumnsOnlySchema = null;
-          List<String> equalityColumns;
-          if (resolvedSchema != null) {
-            physicalColumnsOnlySchema =
-                ResolvedSchema.of(
-                    resolvedSchema.getColumns().stream()
-                        .filter(Column::isPhysical)
-                        .collect(Collectors.toList()));
-
-            equalityColumns =
-                physicalColumnsOnlySchema
-                    .getPrimaryKey()
-                    .map(UniqueConstraint::getColumns)
-                    .orElseGet(ImmutableList::of);
-          } else {
-            equalityColumns =
-                tableSchema
-                    .getPrimaryKey()
-                    
.map(org.apache.flink.table.legacy.api.constraints.UniqueConstraint::getColumns)
-                    .orElseGet(ImmutableList::of);
-          }
-
+          ResolvedSchema physicalColumnsOnlySchema = 
physicalColumnsOnlySchema();
+          List<String> equalityColumns = 
equalityColumns(physicalColumnsOnlySchema);
           if 
(readableConfig.get(FlinkConfigOptions.TABLE_EXEC_ICEBERG_USE_V2_SINK)) {
             return createIcebergSink(dataStream, equalityColumns, 
physicalColumnsOnlySchema);
-          } else {
-            return createLegacySink(dataStream, equalityColumns, 
physicalColumnsOnlySchema);
           }
+
+          return createLegacySink(dataStream, equalityColumns, 
physicalColumnsOnlySchema);
         };
   }
 
+  /**
+   * Whether the sink can be exposed as a {@link SinkV2Provider}, which is 
what it takes to report
+   * sink lineage: the planner reads a FLIP-314 vertex off the {@code Sink} 
object (see {@code
+   * CommonExecSink}), whereas a {@code DataStreamSinkProvider} only hands it 
a built
+   * transformation. Requires {@link IcebergSink}, the only sink that reports 
lineage.
+   *
+   * <p>Also requires {@code TABLE_EXEC_UID_GENERATION=ALWAYS}. {@link 
IcebergSink}'s custom commit
+   * topology puts explicit uids on its operators, so Flink demands one on the 
sink transformation
+   * too ({@code SinkTransformationTranslator.SinkExpander}). The planner only 
sets it under {@code
+   * ALWAYS} — under the default {@code PLAN_ONLY} only for a compiled plan, 
which a connector
+   * cannot detect — so taking this path otherwise would fail job submission 
outright.
+   */
+  private boolean canProvideSinkV2() {
+    if (useDynamicSink || 
!readableConfig.get(FlinkConfigOptions.TABLE_EXEC_ICEBERG_USE_V2_SINK)) {
+      return false;
+    }
+
+    ExecutionConfigOptions.UidGeneration uidGeneration =
+        readableConfig.get(ExecutionConfigOptions.TABLE_EXEC_UID_GENERATION);
+    if (uidGeneration != ExecutionConfigOptions.UidGeneration.ALWAYS) {

Review Comment:
   The toggle is added. 
   
   In terms of error message, for now I actually perserve the behavior of if 
lineage is not there, let the process continue to run and give a `warning` 
instead of `info`. 
   
   I did not want to error out because I don't feel safe to add a breaking 
behavior that now the flink-iceberg process will fail if the lineage is not 
there. 
   
   However, let me know if you prefer we fail here and give error message 
instead of just a warning. 



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


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to