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


##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergIO.java:
##########
@@ -396,6 +397,23 @@ public static WriteRows writeRows(IcebergCatalogConfig 
catalog) {
         .build();
   }
 
+  /**
+   * Returns a {@link WriteCdcRows} transform: the CDC sink, applying 
inserts/updates/deletes from a
+   * {@code PCollection<Row>} of change records (each carrying a {@link
+   * org.apache.beam.sdk.values.ValueKind}) to one or more Iceberg V2+ tables 
via equality deletes;
+   * superseded rows are never written.
+   *
+   * <pre>{@code
+   * input.apply(IcebergIO.writeCdcRows(catalogConfig)
+   *     .to(tableId)
+   *     .withSequenceNumberColumn("seq")
+   *     .withTriggeringFrequency(Duration.standardMinutes(1)));
+   * }</pre>
+   */
+  public static WriteCdcRows writeCdcRows(IcebergCatalogConfig catalog) {

Review Comment:
   Yep will follow up with another PR for the SchemaTransform 



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