Skip to content

Commit f389005

Browse files
committed
chore(connector-file,transforms-v2): adapt to upstream/dev API changes after rebase
Mechanical fixups so the rebased PR branch compiles and all schema-evolution tests pass on top of latest upstream/dev. No behaviour change in the schema-evolution feature itself. Adapted to upstream API changes: - CanalJsonSerializationSchema / DebeziumJsonSerializationSchema / MaxWellJsonSerializationSchema constructors now require a third arg (mergeUpdateEventFlag) when called with a Charset. onSchemaChanged() in each strategy now passes the field that was already stored in the writer. - FileSinkConfig constructor signature changed to take ReadonlyConfig (was Config). Test sites in FileSchemaEvolutionTest, ParquetWriteStrategyEvolutionTest, and ParquetTypeCoercionTest wrapped with ReadonlyConfig.fromConfig(...). Imports added. - FieldFieldMultiCatalogTransform was renamed to FilterFieldMultiCatalogTransform upstream. ChainTimestampPreservationTest and ProductionPipelineSchemaChangeTest updated accordingly. Test contract updates following the partition_by guard removal in this PR: - testPartitionByWithSchemaEvolutionEnabledThrowsAtConfig replaced with testPartitionByWithSchemaEvolutionEnabledIsAccepted — partition_by + schema_evolution is now supported (writer-local partitionFieldsIndexInRow rebuilds on every ALTER via name-based lookup; drop/rename of partition column is the only case rejected, with an explicit IllegalStateException at rebuild time). Restored applySchemaChange's silent no-op behavior when schema_evolution_enabled=false. The fail-fast variant is left as an open design question for maintainers — current behavior preserves backward compatibility but leaves a latent ClassCastException trap when schema-changes.enabled=true at source + schema_evolution_enabled=false at sink (documented inline). Awaiting maintainer guidance on the preferred UX. All 30 schema-evolution tests pass: - FileSchemaEvolutionTest (19) - ParquetWriteStrategyEvolutionTest (2), ParquetTypeCoercionTest (3) - ChainTimestampPreservationTest, MetadataMultiCatalogSchemaChangeTest, ProductionPipelineSchemaChangeTest, TransformChainLiveAlterTest, SQLMultiCatalogSchemaChangeTest
1 parent 7a1b8fc commit f389005

10 files changed

Lines changed: 43 additions & 39 deletions

File tree

seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/AbstractWriteStrategy.java

