ahmedabu98 commented on code in PR #39600:
URL: https://github.com/apache/beam/pull/39600#discussion_r4188672574


##########
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:
   Thanks for flagging. PTAL #40422



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