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();
+        }
+    }
+}

Reply via email to