chamikaramj commented on code in PR #40229:
URL: https://github.com/apache/beam/pull/40229#discussion_r4085185754
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -100,14 +123,33 @@ public static Builder builder() {
"For a streaming pipeline, sets the limit for lifting bundles into the
direct write path.")
public abstract @Nullable Integer getDirectWriteByteLimit();
+ @SchemaFieldDescription(
Review Comment:
Looking at this map I feel like it might be better to do move these
properties to a different SchemaTransform and to the managed API under
Iceberg.WRITE_CDC.
I think a second nested map that groups different config properties
together can end up being confusing the the end user. I think it's much cleaner
to have these at the top level.
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -18,14 +18,23 @@
package org.apache.beam.sdk.io.iceberg;
import static
org.apache.beam.sdk.io.iceberg.IcebergWriteSchemaTransformProvider.Configuration;
Review Comment:
Can we update Managed.java in the same PR ? If not let's keep the Github
issue open to do it in a subsequent PR.
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -180,6 +222,62 @@ public static Builder builder() {
"Sets the number of parallel buckets/workers used to query the Iceberg
catalog during refreshes. Defaults to 1.")
public abstract @Nullable Integer getPollingBuckets();
+ @SchemaFieldDescription(
+ "Columns defining row identity (equality-delete fields). Defaults to
the destination table's "
+ + "identifier (primary-key) fields. Required if the table doesn't
exist yet. Currently only supported in CDC mode.")
Review Comment:
Will this and other properties marked as "Currently only supported in CDC
mode" ever apply to the non-CDC mode ? If not another reason to move CDC stuff
to a different SchemaTransform.
##########
CHANGES.md:
##########
@@ -102,6 +102,7 @@
* ClickHouseIO: support writing `Decimal(P, S)` / `Decimal32/64/128/256`
columns (Java) ([#39840](https://github.com/apache/beam/issues/39840)).
* SolaceIO now supports reading and writing user properties (message metadata)
(Java) ([#40099](https://github.com/apache/beam/issues/40099)).
* [IcebergIO] AddFiles (`IcebergAddFiles` in YAML) can evolve the table schema
before registering files, with `schema_evolution_options`, `required_columns`,
`incompatible_schema_handling` and `unverifiable_file_handling` (Java/YAML,
batch only) ([#40144](https://github.com/apache/beam/issues/40144)).
+* [IcebergIO] Added batch and streaming CDC writes that applies
INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE changes to Iceberg V2+ tables by
primary key. Invoke with `IcebergIO.writeCdcRows` (Java) or by setting the
`cdc` option of Managed `ICEBERG` write (Java, Python, YAML)
([#X](https://github.com/apache/beam/issues/X)).
Review Comment:
`s/#X/#39979`
##########
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergSchemaTransformTranslationTest.java:
##########
@@ -98,6 +99,27 @@ public class IcebergSchemaTransformTranslationTest {
.withFieldValue("keep", Collections.singletonList("str"))
Review Comment:
Also add tests for `validateModeOptions` ?
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -35,28 +44,37 @@
import org.apache.beam.sdk.schemas.transforms.SchemaTransform;
import org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider;
import org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider;
+import org.apache.beam.sdk.schemas.transforms.providers.ErrorHandling;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.transforms.SimpleFunction;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.PCollectionRowTuple;
import org.apache.beam.sdk.values.Row;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
+import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.DistributionMode;
import org.apache.iceberg.FileFormat;
import org.checkerframework.checker.nullness.qual.Nullable;
import org.joda.time.Duration;
/**
- * SchemaTransform implementation for {@link IcebergIO#writeRows}. Writes Beam
Rows to Iceberg and
- * outputs a {@code PCollection<Row>} representing snapshots created in the
process.
+ * SchemaTransform implementation for {@link IcebergIO#writeRows} and, when
the {@code cdc} block is
+ * set, {@link IcebergIO#writeCdcRows}. Outputs a {@code PCollection<Row>}
representing the
+ * snapshots created in the process; CDC writes add a {@code dead_letter}
output of late records.
*/
@AutoService(SchemaTransformProvider.class)
public class IcebergWriteSchemaTransformProvider
extends TypedSchemaTransformProvider<Configuration> {
static final String INPUT_TAG = "input";
static final String SNAPSHOTS_TAG = "snapshots";
+ static final String DEAD_LETTER_TAG = "dead_letter";
+ static final String ERRORS_TAG = "errors";
+
+ /** The default sequence-number column: what the CDC read source emits. */
+ private static final String DEFAULT_SEQUENCE_NUMBER_COLUMN =
Review Comment:
We should make following public and reuse instead of redefining here.
https://github.com/apache/beam/blob/c462e49b115b2c7f2f9159b193d2547c00c6f157/sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/cdc/sink/CdcWriteConfig.java#L34
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -230,6 +348,121 @@ public IcebergCatalogConfig getIcebergCatalog() {
.setConfigProperties(getConfigProperties())
.build();
}
+
+ enum Mode {
+ APPEND,
+ CDC
+ }
+
+ /** The write mode this configuration selects. */
+ Mode mode() {
+ return getCdc() == null ? Mode.APPEND : Mode.CDC;
+ }
+
+ /**
+ * Options that only one mode supports today; everything else applies to
both. Once an option
+ * becomes available to the other mode, just remove it from here.
+ */
+ private static final ImmutableMap<String, Mode> SUPPORTED_MODES =
+ ImmutableMap.<String, Mode>builder()
+ .put("equality_columns", Mode.CDC)
+ .put("num_shards", Mode.CDC)
+ .put("shards_per_partition", Mode.CDC)
+ .put("allowed_lateness_seconds", Mode.CDC)
+ .put("sink_id", Mode.CDC)
+ .put("token_heartbeat_seconds", Mode.CDC)
+ .put("snapshot_properties", Mode.CDC)
+ .put("error_handling", Mode.CDC)
+ .put("sorter_memory_mb", Mode.CDC)
+ .put("direct_write_byte_limit", Mode.APPEND)
+ .put("distribution_mode", Mode.APPEND)
+ .put("autosharding", Mode.APPEND)
+ .put("write_properties", Mode.APPEND)
+ .put("using_side_input_table_cache", Mode.APPEND)
+ .put("table_refresh_interval_seconds", Mode.APPEND)
+ .put("maximum_cache_size", Mode.APPEND)
+ .put("polling_buckets", Mode.APPEND)
+ .build();
+
+ /** Rejects every set option that the selected mode does not support. */
+ void validateModeOptions() {
+ Mode mode = mode();
+ Map<String, @Nullable Object> values = new LinkedHashMap<>();
Review Comment:
We should derive this from the `Configuration` instead of maintaining a
separate local map here and a `SUPPORTED_MODES` map above so that the two
things don't drift in the future.
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/IcebergWriteSchemaTransformProvider.java:
##########
@@ -230,6 +348,121 @@ public IcebergCatalogConfig getIcebergCatalog() {
.setConfigProperties(getConfigProperties())
.build();
}
+
+ enum Mode {
+ APPEND,
+ CDC
+ }
+
+ /** The write mode this configuration selects. */
+ Mode mode() {
+ return getCdc() == null ? Mode.APPEND : Mode.CDC;
+ }
+
+ /**
+ * Options that only one mode supports today; everything else applies to
both. Once an option
+ * becomes available to the other mode, just remove it from here.
+ */
+ private static final ImmutableMap<String, Mode> SUPPORTED_MODES =
+ ImmutableMap.<String, Mode>builder()
+ .put("equality_columns", Mode.CDC)
+ .put("num_shards", Mode.CDC)
+ .put("shards_per_partition", Mode.CDC)
+ .put("allowed_lateness_seconds", Mode.CDC)
+ .put("sink_id", Mode.CDC)
+ .put("token_heartbeat_seconds", Mode.CDC)
+ .put("snapshot_properties", Mode.CDC)
+ .put("error_handling", Mode.CDC)
+ .put("sorter_memory_mb", Mode.CDC)
+ .put("direct_write_byte_limit", Mode.APPEND)
+ .put("distribution_mode", Mode.APPEND)
+ .put("autosharding", Mode.APPEND)
+ .put("write_properties", Mode.APPEND)
+ .put("using_side_input_table_cache", Mode.APPEND)
+ .put("table_refresh_interval_seconds", Mode.APPEND)
+ .put("maximum_cache_size", Mode.APPEND)
+ .put("polling_buckets", Mode.APPEND)
+ .build();
+
+ /** Rejects every set option that the selected mode does not support. */
+ void validateModeOptions() {
+ Mode mode = mode();
+ Map<String, @Nullable Object> values = new LinkedHashMap<>();
+ values.put("equality_columns", getEqualityColumns());
+ values.put("num_shards", getNumShards());
+ values.put("shards_per_partition", getShardsPerPartition());
+ values.put("allowed_lateness_seconds", getAllowedLatenessSeconds());
+ values.put("sink_id", getSinkId());
+ values.put("token_heartbeat_seconds", getTokenHeartbeatSeconds());
+ values.put("snapshot_properties", getSnapshotProperties());
+ values.put("error_handling", getErrorHandling());
+ values.put("sorter_memory_mb", getSorterMemoryMb());
+ values.put("direct_write_byte_limit", getDirectWriteByteLimit());
+ values.put("distribution_mode", getDistributionMode());
+ values.put("autosharding", getAutosharding());
+ values.put("write_properties", getWriteProperties());
+ values.put("using_side_input_table_cache",
getUsingSideInputTableCache());
+ values.put("table_refresh_interval_seconds",
getTableRefreshIntervalSeconds());
+ values.put("maximum_cache_size", getMaximumCacheSize());
+ values.put("polling_buckets", getPollingBuckets());
+ List<String> invalidOptions = new ArrayList<>();
+ for (Map.Entry<String, Mode> option : SUPPORTED_MODES.entrySet()) {
+ String name = option.getKey();
+ if (values.get(name) == null || option.getValue() == mode) {
+ continue;
+ }
+ invalidOptions.add(name);
+ }
+ if (!invalidOptions.isEmpty()) {
+ throw new IllegalArgumentException(
+ String.format(
+ "The following options are not supported by %s writes yet: %s",
Review Comment:
Probably give a better hint regarding potential action to take in this error
message based on the property. For example,
* set the cdc block
* these aren't implemented for CDC yet
--
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]