ahmedabu98 commented on code in PR #39600:
URL: https://github.com/apache/beam/pull/39600#discussion_r4188731897
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergCdcReadSchemaTransformProvider.java:
##########
@@ -193,6 +196,26 @@ static Builder builder() {
"A subset of column names to exclude from reading. If null or empty,
all columns will be read.")
abstract @Nullable List<String> getDrop();
+ @SchemaFieldDescription(
+ "Column used to derive the source's output watermark. "
+ + "Must be an existing, required, top-level column of type 'long'
or 'timestamp'. "
+ + "If not set, the watermark advances according to snapshot commit
timestamp.")
+ abstract @Nullable String getWatermarkColumn();
+
+ @SchemaFieldDescription(
+ "Time unit used to interpret watermark column of type LONG. One of
NANOSECONDS, MICROSECONDS, "
+ + "MILLISECONDS, SECONDS, MINUTES, HOURS, DAYS. Defaults to
MICROSECONDS.")
+ abstract @Nullable String getWatermarkColumnTimeUnit();
+
+ @SchemaFieldDescription(
+ "List of top-level metadata columns to include with CDC output rows.
Supported columns: \n"
Review Comment:
And #40424 for the release
--
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]