Skip to content

Spark 4.2: Implement SupportsSchemaEvolution for automatic schema evolution - #17957

Open
rahulsmahadev wants to merge 6 commits into
apache:mainfrom
rahulsmahadev:spark-4.2-schema-evolution-capability
Open

rahulsmahadev wants to merge 6 commits into
apache:mainfrom
rahulsmahadev:spark-4.2-schema-evolution-capability

Conversation

@rahulsmahadev

Copy link
Copy Markdown
Contributor

Implements Spark 4.2's SupportsSchemaEvolution mix-in on Iceberg's SparkTable, so Iceberg participates in DSv2 automatic schema evolution during writes.

During automatic schema evolution Spark passes each candidate column change to supportsColumnChange(TableChange.ColumnChange) and skips the ones the source reports as unsupported. SparkTable now reports support consistent with what Iceberg's schema update actually applies:

  • add column: only when nullable, without a default, and with a Spark type Iceberg can convert;
  • update type: only when the new type is a primitive and an allowed Iceberg type promotion;
  • update nullability: only when relaxing to nullable;
  • delete column: rejected for identifier fields and for any field whose subtree contains an identifier field (matching SchemaUpdate), and for unknown fields;
  • rename, comment, and position changes are supported.

SparkTable already advertises TableCapability.AUTOMATIC_SCHEMA_EVOLUTION (gated by the write.spark.auto-schema-evolution.enabled table property), which is what makes Spark consult this hook. A unit test is added in TestSparkTable.

