yidawang-shopify commented on code in PR #17668:
URL: https://github.com/apache/iceberg/pull/17668#discussion_r4078093615
##########
flink/v2.1/flink/src/main/java/org/apache/iceberg/flink/source/IcebergTableSource.java:
##########
@@ -204,15 +205,19 @@ public ChangelogMode getChangelogMode() {
@Override
public ScanRuntimeProvider getScanRuntimeProvider(ScanContext
runtimeProviderContext) {
+ // The planner reads a FLIP-314 lineage vertex from a SourceProvider but
not from a
+ // DataStreamScanProvider (see CommonExecTableSourceScan), so expose
IcebergSource
+ // declaratively. The legacy FlinkSource has no lineage to report and
stays on the old path.
+ if
(readableConfig.get(FlinkConfigOptions.TABLE_EXEC_ICEBERG_USE_FLIP27_SOURCE)) {
+ IcebergSource<RowData> source = buildFLIP27Source();
+ return SourceProvider.of(source, scanParallelism(source));
+ }
Review Comment:
This is good call. Yes, I believe this will result the the UID being change
when version is updated.
An explicit config is added for this. Default value is false.
--
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]