ahmedabu98 commented on code in PR #40167:
URL: https://github.com/apache/beam/pull/40167#discussion_r4066824317
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CommitSchemaUnion.java:
##########
@@ -153,6 +179,33 @@ private static String truncate(String json) {
+ " chars truncated)";
}
+ /** Where each distinct file schema of the window stands while the plan is
worked out. */
+ private static final class Verdicts {
+ final List<SchemaToMerge> toMerge = new ArrayList<>();
+ final List<IncompatibleSchema> incompatible = new ArrayList<>();
+
+ void refuse(CollectDistinctSchemas.SchemaGroup group, String reason) {
+ incompatible.add(new IncompatibleSchema(group.getSchemaJson(),
group.getFiles(), reason));
+ }
+
+ void refuse(Conflict conflict) {
+ toMerge.remove(conflict.schema);
+ incompatible.add(
+ new IncompatibleSchema(conflict.schema.json, conflict.schema.files,
conflict.reason));
+ }
+ }
+
+ /** A schema that passed classification but cannot be staged after the ones
before it. */
+ private static final class Conflict {
Review Comment:
> but cannot be staged after the ones before it.
Can you clarify what this means?
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CommitSchemaUnion.java:
##########
@@ -191,98 +255,293 @@ static long commit(
Thread.currentThread().interrupt();
throw e;
}
- LOG.info(
- "Schema commit attempt {}/{} for {} failed; reloading and
rebuilding",
- attempt,
- MAX_ATTEMPTS,
- tableId,
- e);
}
}
}
- private static long commitOnce(
+ /**
+ * What one schema commit would do, as data: which distinct file schemas it
merges into the table,
+ * which it refuses and why, and the schema the table ends up with. Computed
on scratch
+ * transactions; the commit executes it and the dry run reports it. Adding a
check here is the
+ * only way to add one, so the two cannot drift.
+ */
+ abstract static class Plan {
Review Comment:
nit: Worth putting all these `Plan` related classes and logic in a separate
file
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/CommitSchemaUnion.java:
##########
@@ -191,98 +255,293 @@ static long commit(
Thread.currentThread().interrupt();
throw e;
}
- LOG.info(
- "Schema commit attempt {}/{} for {} failed; reloading and
rebuilding",
- attempt,
- MAX_ATTEMPTS,
- tableId,
- e);
}
}
}
- private static long commitOnce(
+ /**
+ * What one schema commit would do, as data: which distinct file schemas it
merges into the table,
+ * which it refuses and why, and the schema the table ends up with. Computed
on scratch
+ * transactions; the commit executes it and the dry run reports it. Adding a
check here is the
+ * only way to add one, so the two cannot drift.
+ */
+ abstract static class Plan {
+ final List<SchemaToMerge> schemasToMerge;
+ final List<IncompatibleSchema> incompatibleSchemas;
+
+ /** The union to replay, or the schema to create the table with; null when
nothing changes. */
+ final @Nullable Schema newSchema;
+
+ /** Problems the configuration raises against the planned schema, as
messages. */
+ final List<String> configProblems;
+
+ private Plan(Verdicts verdicts, @Nullable Schema newSchema, List<String>
configProblems) {
+ this.schemasToMerge = verdicts.toMerge;
+ this.incompatibleSchemas = verdicts.incompatible;
+ this.newSchema = newSchema;
+ this.configProblems = configProblems;
+ }
+
+ /** Null when the schema is incompatible, or when the existing table
already covers it. */
+ @Nullable SchemaToMerge toMerge(String schemaJson) {
+ for (SchemaToMerge item : schemasToMerge) {
+ if (item.json.equals(schemaJson)) {
+ return item;
+ }
+ }
+ return null;
+ }
+
+ /** Null when the schema is not incompatible. */
+ @Nullable String incompatibleReason(String schemaJson) {
Review Comment:
javadoc seems contradictory to the method?
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/AddFilesSchemaTransformProvider.java:
##########
@@ -145,6 +145,25 @@ public static Builder builder() {
+ " Requires schema_evolution_options.")
public abstract @Nullable List<String> getRequiredColumns();
+ @SchemaFieldDescription(
+ "When true, nothing is committed or registered: the transform reads
the files' schemas"
+ + " and emits a dry_run_report output. Its rows are told apart by
row_type: schema"
+ + " (one per distinct file schema: files, changes a real run would
make, whether they"
+ + " are allowed and why not), create (the table a real run would
create), unreadable,"
Review Comment:
Extremely long and convoluted for a public facing description. Can you
please trim it down and simplify for readability?
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/DryRunReport.java:
##########
@@ -0,0 +1,450 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.iceberg;
+
+import static org.apache.beam.sdk.metrics.Metrics.counter;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Locale;
+import
org.apache.beam.sdk.io.iceberg.SchemaEvolutionConfig.IncompatibleSchemaHandling;
+import org.apache.beam.sdk.metrics.Counter;
+import org.apache.beam.sdk.schemas.Schema;
+import org.apache.beam.sdk.transforms.DoFn;
+import org.apache.beam.sdk.values.Row;
+import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing;
+import org.apache.iceberg.SchemaParser;
+import org.apache.iceberg.catalog.Catalog;
+import org.apache.iceberg.catalog.TableIdentifier;
+import org.apache.iceberg.types.Type;
+import org.checkerframework.checker.nullness.qual.MonotonicNonNull;
+import org.checkerframework.checker.nullness.qual.Nullable;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * Dry run of the schema pre-pass: computes the {@link CommitSchemaUnion#plan}
a real commit would,
+ * on scratch transactions only, and reports it without committing or
registering anything. Planning
+ * the real fold rather than classifying each schema alone is what makes a
schema that is fine
+ * against the table but conflicts with another schema of the input come out
with the blame a real
+ * run would assign.
+ *
+ * <p>Rows, by {@code row_type}: {@code schema}, one per distinct schema, with
the changes a real
+ * run would make on an existing table ({@code schema_key} is a short hash to
group by); {@code
+ * create} for the table a real run would create, once, from the union of the
allowed schemas;
+ * {@code unreadable} for files whose schema could not be read; {@code
unchecked} for ORC and Avro
+ * files the per-file checks cannot verify; and {@code summary}. The summary
is {@code allowed} when
+ * no schema is incompatible and the configuration raises no problem; its
{@code changes} hold the
+ * totals line and any table-level change a real run would make without a
schema change, such as
+ * regenerating the name mapping. {@code would_create_table} is false whenever
a real run would not
+ * create the table, including when it would fail first.
+ *
+ * <p>The counters and the totals count checked Parquet files by their
schema's verdict; unchecked
+ * files count separately even when ACCEPT registers them, so the unchecked
row can be {@code
+ * allowed} while its files are outside {@code numDryRunFilesAllowed}. Pin
evidence is per file (the
+ * footer's null counts), so a pin violation or an unproven pin is not
predicted here.
+ */
+class DryRunReport extends DoFn<List<CollectDistinctSchemas.SchemaGroup>, Row>
{
+ private static final Logger LOG =
LoggerFactory.getLogger(DryRunReport.class);
+
+ static final String FILES_ALLOWED_COUNTER = "numDryRunFilesAllowed";
+ static final String FILES_INCOMPATIBLE_COUNTER =
"numDryRunFilesIncompatible";
+ static final String FILES_UNREADABLE_COUNTER = "numDryRunFilesUnreadable";
+ static final String FILES_UNCHECKED_COUNTER = "numDryRunFilesUnchecked";
+ static final String CONFIG_PROBLEMS_COUNTER = "numDryRunConfigProblems";
+
+ private static final Counter numFilesAllowed = counter(DryRunReport.class,
FILES_ALLOWED_COUNTER);
+ private static final Counter numFilesIncompatible =
+ counter(DryRunReport.class, FILES_INCOMPATIBLE_COUNTER);
+ private static final Counter numFilesUnreadable =
+ counter(DryRunReport.class, FILES_UNREADABLE_COUNTER);
+ private static final Counter numFilesUnchecked =
+ counter(DryRunReport.class, FILES_UNCHECKED_COUNTER);
+ private static final Counter numConfigProblems =
+ counter(DryRunReport.class, CONFIG_PROBLEMS_COUNTER);
+
+ static final Schema REPORT_SCHEMA =
+ Schema.builder()
+ .addStringField("row_type")
+ .addStringField("schema_key")
+ .addStringField("schema")
+ .addInt64Field("num_files")
+ .addArrayField("changes", Schema.FieldType.STRING)
+ .addBooleanField("allowed")
+ .addStringField("reason")
+ .addBooleanField("would_create_table")
+ .build();
+
+ static final String SCHEMA_ROW = "schema";
+ static final String CREATE_ROW = "create";
+ static final String UNREADABLE_ROW = "unreadable";
+ static final String UNCHECKED_ROW = "unchecked";
+ static final String SUMMARY_ROW = "summary";
+
+ static final String NAME_MAPPING_CHANGE =
+ "regenerate the name mapping property to cover the schema";
+
+ /** Wide inputs would otherwise put every column of every schema into one
log entry. */
+ private static final int MAX_RENDERED_ROWS = 50;
+
+ private static final int MAX_RENDERED_CHANGES = 10;
+
+ static final String UNREADABLE_REASON =
+ "the file schema could not be read (unknown format, or an unreadable
footer);"
+ + " a real run routes these files to the error output";
+
+ static final String UNCHECKED_REJECTED_REASON =
+ "ORC and Avro files cannot be checked; a real run routes these files to
the error output"
+ + " (UnverifiableFileHandling.REJECT)";
+
+ static final String UNCHECKED_ACCEPTED_REASON =
+ "ORC and Avro files cannot be checked; a real run registers these files
unchecked"
+ + " (UnverifiableFileHandling.ACCEPT)";
+
+ private final IcebergCatalogConfig catalogConfig;
+ private final String identifier;
+ private final CommitSchemaUnion.Settings settings;
+ private transient @MonotonicNonNull Catalog catalog;
+
+ DryRunReport(
+ IcebergCatalogConfig catalogConfig, String identifier,
CommitSchemaUnion.Settings settings) {
+ this.catalogConfig = catalogConfig;
+ this.identifier = identifier;
+ this.settings = settings;
+ }
+
+ /** One report row before it is a Row: the summary and the rendered log read
these. */
+ private static final class Line {
+ final String rowType;
+ final String schemaKey;
+ final String schema;
+ final long files;
+ final List<String> changes;
+ final boolean allowed;
+ final String reason;
+
+ Line(
+ String rowType,
+ String schemaKey,
+ String schema,
+ long files,
+ List<String> changes,
+ boolean allowed,
+ String reason) {
+ this.rowType = rowType;
+ this.schemaKey = schemaKey;
+ this.schema = schema;
+ this.files = files;
+ this.changes = changes;
+ this.allowed = allowed;
+ this.reason = reason;
+ }
+
+ Row toRow(boolean wouldCreateTable) {
+ return Row.withSchema(REPORT_SCHEMA)
+ .withFieldValue("row_type", rowType)
+ .withFieldValue("schema_key", schemaKey)
+ .withFieldValue("schema", schema)
+ .withFieldValue("num_files", files)
+ .withFieldValue("changes", changes)
+ .withFieldValue("allowed", allowed)
+ .withFieldValue("reason", reason)
+ .withFieldValue("would_create_table", wouldCreateTable)
+ .build();
+ }
+ }
+
+ /** The window's schema groups: the ones the plan sees, and the files that
contribute none. */
+ private static final class Input {
+ final List<CollectDistinctSchemas.SchemaGroup> readable = new
ArrayList<>();
+ long unreadableFiles;
+ long uncheckedFiles;
+
+ static Input of(List<CollectDistinctSchemas.SchemaGroup> schemas) {
+ Input input = new Input();
+ for (CollectDistinctSchemas.SchemaGroup group : schemas) {
+ if (group.getSchemaJson().equals(ReadFooterSchema.UNREADABLE_KEY)) {
+ input.unreadableFiles += group.getFiles();
+ } else if
(group.getSchemaJson().equals(ReadFooterSchema.UNCHECKED_FORMAT_KEY)) {
+ input.uncheckedFiles += group.getFiles();
+ } else {
+ input.readable.add(group);
+ }
+ }
+ return input;
+ }
+ }
+
+ private static final class Totals {
+ int allowedSchemas;
+ long allowedFiles;
+ int incompatibleSchemas;
+ long incompatibleFiles;
+
+ static Totals of(List<Line> schemaLines) {
+ Totals totals = new Totals();
+ for (Line line : schemaLines) {
+ if (line.allowed) {
+ totals.allowedSchemas++;
+ totals.allowedFiles += line.files;
+ } else {
+ totals.incompatibleSchemas++;
+ totals.incompatibleFiles += line.files;
+ }
+ }
+ return totals;
+ }
+ }
+
+ @ProcessElement
+ public void process(
+ @Element List<CollectDistinctSchemas.SchemaGroup> schemas,
OutputReceiver<Row> out) {
+ if (catalog == null) {
+ catalog = catalogConfig.catalog();
+ }
+ TableIdentifier tableId = IcebergUtils.parseTableIdentifier(identifier);
+ Input input = Input.of(schemas);
+ CommitSchemaUnion.Plan plan =
+ CommitSchemaUnion.plan(catalog, tableId, input.readable, settings);
+
+ List<Line> lines = new ArrayList<>();
+ for (CollectDistinctSchemas.SchemaGroup group : input.readable) {
+ lines.add(schemaLine(group, plan));
+ }
+ Totals totals = Totals.of(lines);
+
+ boolean repairsNameMapping =
+ plan instanceof CommitSchemaUnion.EvolutionPlan
+ && ((CommitSchemaUnion.EvolutionPlan) plan).repairsNameMapping;
+ CommitSchemaUnion.@Nullable CreationPlan creation = null;
+ if (plan instanceof CommitSchemaUnion.CreationPlan) {
+ creation = (CommitSchemaUnion.CreationPlan) plan;
+ }
+ @Nullable String creationProblem = creation == null ? null :
creation.problem;
+ boolean wouldFail =
+ creationProblem != null
+ || (settings.handling == IncompatibleSchemaHandling.FAIL_PIPELINE
+ && (totals.incompatibleSchemas > 0 ||
!plan.configProblems.isEmpty()));
+ boolean wouldCreateTable = creation != null && creation.canCreate() &&
!wouldFail;
+
+ if (creation != null && creation.newSchema != null) {
+ lines.add(
+ createLine(creation.newSchema, totals.allowedFiles,
wouldCreateTable, creationProblem));
+ }
+ if (input.unreadableFiles > 0) {
+ lines.add(unreadableLine(input.unreadableFiles));
+ }
+ if (input.uncheckedFiles > 0) {
+ lines.add(uncheckedLine(input.uncheckedFiles));
+ }
Review Comment:
Would it be better to have one single output row that contains the full
report?
It would let the user expect only one output instead of potentially many
different outputs.
Also not all columns are relevant for the different outputs we currently
have.
Haven't thought too deeply about this though, there could be downsides and
increased complexity to putting everything in one output
--
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]