This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new ce196d8faa [Fix][Connector-V2] Close Parquet and ORC writers when
rolling files (#11832)
ce196d8faa is described below
commit ce196d8faa098d4e166bcb12d28bf9a0a2bad693
Author: xuepeng <[email protected]>
AuthorDate: Wed Aug 19 13:36:15 2026 +0800
[Fix][Connector-V2] Close Parquet and ORC writers when rolling files
(#11832)
---
.../file/sink/writer/OrcWriteStrategy.java | 6 +++
.../file/sink/writer/ParquetWriteStrategy.java | 6 +++
.../file/writer/OrcWriteStrategyTest.java | 57 ++++++++++++++++++++++
.../file/writer/ParquetWriteStrategyTest.java | 57 ++++++++++++++++++++++
4 files changed, 126 insertions(+)
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/OrcWriteStrategy.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/OrcWriteStrategy.java
index 50bb69eea8..48424d66ed 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/OrcWriteStrategy.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/OrcWriteStrategy.java
@@ -96,6 +96,12 @@ public class OrcWriteStrategy extends
AbstractWriteStrategy<Writer> {
}
}
+ @Override
+ public synchronized void newFilePart() {
+ finishAndCloseFile();
+ super.newFilePart();
+ }
+
@Override
public void finishAndCloseFile() {
List<FileConnectorException> closeErrors = new ArrayList<>();
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
index 43bd1bc8c1..e72436338f 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/ParquetWriteStrategy.java
@@ -140,6 +140,12 @@ public class ParquetWriteStrategy extends
AbstractWriteStrategy<ParquetWriter<Ge
}
}
+ @Override
+ public synchronized void newFilePart() {
+ finishAndCloseFile();
+ super.newFilePart();
+ }
+
@Override
public void finishAndCloseFile() {
List<FileConnectorException> closeErrors = new ArrayList<>();
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/OrcWriteStrategyTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/OrcWriteStrategyTest.java
index d75c04ea03..aaca9476c3 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/OrcWriteStrategyTest.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/OrcWriteStrategyTest.java
@@ -49,6 +49,49 @@ public class OrcWriteStrategyTest {
private static final String TMP_PATH =
"file:///tmp/seatunnel/orc/batch/test";
private static final int ORC_WRITE_NUMBER = 2000;
+ @DisabledOnOs(OS.WINDOWS)
+ @Test
+ public void testCloseWriterWhenRollingFile() throws Exception {
+ String tmpPath = "file:///tmp/seatunnel/orc/rolling/test-" +
System.nanoTime();
+ Map<String, Object> writeConfig = new HashMap<>();
+ writeConfig.put("tmp_path", tmpPath);
+ writeConfig.put("path", "file:///tmp/seatunnel/orc/rolling");
+ writeConfig.put("file_format_type", FileFormat.ORC.name());
+ writeConfig.put("batch_size", 1);
+
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"value"}, new SeaTunnelDataType[]
{BasicType.STRING_TYPE});
+ OrcWriteStrategy writeStrategy =
+ new OrcWriteStrategy(
+ new
FileSinkConfig(ReadonlyConfig.fromMap(writeConfig), rowType));
+ LocalFileSystemConf.LocalConf hadoopConf =
+ new LocalFileSystemConf.LocalConf(FS_DEFAULT_NAME_DEFAULT);
+ writeStrategy.setCatalogTable(
+ CatalogTableUtil.getCatalogTable("test", null, null, "test",
rowType));
+ writeStrategy.init(hadoopConf, "rolling-test", "rolling-test", 0);
+ writeStrategy.beginTransaction(1L);
+
+ writeStrategy.write(new SeaTunnelRow(new Object[] {"first"}));
+ writeStrategy.write(new SeaTunnelRow(new Object[] {"second"}));
+
+ OrcReadStrategy readStrategy = new OrcReadStrategy();
+ readStrategy.init(hadoopConf);
+ List<String> readFiles = readStrategy.getFileNamesByPath(tmpPath);
+ Assertions.assertEquals(1, readFiles.size());
+ String rolledFile = readFiles.get(0);
+ Assertions.assertTrue(rolledFile.endsWith("_0.orc"));
+ readStrategy.getSeaTunnelRowTypeInfo(rolledFile);
+ List<String> values = new ArrayList<>();
+ readStrategy.read(rolledFile, "test", collectorForFirstField(values));
+ Assertions.assertEquals(java.util.Arrays.asList("first"), values);
+
+ writeStrategy.finishAndCloseFile();
+ writeStrategy.close();
+ Assertions.assertEquals(2,
readStrategy.getFileNamesByPath(tmpPath).size());
+ readStrategy.close();
+ }
+
@DisabledOnOs(OS.WINDOWS)
@Test
public void testOrcWriteWithBatch() throws Exception {
@@ -106,4 +149,18 @@ public class OrcWriteStrategyTest {
Assertions.assertEquals(ORC_WRITE_NUMBER, readRows.size());
readStrategy.close();
}
+
+ private static Collector<SeaTunnelRow> collectorForFirstField(List<String>
values) {
+ return new Collector<SeaTunnelRow>() {
+ @Override
+ public void collect(SeaTunnelRow record) {
+ values.add((String) record.getField(0));
+ }
+
+ @Override
+ public Object getCheckpointLock() {
+ return null;
+ }
+ };
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
index eea84414c9..7c1a71eef2 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/ParquetWriteStrategyTest.java
@@ -60,6 +60,49 @@ import static
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_DEFAULT_NAME
public class ParquetWriteStrategyTest {
private static final String TMP_PATH =
"file:///tmp/seatunnel/parquet/int96/test";
+ @DisabledOnOs(OS.WINDOWS)
+ @Test
+ public void testCloseWriterWhenRollingFile() throws Exception {
+ String tmpPath = "file:///tmp/seatunnel/parquet/rolling/test-" +
System.nanoTime();
+ Map<String, Object> writeConfig = new HashMap<>();
+ writeConfig.put("tmp_path", tmpPath);
+ writeConfig.put("path", "file:///tmp/seatunnel/parquet/rolling");
+ writeConfig.put("file_format_type", FileFormat.PARQUET.name());
+ writeConfig.put("batch_size", 1);
+
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"value"}, new SeaTunnelDataType[]
{BasicType.STRING_TYPE});
+ ParquetWriteStrategy writeStrategy =
+ new ParquetWriteStrategy(
+ new
FileSinkConfig(ReadonlyConfig.fromMap(writeConfig), rowType));
+ LocalFileSystemConf.LocalConf hadoopConf =
+ new LocalFileSystemConf.LocalConf(FS_DEFAULT_NAME_DEFAULT);
+ writeStrategy.setCatalogTable(
+ CatalogTableUtil.getCatalogTable("test", null, null, "test",
rowType));
+ writeStrategy.init(hadoopConf, "rolling-test", "rolling-test", 0);
+ writeStrategy.beginTransaction(1L);
+
+ writeStrategy.write(new SeaTunnelRow(new Object[] {"first"}));
+ writeStrategy.write(new SeaTunnelRow(new Object[] {"second"}));
+
+ ParquetReadStrategy readStrategy = new ParquetReadStrategy();
+ readStrategy.init(hadoopConf);
+ List<String> readFiles = readStrategy.getFileNamesByPath(tmpPath);
+ Assertions.assertEquals(1, readFiles.size());
+ String rolledFile = readFiles.get(0);
+ Assertions.assertTrue(rolledFile.endsWith("_0.parquet"));
+ readStrategy.getSeaTunnelRowTypeInfo(rolledFile);
+ List<String> values = new ArrayList<>();
+ readStrategy.read(rolledFile, "test", collectorForFirstField(values));
+ Assertions.assertEquals(Arrays.asList("first"), values);
+
+ writeStrategy.finishAndCloseFile();
+ writeStrategy.close();
+ Assertions.assertEquals(2,
readStrategy.getFileNamesByPath(tmpPath).size());
+ readStrategy.close();
+ }
+
@DisabledOnOs(OS.WINDOWS)
@Test
public void testParquetWriteInt96() throws Exception {
@@ -193,4 +236,18 @@ public class ParquetWriteStrategyTest {
}
readStrategy.close();
}
+
+ private static Collector<SeaTunnelRow> collectorForFirstField(List<String>
values) {
+ return new Collector<SeaTunnelRow>() {
+ @Override
+ public void collect(SeaTunnelRow record) {
+ values.add((String) record.getField(0));
+ }
+
+ @Override
+ public Object getCheckpointLock() {
+ return null;
+ }
+ };
+ }
}