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 deb1622e2b [Improve][Connector-V2][File] Parse the Hadoop
configuration once per subtask instead of once per output file (#11661)
deb1622e2b is described below
commit deb1622e2bf04700e22710daf2493231e5945363
Author: Dongyeon Lee <[email protected]>
AuthorDate: Wed Aug 12 00:19:17 2026 +0900
[Improve][Connector-V2][File] Parse the Hadoop configuration once per
subtask instead of once per output file (#11661)
Co-authored-by: Dongyeon <[email protected]>
Co-authored-by: Claude Opus 5 (1M context) <[email protected]>
---
.../file/sink/writer/AbstractWriteStrategy.java | 36 +++-
.../AbstractWriteStrategyConfigurationTest.java | 191 +++++++++++++++++++++
2 files changed, 226 insertions(+), 1 deletion(-)
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/AbstractWriteStrategy.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/AbstractWriteStrategy.java
index b029c75144..143e2a3eab 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/AbstractWriteStrategy.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/writer/AbstractWriteStrategy.java
@@ -86,6 +86,9 @@ public abstract class AbstractWriteStrategy<T> implements
WriteStrategy<T> {
protected int subTaskIndex;
protected HadoopConf hadoopConf;
protected HadoopFileSystemProxy hadoopFileSystemProxy;
+ /** Template whose resources are already parsed; see {@link
#getConfiguration(HadoopConf)}. */
+ private transient Configuration parsedConfiguration;
+
protected String transactionId;
/** The uuid prefix to make sure same job different file sink will not
conflict. */
protected String uuidPrefix;
@@ -167,13 +170,44 @@ public abstract class AbstractWriteStrategy<T> implements
WriteStrategy<T> {
/**
* use hadoop conf generate hadoop configuration
*
+ * <p>Callers get a fresh, independently mutable Configuration, because
some of them mutate it
+ * ({@code ParquetWriteStrategy#init} sets {@code
AvroWriteSupport.WRITE_FIXED_AS_INT96}). The
+ * expensive part is not the object but the resource load: the first
property access on a new
+ * Configuration parses {@code core-default.xml} and friends. Since this
is called once per
+ * output file, that parse used to be repeated for every file the subtask
wrote. So parse once
+ * and hand out copies — Hadoop's copy constructor clones the
already-loaded properties instead
+ * of re-reading the XML.
+ *
* @param hadoopConf hadoop conf
* @return Configuration
*/
@Override
public Configuration getConfiguration(HadoopConf hadoopConf) {
+ // The cache is only valid for the HadoopConf this strategy was
initialised with.
+ if (hadoopConf != this.hadoopConf) {
+ return buildConfiguration(hadoopConf);
+ }
+ if (parsedConfiguration == null) {
+ parsedConfiguration = buildConfiguration(hadoopConf);
+ }
+ return new Configuration(parsedConfiguration);
+ }
+
+ /**
+ * Does the full, expensive work — the resource parse plus the extra
options — for one
+ * HadoopConf. This is the thing {@link #getConfiguration(HadoopConf)}
caches, so it must stay
+ * free of any per-output-file state.
+ *
+ * <p>Both steps read the same {@code hadoopConf}. That matters because
{@link
+ * HadoopConf#setExtraOptionsForConfiguration} does not only copy {@code
extraOptions} in: it
+ * also decides which keys an {@code hdfs-site.xml} resource may not
overwrite, and those keys
+ * are derived from the conf's own {@code getSchema()}, which every
filesystem subclass
+ * overrides. Taking them from a different conf would let that resource
overwrite the very
+ * properties {@code unsetUnwantedOverwritingProps} exists to protect.
+ */
+ private Configuration buildConfiguration(HadoopConf hadoopConf) {
Configuration configuration = hadoopConf.toConfiguration();
- this.hadoopConf.setExtraOptionsForConfiguration(configuration);
+ hadoopConf.setExtraOptionsForConfiguration(configuration);
return configuration;
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/AbstractWriteStrategyConfigurationTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/AbstractWriteStrategyConfigurationTest.java
new file mode 100644
index 0000000000..7ce26a3b31
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/AbstractWriteStrategyConfigurationTest.java
@@ -0,0 +1,191 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.file.writer;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileFormat;
+import
org.apache.seatunnel.connectors.seatunnel.file.sink.config.FileSinkConfig;
+import
org.apache.seatunnel.connectors.seatunnel.file.sink.writer.ParquetWriteStrategy;
+import org.apache.seatunnel.connectors.seatunnel.file.util.LocalFileSystemConf;
+
+import org.apache.hadoop.conf.Configuration;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+import static
org.apache.hadoop.fs.CommonConfigurationKeysPublic.FS_DEFAULT_NAME_DEFAULT;
+
+/**
+ * getConfiguration is called once per output file, so it caches the parsed
resources. These tests
+ * pin the part callers depend on: every call still yields a Configuration
they may freely mutate.
+ */
+public class AbstractWriteStrategyConfigurationTest {
+
+ private static ParquetWriteStrategy strategy(LocalFileSystemConf.LocalConf
hadoopConf) {
+ Map<String, Object> writeConfig = new HashMap<>();
+ writeConfig.put("tmp_path", "file:///tmp/seatunnel/conf-cache/tmp");
+ writeConfig.put("path", "file:///tmp/seatunnel/conf-cache");
+ writeConfig.put("file_format_type", FileFormat.PARQUET.name());
+
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"f1_text"}, new SeaTunnelDataType[]
{BasicType.STRING_TYPE});
+ ParquetWriteStrategy strategy =
+ new ParquetWriteStrategy(
+ new
FileSinkConfig(ReadonlyConfig.fromMap(writeConfig), rowType));
+ strategy.setCatalogTable(
+ CatalogTableUtil.getCatalogTable("test", null, null, "test",
rowType));
+ strategy.init(hadoopConf, "job1", "job1", 0);
+ return strategy;
+ }
+
+ @Test
+ public void testEachCallReturnsAnIndependentlyMutableConfiguration() {
+ LocalFileSystemConf.LocalConf hadoopConf =
+ new LocalFileSystemConf.LocalConf(FS_DEFAULT_NAME_DEFAULT);
+ ParquetWriteStrategy strategy = strategy(hadoopConf);
+
+ Configuration first = strategy.getConfiguration(hadoopConf);
+ Configuration second = strategy.getConfiguration(hadoopConf);
+
+ Assertions.assertNotSame(first, second);
+
+ // Callers do mutate what they are handed - ParquetWriteStrategy#init
sets
+ // AvroWriteSupport.WRITE_FIXED_AS_INT96 on it - so a mutation must
not be visible
+ // to any later caller.
+ first.set("seatunnel.test.marker", "written-by-first-caller");
+
Assertions.assertNull(strategy.getConfiguration(hadoopConf).get("seatunnel.test.marker"));
+ Assertions.assertNull(second.get("seatunnel.test.marker"));
+ }
+
+ @Test
+ public void testConfigurationCarriesHadoopConfValues() {
+ LocalFileSystemConf.LocalConf hadoopConf =
+ new LocalFileSystemConf.LocalConf(FS_DEFAULT_NAME_DEFAULT);
+ ParquetWriteStrategy strategy = strategy(hadoopConf);
+
+ // Same assertions on the first and a later call: the cached copy must
not lose anything
+ // that toConfiguration()/setExtraOptionsForConfiguration() put there.
+ for (int call = 0; call < 3; call++) {
+ Configuration configuration =
strategy.getConfiguration(hadoopConf);
+ Assertions.assertEquals(
+ FS_DEFAULT_NAME_DEFAULT,
configuration.get("fs.defaultFS"), "call " + call);
+ Assertions.assertEquals(
+ "org.apache.hadoop.fs.LocalFileSystem",
+ configuration.get("fs.file.impl"),
+ "call " + call);
+ Assertions.assertTrue(
+ configuration.getBoolean("fs.file.impl.disable.cache",
false), "call " + call);
+ // A core-default.xml value, i.e. proof the copy carries the
parsed resources and not
+ // just the six properties toConfiguration() sets explicitly.
+ Assertions.assertNotNull(configuration.get("io.file.buffer.size"),
"call " + call);
+ }
+ }
+
+ @Test
+ public void testForeignHadoopConfIsNotServedFromTheCache() {
+ LocalFileSystemConf.LocalConf hadoopConf =
+ new LocalFileSystemConf.LocalConf(FS_DEFAULT_NAME_DEFAULT);
+ ParquetWriteStrategy strategy = strategy(hadoopConf);
+ // Prime the cache.
+ strategy.getConfiguration(hadoopConf);
+
+ LocalFileSystemConf.LocalConf other =
+ new
LocalFileSystemConf.LocalConf("file:///tmp/seatunnel/other");
+ Assertions.assertEquals(
+ "file:///tmp/seatunnel/other",
+ strategy.getConfiguration(other).get("fs.defaultFS"));
+ }
+
+ /**
+ * The point of the cache is that the resource parse happens once. The
assertions above prove
+ * the returned Configuration is correct, which would hold just as well if
nothing were cached
+ * at all - so count the expensive builds directly.
+ */
+ @Test
+ public void testTheExpensiveBuildHappensOncePerHadoopConf() {
+ CountingConf hadoopConf = new CountingConf(FS_DEFAULT_NAME_DEFAULT);
+ ParquetWriteStrategy strategy = strategy(hadoopConf);
+
+ // Deliberately not asserting an absolute number here. Two independent
builds happen during
+ // init - HadoopFileSystemProxy's constructor eagerly builds its own
Configuration, and
+ // ParquetWriteStrategy#init calls getConfiguration - and that split
is init's business, not
+ // this cache's. What this cache promises is that the count stops
growing afterwards.
+ int afterInit = hadoopConf.toConfigurationCalls;
+ Assertions.assertTrue(afterInit >= 1, "init should have built at least
one Configuration");
+
+ for (int i = 0; i < 5; i++) {
+ Assertions.assertNotNull(strategy.getConfiguration(hadoopConf));
+ }
+ Assertions.assertEquals(
+ afterInit,
+ hadoopConf.toConfigurationCalls,
+ "5 further calls must all be served from the cache");
+
+ // A foreign conf is built from itself and must not disturb the cached
template.
+ CountingConf other = new CountingConf("file:///tmp/seatunnel/other");
+ strategy.getConfiguration(other);
+ Assertions.assertEquals(1, other.toConfigurationCalls);
+ Assertions.assertEquals(afterInit, hadoopConf.toConfigurationCalls);
+ }
+
+ /**
+ * A foreign HadoopConf must be configured entirely from itself, never
from the strategy's own
+ * conf. This is easy to get wrong in a way that is hard to notice: the
keys
+ * setExtraOptionsForConfiguration protects against an hdfs-site.xml
overwrite are derived from
+ * getSchema(), and every filesystem subclass overrides that - so mixing
the two confs would let
+ * that resource overwrite the properties the protection exists for.
+ */
+ @Test
+ public void testForeignHadoopConfGetsItsOwnExtraOptions() {
+ LocalFileSystemConf.LocalConf hadoopConf =
+ new LocalFileSystemConf.LocalConf(FS_DEFAULT_NAME_DEFAULT);
+
hadoopConf.setExtraOptions(Collections.singletonMap("seatunnel.test.owner",
"strategy"));
+ ParquetWriteStrategy strategy = strategy(hadoopConf);
+
+ LocalFileSystemConf.LocalConf other =
+ new
LocalFileSystemConf.LocalConf("file:///tmp/seatunnel/other");
+ other.setExtraOptions(Collections.singletonMap("seatunnel.test.owner",
"foreign"));
+
+ Configuration configuration = strategy.getConfiguration(other);
+ Assertions.assertEquals("foreign",
configuration.get("seatunnel.test.owner"));
+ }
+
+ /** Counts how many times the expensive build path actually ran. */
+ private static class CountingConf extends LocalFileSystemConf.LocalConf {
+ private int toConfigurationCalls;
+
+ private CountingConf(String hdfsNameKey) {
+ super(hdfsNameKey);
+ }
+
+ @Override
+ public Configuration toConfiguration() {
+ toConfigurationCalls++;
+ return super.toConfiguration();
+ }
+ }
+}