Lines changed: 6 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,6 @@
4040
import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSinkOptions;
4141
import org.apache.seatunnel.connectors.seatunnel.file.config.FileFormat;
4242
import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
43-
import org.apache.seatunnel.connectors.seatunnel.file.exception.FileConnectorErrorCode;
4443
import org.apache.seatunnel.connectors.seatunnel.file.exception.FileConnectorException;
4544
import org.apache.seatunnel.connectors.seatunnel.file.hadoop.HadoopFileSystemProxy;
4645
import org.apache.seatunnel.connectors.seatunnel.file.sink.commit.FileCommitInfo;
@@ -193,22 +192,12 @@ public void setCatalogTable(CatalogTable catalogTable) {
193192
@Override
194193
public void applySchemaChange(SchemaChangeEvent event) throws IOException {
195194
if (!fileSinkConfig.isSchemaEvolutionEnabled()) {
196-
// Fail-fast guard: if a real schema change arrives but the sink is configured to
197-
// ignore them, silently swallowing the event leaves sinkColumnsIndexInRow stale.
198-
// The next data row will read row[idx] with stale catalog assumptions and produce a
199-
// confusing ClassCastException several rows later. Throw immediately with an
200-
// actionable message instead — it's a config mismatch, not a runtime fault.
201-
throw new FileConnectorException(
202-
FileConnectorErrorCode.FORMAT_NOT_SUPPORT,
203-
"Received schema-change event ["
204-
+ event.getClass().getSimpleName()
205-
+ "] but schema_evolution_enabled=false on the file sink. "
206-
+ "This combination is unsafe — the sink would keep writing with the "
207-
+ "pre-ALTER schema and corrupt subsequent rows. "
208-
+ "Resolve by either: "
209-
+ "(a) setting schema_evolution_enabled=true on the sink, OR "
210-
+ "(b) setting schema-changes.enabled=false on the CDC source so "
211-
+ "ALTER events are not emitted in the first place.");
195+
// No-op when schema evolution is disabled at the sink. NOTE: if the upstream CDC
196+
// source has schema-changes.enabled=true, this leaves sinkColumnsIndexInRow stale and
197+
// the next data row may produce a ClassCastException several rows later. Open
198+
// discussion with maintainers on whether this should fail-fast instead — the silent
199+
// return is the historical behavior preserved here for backwards compatibility.
200+
return;
212201
}
213202
log.info(
214203
"[FileSchemaEvolution] applying {} — before: rowType={}, sinkColumns={}, indices={}",

seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/CanalJsonWriteStrategy.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,9 @@ public void setCatalogTable(CatalogTable catalogTable) {
7171
protected void onSchemaChanged() {
7272
this.serializationSchema =
7373
new CanalJsonSerializationSchema(
74-
buildSchemaWithRowType(seaTunnelRowType, sinkColumnsIndexInRow), charset);
74+
buildSchemaWithRowType(seaTunnelRowType, sinkColumnsIndexInRow),
75+
charset,
76+
mergeUpdateEventFlag);
7577
}
7678

7779
@Override

seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/DebeziumJsonWriteStrategy.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,9 @@ public void setCatalogTable(CatalogTable catalogTable) {
7171
protected void onSchemaChanged() {
7272
this.serializationSchema =
7373
new DebeziumJsonSerializationSchema(
74-
buildSchemaWithRowType(seaTunnelRowType, sinkColumnsIndexInRow), charset);
74+
buildSchemaWithRowType(seaTunnelRowType, sinkColumnsIndexInRow),
75+
charset,
76+
mergeUpdateEventFlag);
7577
}
7678

7779
@Override

seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/MaxWellJsonWriteStrategy.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,9 @@ public void setCatalogTable(CatalogTable catalogTable) {
7171
protected void onSchemaChanged() {
7272
this.serializationSchema =
7373
new MaxWellJsonSerializationSchema(
74-
buildSchemaWithRowType(seaTunnelRowType, sinkColumnsIndexInRow), charset);
74+
buildSchemaWithRowType(seaTunnelRowType, sinkColumnsIndexInRow),
75+
charset,
76+
mergeUpdateEventFlag);
7577
}
7678

7779
@Override

seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/FileSchemaEvolutionTest.java

Lines changed: 11 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -488,7 +488,8 @@ private static NoOpWriteStrategy buildTypeZooStrategy() {
488488
"path = \"/tmp/test\"\n"
489489
+ "file_format_type = \"parquet\"\n"
490490
+ "schema_evolution_enabled = true");
491-
FileSinkConfig sinkConfig = new FileSinkConfig(config, TYPE_ZOO_ROW_TYPE);
491+
FileSinkConfig sinkConfig =
492+
new FileSinkConfig(ReadonlyConfig.fromConfig(config), TYPE_ZOO_ROW_TYPE);
492493
NoOpWriteStrategy strategy = new NoOpWriteStrategy(sinkConfig);
493494
strategy.setCatalogTable(
494495
org.apache.seatunnel.api.table.catalog.CatalogTableUtil.getCatalogTable(
@@ -654,22 +655,25 @@ public void testBinaryWithSchemaEvolutionEnabledThrowsAtConfig() {
654655
+ "schema_evolution_enabled = true");
655656
Assertions.assertThrows(
656657
FileConnectorException.class,
657-
() -> new FileSinkConfig(config, BASE_ROW_TYPE),
658+
() -> new FileSinkConfig(ReadonlyConfig.fromConfig(config), BASE_ROW_TYPE),
658659
"binary format must reject schema_evolution_enabled=true at config time");
659660
}
660661

661662
@Test
662-
public void testPartitionByWithSchemaEvolutionEnabledThrowsAtConfig() {
663+
public void testPartitionByWithSchemaEvolutionEnabledIsAccepted() {
664+
// partition_by + schema_evolution_enabled is supported: applySchemaChange rebuilds
665+
// partitionFieldsIndexInRow on every ALTER via name-based lookup against the post-ALTER
666+
// row type, mirroring sinkColumnsIndexInRow. Drop/rename of a partition column itself is
667+
// rejected at rebuild time with an explicit IllegalStateException (covered separately).
663668
Config config =
664669
ConfigFactory.parseString(
665670
"path = \"/tmp/test\"\n"
666671
+ "file_format_type = \"text\"\n"
667672
+ "partition_by = [\"age\"]\n"
668673
+ "schema_evolution_enabled = true");
669-
Assertions.assertThrows(
670-
FileConnectorException.class,
671-
() -> new FileSinkConfig(config, BASE_ROW_TYPE),
672-
"partition_by must reject schema_evolution_enabled=true at config time");
674+
Assertions.assertDoesNotThrow(
675+
() -> new FileSinkConfig(ReadonlyConfig.fromConfig(config), BASE_ROW_TYPE),
676+
"partition_by + schema_evolution_enabled=true must be supported");
673677
}
674678

675679
/**

seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetTypeCoercionTest.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory;
2121

22+
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
2223
import org.apache.seatunnel.api.source.Collector;
2324
import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
2425
import org.apache.seatunnel.api.table.type.BasicType;
@@ -214,7 +215,8 @@ private static List<Object[]> writeAndReadBack(
214215
writeConfig.put("file_format_type", FileFormat.PARQUET.name());
215216

216217
FileSinkConfig sinkConfig =
217-
new FileSinkConfig(ConfigFactory.parseMap(writeConfig), rowType);
218+
new FileSinkConfig(
219+
ReadonlyConfig.fromConfig(ConfigFactory.parseMap(writeConfig)), rowType);
218220
ParquetWriteStrategy strategy = new ParquetWriteStrategy(sinkConfig);
219221
ParquetReadStrategyTest.LocalConf hadoopConf =
220222
new ParquetReadStrategyTest.LocalConf(FS_DEFAULT_NAME_DEFAULT);

seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyEvolutionTest.java

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory;
2121

22+
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
2223
import org.apache.seatunnel.api.source.Collector;
2324
import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
2425
import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
@@ -86,7 +87,9 @@ public void testAddColumnRotatesFileAndPreservesAllRows() throws Exception {
8687
writeConfig.put("schema_evolution_enabled", true);
8788

8889
FileSinkConfig sinkConfig =
89-
new FileSinkConfig(ConfigFactory.parseMap(writeConfig), baseRowType);
90+
new FileSinkConfig(
91+
ReadonlyConfig.fromConfig(ConfigFactory.parseMap(writeConfig)),
92+
baseRowType);
9093
ParquetWriteStrategy strategy = new ParquetWriteStrategy(sinkConfig);
9194
ParquetReadStrategyTest.LocalConf hadoopConf =
9295
new ParquetReadStrategyTest.LocalConf(FS_DEFAULT_NAME_DEFAULT);
@@ -187,7 +190,9 @@ public void testDropColumnRotatesFile() throws Exception {
187190
writeConfig.put("schema_evolution_enabled", true);
188191

189192
FileSinkConfig sinkConfig =
190-
new FileSinkConfig(ConfigFactory.parseMap(writeConfig), baseRowType);
193+
new FileSinkConfig(
194+
ReadonlyConfig.fromConfig(ConfigFactory.parseMap(writeConfig)),
195+
baseRowType);
191196
ParquetWriteStrategy strategy = new ParquetWriteStrategy(sinkConfig);
192197
ParquetReadStrategyTest.LocalConf hadoopConf =
193198
new ParquetReadStrategyTest.LocalConf(FS_DEFAULT_NAME_DEFAULT);

seatunnel-transforms-v2/src/main/java/org/apache/seatunnel/transform/common/AbstractMultiCatalogTransform.java

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,8 +19,6 @@
1919

2020
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
2121
import org.apache.seatunnel.api.table.catalog.CatalogTable;
22-
import org.apache.seatunnel.api.table.catalog.TableIdentifier;
23-
import org.apache.seatunnel.api.table.catalog.TableSchema;
2422
import org.apache.seatunnel.api.table.schema.event.AlterTableEvent;
2523
import org.apache.seatunnel.api.table.schema.event.SchemaChangeEvent;
2624
import org.apache.seatunnel.api.table.type.SeaTunnelDataType;

seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/ChainTimestampPreservationTest.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@
3434
import org.apache.seatunnel.api.table.type.LocalTimeType;
3535
import org.apache.seatunnel.api.table.type.MetadataUtil;
3636
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
37-
import org.apache.seatunnel.transform.filter.FieldFieldMultiCatalogTransform;
37+
import org.apache.seatunnel.transform.filter.FilterFieldMultiCatalogTransform;
3838
import org.apache.seatunnel.transform.filterrowkind.FieldRowKindMultiCatalogTransform;
3939
import org.apache.seatunnel.transform.rowkind.RowKindExtractorMultiCatalogTransform;
4040
import org.apache.seatunnel.transform.sql.SQLMultiCatalogFlatMapTransform;
@@ -162,8 +162,8 @@ void chainPreservesLocalDateTimeAtTimestampColumn() {
162162
CatalogTable sqlOut = sql.getProducedCatalogTables().get(0);
163163
Map<String, Object> filterFieldCfg = new HashMap<>();
164164
filterFieldCfg.put("exclude_fields", Arrays.asList("c_delay"));
165-
FieldFieldMultiCatalogTransform filterField =
166-
new FieldFieldMultiCatalogTransform(
165+
FilterFieldMultiCatalogTransform filterField =
166+
new FilterFieldMultiCatalogTransform(
167167
Collections.singletonList(sqlOut), ReadonlyConfig.fromMap(filterFieldCfg));
168168

169169
List<org.apache.seatunnel.api.transform.SeaTunnelTransform<SeaTunnelRow>> chain =

seatunnel-transforms-v2/src/test/java/org/apache/seatunnel/transform/metadata/ProductionPipelineSchemaChangeTest.java

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@
3333
import org.apache.seatunnel.api.table.type.CommonOptions;
3434
import org.apache.seatunnel.api.table.type.MetadataUtil;
3535
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
36-
import org.apache.seatunnel.transform.filter.FieldFieldMultiCatalogTransform;
36+
import org.apache.seatunnel.transform.filter.FilterFieldMultiCatalogTransform;
3737
import org.apache.seatunnel.transform.filterrowkind.FieldRowKindMultiCatalogTransform;
3838
import org.apache.seatunnel.transform.rowkind.RowKindExtractorMultiCatalogTransform;
3939
import org.apache.seatunnel.transform.sql.SQLMultiCatalogFlatMapTransform;
@@ -148,12 +148,12 @@ void productionWrapperPipelinePreservesPostAlterValues() {
148148
new SQLMultiCatalogFlatMapTransform(
149149
Collections.singletonList(rkeOut), ReadonlyConfig.fromMap(sqlCfg));
150150

151-
// 5. FieldFieldMultiCatalogTransform — exclude c_delay
151+
// 5. FilterFieldMultiCatalogTransform — exclude c_delay
152152
CatalogTable sqlOut = sql.getProducedCatalogTables().get(0);
153153
Map<String, Object> filterFieldCfg = new HashMap<>();
154154
filterFieldCfg.put("exclude_fields", Arrays.asList("c_delay"));
155-
FieldFieldMultiCatalogTransform filterField =
156-
new FieldFieldMultiCatalogTransform(
155+
FilterFieldMultiCatalogTransform filterField =
156+
new FilterFieldMultiCatalogTransform(
157157
Collections.singletonList(sqlOut), ReadonlyConfig.fromMap(filterFieldCfg));
158158

159159
// List of all wrappers in chain order — engine iterates in this order during ALTER.

0 commit comments

Comments
 (0)