public boolean supportsColumnChange(TableChange.ColumnChange change) {
if (change instanceof TableChange.AddColumn) {
TableChange.AddColumn add = (TableChange.AddColumn) change;
return add.isNullable() && add.defaultValue() == null && canConvert(add.dataType());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Spark 4.2 recursively generates schema changes for map keys.

For example, evolving a target map<int, string> from a source map<long, string> produces updateColumnType(["m", "key"], LongType). T

his implementation returns true, but SchemaUpdate.apply() explicitly rejects(as Iceberg don't allow it) updates to map keys (Cannot update map keys), so the write fails during alterTable instead of letting Spark skip the unsupported evolution and attempt a cast.

The same false positive occurs for additions or type changes inside a struct-valued map key. Could we reject changes whose field path targets the map-key subtree and add regression coverage for it?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed, good point supportsColumnChange now rejects map-key changes up front

@szehon-ho

szehon-ho commented Sep 29, 2026 •

Copy link
Copy Markdown
Member

Could we log a warning when this hook is invoked with write.spark.accept-any-schema=true (the deprecated one), deduplicated so it isn't emitted for every column?

Suggested text:

Spark-native schema evolution is being used with write.spark.accept-any-schema=true, which skips Spark’s normal schema checks and casts. For native evolution, use accept-any-schema=false.

I believe it would run both schema evolutions which is weird. We can also add a test for this case?

Comment on lines +229 to +232
Type newType = tryConvert(update.newDataType());
return newType != null
&& newType.isPrimitiveType()
&& TypeUtil.isPromotionAllowed(field.type(), newType.asPrimitiveType());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we reject type updates when the converted Iceberg type does not map back to the requested Spark type, and add an INSERT or MERGE test? For an INT target with a SMALLINT/TINYINT source, ShortType/ByteType converts to Iceberg int, so isPromotionAllowed(int, int) returns true and updateColumn is a no-op. Spark reloads the table, sees the same pending change, and fails with UNSUPPORTED_AUTO_SCHEMA_EVOLUTION_CHANGES instead of casting the source to INT.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good pointt, supportsTypeUpdate now also requires the converted Iceberg type to round-trip back to the requested Spark type

@rahulsmahadev

Copy link
Copy Markdown
Contributor Author

Could we log a warning when this hook is invoked with write.spark.accept-any-schema=true (the deprecated one), deduplicated so it isn't emitted for every column?

Suggested text:

Spark-native schema evolution is being used with write.spark.accept-any-schema=true, which skips Spark’s normal schema checks and casts. For native evolution, use accept-any-schema=false.

I believe it would run both schema evolutions which is weird. We can also add a test for this case?

added

…lution

Implement Spark 4.2's SupportsSchemaEvolution mix-in on SparkTable so Iceberg participates in DSv2 automatic schema evolution during writes. supportsColumnChange reports support consistent with what Iceberg's schema update applies: nullable, no-default column additions of a convertible type; primitive type promotions allowed by Iceberg; relaxing nullability; column deletes except identifier fields and any field whose subtree contains one; and rename, comment, and position changes. Adds a TestSparkTable unit test.
supportsColumnChange returned true for changes targeting a map-key subtree, but Iceberg's SchemaUpdate.apply() rejects map-key updates, so the write failed during alterTable instead of letting Spark skip the unsupported evolution and cast. Return false when a change resolves into a map key, including additions and type changes inside a struct-valued map key.

Add regression coverage in TestSparkTable for map-key rejection and map-value acceptance.
@rahulsmahadev
rahulsmahadev force-pushed the spark-4.2-schema-evolution-capability branch from e91c4bd to 38e2062 Compare September 30, 2026 05:36
private final String branch; // set if table is loaded for specific branch
private final TimeTravel timeTravel; // set if table is loaded for time travel
private final Set<TableCapability> capabilities;
private final AtomicBoolean acceptAnySchemaWarningLogged = new AtomicBoolean(false);

@szehon-ho szehon-ho Oct 1, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is a bit hacky and has some issue:
There are two scope limitations:

  • Reusing the same wrapper suppresses subsequent warnings throughout its lifetime.
  • Reloading creates a fresh flag. Spark reloads after schema evolution, so one query could warn again if it checks remaining candidate changes.

Maybe we should just drop this idea, sorry. As long as we have test coverage that nothing bad happens.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'll look on spark side to log the warning there.


if (change instanceof TableChange.AddColumn) {
TableChange.AddColumn add = (TableChange.AddColumn) change;
return add.isNullable() && add.defaultValue() == null && canConvert(add.dataType());

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please check the converted type against the table’s format version before returning true, including nested types. VariantType and NullType convert successfully, but v1/v2 tables reject their Iceberg types in Schema.checkCompatibility during commit. Tests for v2 rejection and v3 acceptance would cover this.

} else if (change instanceof TableChange.UpdateColumnType) {
return supportsTypeUpdate((TableChange.UpdateColumnType) change);
} else if (change instanceof TableChange.UpdateColumnNullability) {
return ((TableChange.UpdateColumnNullability) change).nullable();

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please also reject making an identifier field, or a struct containing one, nullable. SchemaUpdate.apply() rejects both through Schema.validateIdentifierField. Spark 4.2 currently generates only additions and type updates, so this is a hook-contract edge case that could be covered with a unit test.

Table table, Schema schema, Snapshot snapshot, String branch, TimeTravel timeTravel) {
super(table, schema);
this.schema = schema;
this.mapKeyFieldIds = mapKeyFieldIds(schema);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we lazily compute and cache mapKeyFieldIds when it is first needed? Eager initialization traverses the schema and allocates indexes for every SparkTable load, including reads and tables with schema evolution disabled. The schema is pinned, so the cache needs no invalidation.

… nullability, and lazily cache map-key field IDs
}
}

private Set<Integer> mapKeyFieldIds() {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we make this accessor synchronized, or declare mapKeyFieldIds volatile, so the initialized set is safely published if the same SparkTable is accessed concurrently? With the plain field, another thread is not guaranteed to see all of the HashSet contents. volatile is sufficient if duplicate initialization is acceptable, since the computed set is never mutated afterward.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

made it volatile

@szehon-ho

Copy link
Copy Markdown
Member

Could we add native INSERT WITH SCHEMA EVOLUTION integration tests? The current end-to-end coverage is for MERGE; the existing DataFrameWriter tests exercise legacy merge-schema behavior.

With write.spark.accept-any-schema=false, cover:

  1. Adding a column while casting a SMALLINT source into an existing INT column in the same write. Assert that the target stays INT, the new column is added, existing rows get null for it, and inserted values are correct.
  2. Widening INT to BIGINT using a value above Integer.MAX_VALUE. Assert the resulting type and preservation of old and new values.
  3. Writing map<bigint,string> into map<int,string> with keys that fit INT. Assert that the key type stays INT and the write succeeds with the expected values.

Please parameterize these for by-name and by-position INSERT. Also, the accept-any-schema test currently only calls the hook twice; an actual write with native evolution and legacy merge-schema enabled would verify their interaction. Assert the final schema and rows.


Type newType = tryConvert(update.newDataType());
return newType != null
&& isSupportedAtFormatVersion(newType)

@szehon-ho szehon-ho Oct 2, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

As an optional cleanup, consider keeping isSupportedAtFormatVersion for additions and removing it from type updates. With the current promotion rules (INT to LONG, FLOAT to DOUBLE, decimal precision widening, or unchanged types), an update cannot introduce a type requiring a newer format version when the existing schema is valid. The primitive-type, promotion, and Spark round-trip checks are sufficient here.

    return newType != null
        && newType.isPrimitiveType()
        && TypeUtil.isPromotionAllowed(field.type(), newType.asPrimitiveType())
        && SparkSchemaUtil.convert(newType).equals(update.newDataType());

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ok, the format version check now only runs for additions.

Comment on lines +183 to +188
TableChange.AddColumn add = (TableChange.AddColumn) change;
Type type = tryConvert(add.dataType());
return add.isNullable()
&& add.defaultValue() == null
&& type != null
&& isSupportedAtFormatVersion(type);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Optional cleanup: return early for required columns or columns with defaults, avoiding type conversion for additions that will be rejected anyway.

Suggested change
TableChange.AddColumn add = (TableChange.AddColumn) change;
Type type = tryConvert(add.dataType());
return add.isNullable()
&& add.defaultValue() == null
&& type != null
&& isSupportedAtFormatVersion(type);
TableChange.AddColumn add = (TableChange.AddColumn) change;
if (!add.isNullable() || add.defaultValue() != null) {
return false;
}
Type type = tryConvert(add.dataType());
return type != null && isSupportedAtFormatVersion(type);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

done

… checks

- Make the lazily computed map key IDs volatile so they are safely
  published when a SparkTable is shared across threads.
- Reject required or defaulted column additions before converting the
  type, and drop the format version check from type updates, since an
  allowed promotion cannot need a newer format version.
- Add end-to-end INSERT WITH SCHEMA EVOLUTION tests, by name and by
  position, for adding a column alongside a SMALLINT to INT cast, widening
  INT to BIGINT, writing BIGINT map keys to INT keys, and native evolution
  combined with legacy accept-any-schema and merge-schema.
}

// Iceberg rejects dropping an identifier field or a field whose subtree contains one
return !containsIdentifierField(field);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: Could we return false when deleting a list element or map value itself, and add unit coverage? For example, items.element and m.value currently return true, but SchemaUpdate.apply() rejects both. Spark 4.2 only generates additions and type updates today, so this is a hook-contract edge case.

}

@Override
public boolean supportsColumnChange(TableChange.ColumnChange change) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Optional: Could we move the schema-evolution checks and lazy map-key index into a package-private SparkSchemaEvolution helper in this package, leaving supportsColumnChange as a delegate? This would group the rules and make them easier to unit-test directly. There’s precedent for stateful helpers such as SparkMicroBatchPlanner and SparkConfParser, though keeping these checks here is also reasonable.

@szehon-ho szehon-ho left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

looks good to me, just a minor nit/optional comment

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants