Schema Evolution in Pathway Data Pipelines: Versioned Metadata and the InnerSchemaField Abstraction
Pathway handles schema evolution through versioned metadata snapshots, Arrow schema-aware column buffers, and the InnerSchemaField abstraction, allowing backward-compatible changes without manual migrations.
Pathway (pathwaycom/pathway) treats schema evolution as a first-class concern, isolating schema management from its runtime execution engine. This architecture enables seamless changes—adding, removing, or modifying fields—while guaranteeing that historic data remains accessible and new data follows updated layouts without manual migration steps.
Versioned Metadata in the Persistence Layer
Pathway’s persistence layer implements versioned metadata to track schema changes across time. Every persisted operator snapshot carries a metadata version generated via MetadataKey::from_components(version, worker_id, rotation_id) in src/persistence/state.rs. When the schema changes, the system writes a new version while keeping old versions readable.
The engine determines which version to use through the compute_threshold_time_and_versions routine. This function selects the latest stable version that possesses full metadata for all workers, ensuring the pipeline never operates on a partial schema. The persistence layer can automatically shrink or split batches to align with the new version without breaking existing data streams.
Schema-Aware Column Buffers
Column buffers in Pathway embed Arrow schemas generated from the current InnerSchemaField map. The AppendOnlyColumnBuffer and SnapshotColumnBuffer structures in src/persistence/cached_object_storage.rs manage these schemas. When schema evolution occurs, the system constructs a new Arrow schema, and the buffer continues writing using the updated layout while preserving earlier rows in their original format.
This design allows the storage layer to maintain multiple physical layouts within the same logical table, with the runtime engine handling the translation between versions during query execution.
Backward Compatibility via InnerSchemaField
The core of Pathway’s schema flexibility lies in the InnerSchemaField abstraction defined in src/connectors/data_format.rs. Users define schemas as a HashMap<String, InnerSchemaField>, where each field stores a Type (including Optional, Any, Json, etc.) and an optional default Value.
Adding a new field with a default value or marking it as Optional makes the schema backward-compatible. When reading older rows that lack the new field, the engine automatically fills missing values with the specified default or null, preventing runtime failures.
Practical Example: Adding a New Field
The following example demonstrates adding a currency field to an existing Order schema:
import pathway as pw
# Original schema
class Order(pw.Schema):
order_id: int
amount: float
orders = pw.io.csv.read("orders_v1/", schema=Order)
# Evolved schema with new field and default value
class OrderV2(pw.Schema):
order_id: int
amount: float
currency: str = "USD" # Optional with default
# Trigger schema evolution
orders = orders.cast(OrderV2)
Behind the scenes, the cast call builds a new HashMap<String, InnerSchemaField> (see InnerSchemaField::new in src/python_api.rs). The persistence layer in src/persistence/state.rs generates a fresh metadata version, while src/persistence/cached_object_storage.rs handles batch shrinking if needed. The Arrow schema supplied to the column buffers now includes the currency field, with existing rows logically inheriting the "USD" default when queried.
Automatic Compatibility Checks and External Registries
Connector Validation
When connecting to external systems, Pathway validates table schemas against the supplied definition. In src/connectors/sqlite.rs and src/connectors/postgres.rs, the engine checks for missing columns. If fields are absent but marked as Optional or possess default values, the pipeline continues execution; otherwise, it raises a clear error before processing begins.
Schema Registry Integration
For organizations using centralized schema management, Pathway fetches Avro or Protobuf schemas from Confluent-compatible Schema Registries. The PySchemaRegistrySettings structure in src/python_api.rs handles this integration using the schema_registry_converter crate. Fetched schemas convert into InnerSchemaField definitions, allowing seamless evolution when the registry version changes, with the same versioned persistence path applying automatically.
Runtime Re-Evaluation and Default Values
During execution, the Rust engine processes new rows using the current schema while older rows retain their original version metadata. When a query accesses a column that did not exist in older schema versions, the engine automatically substitutes the column’s default value (or null for Optional types) without requiring data migration. This logic operates within the engine::timestamp and engine::dataflow::persist modules, ensuring consistent results across heterogeneous data versions.
Summary
- Versioned snapshots: Every operator state carries a
MetadataKeyversion insrc/persistence/state.rs, enabling time-travel reads and safe updates. - Arrow-backed buffers:
AppendOnlyColumnBufferandSnapshotColumnBufferinsrc/persistence/cached_object_storage.rsseparate physical storage layout from logical schema. - Flexible field definitions: The
InnerSchemaFieldabstraction insrc/connectors/data_format.rssupports defaults andOptionaltypes for backward compatibility. - Registry integration:
src/python_api.rsprovidesPySchemaRegistrySettingsfor external Confluent Schema Registry support. - Runtime resolution: The engine resolves missing columns using defaults or nulls, eliminating the need for ETL reprocessing when schemas evolve.
Frequently Asked Questions
How does Pathway handle missing columns when reading old data with a new schema?
When a query accesses historical data that lacks columns defined in the current schema, Pathway resolves the missing values using the default value specified in the InnerSchemaField definition, or null if the type is Optional. This behavior is implemented in the runtime engine and prevents crashes during schema evolution.
Can Pathway integrate with Confluent Schema Registry for Avro schemas?
Yes. Pathway supports Confluent-compatible Schema Registries through the PySchemaRegistrySettings class in src/python_api.rs, which utilizes the schema_registry_converter crate to fetch and convert Avro or Protobuf schemas into internal InnerSchemaField definitions.
What happens to existing data when I add a required field without a default?
If you add a new field without a default value and do not mark it as Optional, Pathway’s connector validation (in src/connectors/sqlite.rs and src/connectors/postgres.rs) will raise an error when encountering existing rows that lack the field. For persisted internal state, the pipeline would fail to resolve the missing value, requiring you to provide a default or mark the field as Optional to maintain backward compatibility.
Where is schema version metadata stored in Pathway?
Schema version metadata is stored in the persistence layer via the MetadataKey structure in src/persistence/state.rs. Each operator snapshot embeds a version number created by MetadataKey::from_components(version, worker_id, rotation_id), allowing the system to track which schema version corresponds to each data batch.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →