Skip to main content

Schema Evolution on Write

When write.merge-schema is enabled, Paimon automatically evolves the table schema during write to accommodate new columns in the incoming data, while preserving data integrity.

info

Since the table schema may be updated during writing, catalog caching needs to be disabled to use this feature. Configure spark.sql.catalog.<catalogName>.cache-enabled to false.

How It Evolves the Schema

Three options control how aggressively the schema evolves; each only takes effect when the previous one is enabled:

OptionDescription
write.merge-schemaIf true, evolve the table schema to accept new columns from the incoming data. Existing column types are preserved and incoming values are cast to them; to also widen existing types, enable write.merge-schema.type-widening.
write.merge-schema.type-wideningOnly effective when write.merge-schema is true. If true, widen an existing column type when the incoming data has a wider compatible type (e.g. INT -> BIGINT, DECIMAL precision increase). Lossy changes are still rejected unless write.merge-schema.explicit-cast is also true.
write.merge-schema.explicit-castOnly effective when write.merge-schema.type-widening is true. If true, also allow lossy type changes between compatible types (e.g. BIGINT -> INT, STRING -> DATE).

Examples

DataFrame batch write:

data.write
.format("paimon")
.mode("append")
.option("write.merge-schema", "true")
.saveAsTable("t")

Spark SQL (requires Spark 3.5+ for BY NAME):

SET `spark.paimon.write.merge-schema` = true;

CREATE TABLE t (a INT, b STRING);
INSERT INTO t VALUES (1, '1'), (2, '2');

INSERT INTO t BY NAME SELECT 3 AS a, '3' AS b, 3 AS c;

Streaming write (use the existing table's actual location):

// input is a streaming DataFrame whose columns match or extend table t.
input
.writeStream
.format("paimon")
.option("checkpointLocation", "/path/to/checkpoint")
.option("write.merge-schema", "true")
.start("/path/to/warehouse/default.db/t")

Column Alignment by Write Path

When the source schema doesn't match the target schema exactly, the behavior depends on both write.merge-schema and the write path. For nested struct fields, all byName paths behave the same; at the top level, MERGE INTO * differs from regular byName INSERT because * expansion only references target columns.

Write pathScenariomerge-schema=false (default)merge-schema=true
byName INSERT (INSERT INTO ... BY NAME / saveAsTable / writeTo)Top-level source-extra columnsThrowsEvolved into the target schema
Top-level target columns missing from sourceNULL-filledNULL-filled
Nested struct source-extra fieldsThrowsEvolved into the target schema
Nested struct target-missing fieldsThrowsNULL-filled
MERGE INTO * (UPDATE * / INSERT *)Top-level source-extra columnsSilently dropped (* only covers target columns)Evolved into the target schema
Top-level target columns missing from sourceThrowsUPDATE * preserves current value; INSERT * fills CURRENT_DEFAULT (or NULL when no default)
Nested struct source-extra fieldsThrowsEvolved into the target schema
Nested struct target-missing fieldsThrowsUPDATE * preserves current value; INSERT * fills CURRENT_DEFAULT (or NULL when no default)

Notes:

  • Position-based writes (e.g. INSERT INTO t VALUES (...) without BY NAME) require an exact column count match and don't engage schema evolution; only byName writes are covered above.
  • Top-level target-missing under merge-schema=false for byName INSERT mirrors Spark's INSERT FILL semantics — only nested missing fields throw.
  • Under strict mode (merge-schema=false), nested source-extra fields throw to avoid silent data loss; for MERGE INTO * at the top level, source-extras are silently dropped because * never references them.

For explicit DDL changes, see Alter Tables. For default expressions, see Default Values.