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]