nielifeng commented on code in PR #11434:
URL: https://github.com/apache/seatunnel/pull/11434#discussion_r3584530532


##########
seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigValidationUtils.java:
##########
@@ -0,0 +1,339 @@
+/*
+ * 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.core.starter.utils;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.common.PluginIdentifier;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.SeaTunnelSink;
+import org.apache.seatunnel.api.sink.SupportMultiTableSink;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.factory.FactoryUtil;
+import org.apache.seatunnel.api.table.factory.TableSinkFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactoryContext;
+import org.apache.seatunnel.api.transform.SeaTunnelTransform;
+import org.apache.seatunnel.common.Constants;
+import org.apache.seatunnel.common.config.CheckResult;
+import org.apache.seatunnel.common.config.TypesafeConfigUtils;
+import org.apache.seatunnel.common.constants.EngineType;
+import org.apache.seatunnel.common.constants.PluginType;
+import org.apache.seatunnel.common.utils.ExceptionUtils;
+import org.apache.seatunnel.common.utils.ReflectionUtils;
+import org.apache.seatunnel.core.starter.exception.ConfigCheckException;
+import org.apache.seatunnel.core.starter.execution.RuntimeEnvironment;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelFactoryDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSinkPluginDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSourcePluginDiscovery;
+
+import scala.Tuple2;
+
+import java.io.Serializable;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.BiConsumer;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_INPUT;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_OUTPUT;
+import static 
org.apache.seatunnel.api.table.factory.FactoryUtil.ensureJobModeMatch;
+
+/** Utility methods for validating SeaTunnel job configuration without 
executing a job. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+public final class ConfigValidationUtils {
+
+    private static final BiConsumer<ClassLoader, List<URL>> 
ADD_URL_TO_CLASSLOADER =
+            (classLoader, urls) -> {
+                if (classLoader instanceof URLClassLoader) {
+                    urls.forEach(url -> ReflectionUtils.invoke(classLoader, 
"addURL", url));
+                } else {
+                    try {
+                        Optional<Method> method =
+                                ReflectionUtils.getDeclaredMethod(
+                                        URLClassLoader.class, "addURL", 
URL.class);
+                        if (!method.isPresent()) {
+                            throw new IllegalStateException(
+                                    "Unable to find addURL method from 
URLClassLoader");
+                        }
+                        method.get().setAccessible(true);
+                        for (URL url : urls) {
+                            method.get().invoke(classLoader, url);
+                        }
+                    } catch (Exception e) {
+                        throw new RuntimeException(
+                                "Unsupported classloader: " + 
classLoader.getClass().getName(), e);
+                    }
+                }
+            };
+
+    private ConfigValidationUtils() {}
+
+    public static void validate(Config config) {
+        validate(config, CheckResult.success());
+    }
+
+    public static void validate(Config config, CheckResult checkResult) {
+        if (!checkResult.isSuccess()) {
+            throw new ConfigCheckException(checkResult.getMsg());
+        }
+
+        try {
+            JobContext jobContext = new JobContext();
+            jobContext.setJobMode(RuntimeEnvironment.getJobMode(config));
+            
jobContext.setEnableCheckpoint(RuntimeEnvironment.getEnableCheckpoint(config));
+
+            List<TableInfo> sourceTables =
+                    validateSources(
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.SOURCE, 
Collections.emptyList()),
+                            jobContext);
+            if (sourceTables.isEmpty()) {
+                throw new ConfigCheckException("At least one source plugin 
must be configured.");
+            }
+
+            List<TableInfo> outputTables =
+                    validateTransforms(
+                            sourceTables,
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.TRANSFORM, 
Collections.emptyList()),
+                            jobContext);
+
+            validateSinks(
+                    outputTables,
+                    TypesafeConfigUtils.getConfigList(
+                            config, Constants.SINK, Collections.emptyList()),
+                    jobContext);
+        } catch (ConfigCheckException e) {
+            throw e;
+        } catch (Exception e) {
+            Throwable rootException = ExceptionUtils.getRootException(e);
+            String message = rootException.getMessage();
+            throw new ConfigCheckException(
+                    message == null || message.isEmpty() ? e.getMessage() : 
message, e);
+        }
+    }
+
+    private static List<TableInfo> validateSources(
+            List<? extends Config> sourceConfigs, JobContext jobContext) {
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(
+                        
org.apache.seatunnel.api.table.factory.TableSourceFactory.class,
+                        ADD_URL_TO_CLASSLOADER);
+        SeaTunnelSourcePluginDiscovery sourcePluginDiscovery =
+                new SeaTunnelSourcePluginDiscovery(ADD_URL_TO_CLASSLOADER);
+        Function<PluginIdentifier, SeaTunnelSource> fallbackCreateSource =
+                sourcePluginDiscovery::createPluginInstance;
+
+        List<TableInfo> sourceTables = new ArrayList<>();
+        for (Config sourceConfig : sourceConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(sourceConfig, PluginType.SOURCE);
+            Tuple2<SeaTunnelSource<Object, SourceSplit, Serializable>, 
List<CatalogTable>> source =
+                    FactoryUtil.createAndPrepareSource(
+                            ReadonlyConfig.fromConfig(sourceConfig),
+                            classLoader,
+                            pluginIdentifier.getPluginName(),
+                            fallbackCreateSource,
+                            
(org.apache.seatunnel.api.table.factory.TableSourceFactory)
+                                    factoryDiscovery
+                                            
.createOptionalPluginInstance(pluginIdentifier)
+                                            .orElse(null));
+
+            source._1().setJobContext(jobContext);
+            ensureJobModeMatch(jobContext, source._1());
+            sourceTables.add(
+                    new TableInfo(
+                            source._2(),
+                            
ReadonlyConfig.fromConfig(sourceConfig).get(PLUGIN_OUTPUT)));
+        }
+        return sourceTables;
+    }
+
+    private static List<TableInfo> validateTransforms(
+            List<TableInfo> upstreamTables,
+            List<? extends Config> transformConfigs,
+            JobContext jobContext) {
+        if (transformConfigs.isEmpty()) {
+            return upstreamTables;
+        }
+
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(TableTransformFactory.class, 
ADD_URL_TO_CLASSLOADER);
+        TableInfo defaultInput = upstreamTables.get(0);
+        Map<String, TableInfo> outputTables =
+                upstreamTables.stream()
+                        .collect(
+                                Collectors.toMap(
+                                        TableInfo::getTableName,
+                                        Function.identity(),
+                                        (left, right) -> right,
+                                        LinkedHashMap::new));
+
+        for (Config transformConfig : transformConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(transformConfig, PluginType.TRANSFORM);
+
+            TableInfo inputTable =
+                    resolveInputTable(
+                                    transformConfig,
+                                    new ArrayList<>(outputTables.values()),
+                                    "Multiple input tables are not supported 
in the current version")
+                            .orElse(defaultInput);
+
+            TableTransformFactory factory =
+                    (TableTransformFactory) 
factoryDiscovery.createPluginInstance(pluginIdentifier);
+            TableTransformFactoryContext context =
+                    new TableTransformFactoryContext(
+                            inputTable.getCatalogTables(),
+                            ReadonlyConfig.fromConfig(transformConfig),
+                            classLoader);

Review Comment:
   Fixed in 3646fc626. TableTransformFactoryContext now receives the 
classloader that loaded the discovered transform factory. A regression test 
covers selecting a factory loaded by a dedicated child URLClassLoader.



##########
seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigValidationUtils.java:
##########
@@ -0,0 +1,339 @@
+/*
+ * 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.core.starter.utils;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.common.PluginIdentifier;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.SeaTunnelSink;
+import org.apache.seatunnel.api.sink.SupportMultiTableSink;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.factory.FactoryUtil;
+import org.apache.seatunnel.api.table.factory.TableSinkFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactoryContext;
+import org.apache.seatunnel.api.transform.SeaTunnelTransform;
+import org.apache.seatunnel.common.Constants;
+import org.apache.seatunnel.common.config.CheckResult;
+import org.apache.seatunnel.common.config.TypesafeConfigUtils;
+import org.apache.seatunnel.common.constants.EngineType;
+import org.apache.seatunnel.common.constants.PluginType;
+import org.apache.seatunnel.common.utils.ExceptionUtils;
+import org.apache.seatunnel.common.utils.ReflectionUtils;
+import org.apache.seatunnel.core.starter.exception.ConfigCheckException;
+import org.apache.seatunnel.core.starter.execution.RuntimeEnvironment;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelFactoryDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSinkPluginDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSourcePluginDiscovery;
+
+import scala.Tuple2;
+
+import java.io.Serializable;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.BiConsumer;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_INPUT;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_OUTPUT;
+import static 
org.apache.seatunnel.api.table.factory.FactoryUtil.ensureJobModeMatch;
+
+/** Utility methods for validating SeaTunnel job configuration without 
executing a job. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+public final class ConfigValidationUtils {
+
+    private static final BiConsumer<ClassLoader, List<URL>> 
ADD_URL_TO_CLASSLOADER =
+            (classLoader, urls) -> {
+                if (classLoader instanceof URLClassLoader) {
+                    urls.forEach(url -> ReflectionUtils.invoke(classLoader, 
"addURL", url));
+                } else {
+                    try {
+                        Optional<Method> method =
+                                ReflectionUtils.getDeclaredMethod(
+                                        URLClassLoader.class, "addURL", 
URL.class);
+                        if (!method.isPresent()) {
+                            throw new IllegalStateException(
+                                    "Unable to find addURL method from 
URLClassLoader");
+                        }
+                        method.get().setAccessible(true);
+                        for (URL url : urls) {
+                            method.get().invoke(classLoader, url);
+                        }
+                    } catch (Exception e) {
+                        throw new RuntimeException(
+                                "Unsupported classloader: " + 
classLoader.getClass().getName(), e);
+                    }
+                }
+            };
+
+    private ConfigValidationUtils() {}
+
+    public static void validate(Config config) {
+        validate(config, CheckResult.success());
+    }
+
+    public static void validate(Config config, CheckResult checkResult) {
+        if (!checkResult.isSuccess()) {
+            throw new ConfigCheckException(checkResult.getMsg());
+        }
+
+        try {
+            JobContext jobContext = new JobContext();
+            jobContext.setJobMode(RuntimeEnvironment.getJobMode(config));
+            
jobContext.setEnableCheckpoint(RuntimeEnvironment.getEnableCheckpoint(config));
+
+            List<TableInfo> sourceTables =
+                    validateSources(
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.SOURCE, 
Collections.emptyList()),
+                            jobContext);
+            if (sourceTables.isEmpty()) {
+                throw new ConfigCheckException("At least one source plugin 
must be configured.");
+            }
+
+            List<TableInfo> outputTables =
+                    validateTransforms(
+                            sourceTables,
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.TRANSFORM, 
Collections.emptyList()),
+                            jobContext);
+
+            validateSinks(
+                    outputTables,
+                    TypesafeConfigUtils.getConfigList(
+                            config, Constants.SINK, Collections.emptyList()),
+                    jobContext);
+        } catch (ConfigCheckException e) {
+            throw e;
+        } catch (Exception e) {
+            Throwable rootException = ExceptionUtils.getRootException(e);
+            String message = rootException.getMessage();
+            throw new ConfigCheckException(
+                    message == null || message.isEmpty() ? e.getMessage() : 
message, e);
+        }
+    }
+
+    private static List<TableInfo> validateSources(
+            List<? extends Config> sourceConfigs, JobContext jobContext) {
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(
+                        
org.apache.seatunnel.api.table.factory.TableSourceFactory.class,
+                        ADD_URL_TO_CLASSLOADER);
+        SeaTunnelSourcePluginDiscovery sourcePluginDiscovery =
+                new SeaTunnelSourcePluginDiscovery(ADD_URL_TO_CLASSLOADER);
+        Function<PluginIdentifier, SeaTunnelSource> fallbackCreateSource =
+                sourcePluginDiscovery::createPluginInstance;
+
+        List<TableInfo> sourceTables = new ArrayList<>();
+        for (Config sourceConfig : sourceConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(sourceConfig, PluginType.SOURCE);
+            Tuple2<SeaTunnelSource<Object, SourceSplit, Serializable>, 
List<CatalogTable>> source =
+                    FactoryUtil.createAndPrepareSource(
+                            ReadonlyConfig.fromConfig(sourceConfig),
+                            classLoader,
+                            pluginIdentifier.getPluginName(),
+                            fallbackCreateSource,
+                            
(org.apache.seatunnel.api.table.factory.TableSourceFactory)
+                                    factoryDiscovery
+                                            
.createOptionalPluginInstance(pluginIdentifier)
+                                            .orElse(null));

Review Comment:
   Fixed in 3646fc626. The source call now passes the trailing metadata 
argument (null, matching the Flink/Spark execution paths) and uses the 
discovered source factory classloader, with the validation classloader only as 
the legacy fallback.



##########
seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigValidationUtils.java:
##########
@@ -0,0 +1,339 @@
+/*
+ * 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.core.starter.utils;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.common.PluginIdentifier;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.SeaTunnelSink;
+import org.apache.seatunnel.api.sink.SupportMultiTableSink;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.factory.FactoryUtil;
+import org.apache.seatunnel.api.table.factory.TableSinkFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactoryContext;
+import org.apache.seatunnel.api.transform.SeaTunnelTransform;
+import org.apache.seatunnel.common.Constants;
+import org.apache.seatunnel.common.config.CheckResult;
+import org.apache.seatunnel.common.config.TypesafeConfigUtils;
+import org.apache.seatunnel.common.constants.EngineType;
+import org.apache.seatunnel.common.constants.PluginType;
+import org.apache.seatunnel.common.utils.ExceptionUtils;
+import org.apache.seatunnel.common.utils.ReflectionUtils;
+import org.apache.seatunnel.core.starter.exception.ConfigCheckException;
+import org.apache.seatunnel.core.starter.execution.RuntimeEnvironment;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelFactoryDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSinkPluginDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSourcePluginDiscovery;
+
+import scala.Tuple2;
+
+import java.io.Serializable;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.BiConsumer;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_INPUT;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_OUTPUT;
+import static 
org.apache.seatunnel.api.table.factory.FactoryUtil.ensureJobModeMatch;
+
+/** Utility methods for validating SeaTunnel job configuration without 
executing a job. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+public final class ConfigValidationUtils {
+
+    private static final BiConsumer<ClassLoader, List<URL>> 
ADD_URL_TO_CLASSLOADER =
+            (classLoader, urls) -> {
+                if (classLoader instanceof URLClassLoader) {
+                    urls.forEach(url -> ReflectionUtils.invoke(classLoader, 
"addURL", url));
+                } else {
+                    try {
+                        Optional<Method> method =
+                                ReflectionUtils.getDeclaredMethod(
+                                        URLClassLoader.class, "addURL", 
URL.class);
+                        if (!method.isPresent()) {
+                            throw new IllegalStateException(
+                                    "Unable to find addURL method from 
URLClassLoader");
+                        }
+                        method.get().setAccessible(true);
+                        for (URL url : urls) {
+                            method.get().invoke(classLoader, url);
+                        }
+                    } catch (Exception e) {
+                        throw new RuntimeException(
+                                "Unsupported classloader: " + 
classLoader.getClass().getName(), e);
+                    }
+                }
+            };
+
+    private ConfigValidationUtils() {}
+
+    public static void validate(Config config) {
+        validate(config, CheckResult.success());
+    }
+
+    public static void validate(Config config, CheckResult checkResult) {
+        if (!checkResult.isSuccess()) {
+            throw new ConfigCheckException(checkResult.getMsg());
+        }
+
+        try {
+            JobContext jobContext = new JobContext();
+            jobContext.setJobMode(RuntimeEnvironment.getJobMode(config));
+            
jobContext.setEnableCheckpoint(RuntimeEnvironment.getEnableCheckpoint(config));
+
+            List<TableInfo> sourceTables =
+                    validateSources(
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.SOURCE, 
Collections.emptyList()),
+                            jobContext);
+            if (sourceTables.isEmpty()) {
+                throw new ConfigCheckException("At least one source plugin 
must be configured.");
+            }
+
+            List<TableInfo> outputTables =
+                    validateTransforms(
+                            sourceTables,
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.TRANSFORM, 
Collections.emptyList()),
+                            jobContext);
+
+            validateSinks(
+                    outputTables,
+                    TypesafeConfigUtils.getConfigList(
+                            config, Constants.SINK, Collections.emptyList()),
+                    jobContext);
+        } catch (ConfigCheckException e) {
+            throw e;
+        } catch (Exception e) {
+            Throwable rootException = ExceptionUtils.getRootException(e);
+            String message = rootException.getMessage();
+            throw new ConfigCheckException(
+                    message == null || message.isEmpty() ? e.getMessage() : 
message, e);
+        }
+    }
+
+    private static List<TableInfo> validateSources(
+            List<? extends Config> sourceConfigs, JobContext jobContext) {
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(
+                        
org.apache.seatunnel.api.table.factory.TableSourceFactory.class,
+                        ADD_URL_TO_CLASSLOADER);
+        SeaTunnelSourcePluginDiscovery sourcePluginDiscovery =
+                new SeaTunnelSourcePluginDiscovery(ADD_URL_TO_CLASSLOADER);
+        Function<PluginIdentifier, SeaTunnelSource> fallbackCreateSource =
+                sourcePluginDiscovery::createPluginInstance;
+
+        List<TableInfo> sourceTables = new ArrayList<>();
+        for (Config sourceConfig : sourceConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(sourceConfig, PluginType.SOURCE);
+            Tuple2<SeaTunnelSource<Object, SourceSplit, Serializable>, 
List<CatalogTable>> source =
+                    FactoryUtil.createAndPrepareSource(
+                            ReadonlyConfig.fromConfig(sourceConfig),
+                            classLoader,
+                            pluginIdentifier.getPluginName(),
+                            fallbackCreateSource,
+                            
(org.apache.seatunnel.api.table.factory.TableSourceFactory)
+                                    factoryDiscovery
+                                            
.createOptionalPluginInstance(pluginIdentifier)
+                                            .orElse(null));
+
+            source._1().setJobContext(jobContext);
+            ensureJobModeMatch(jobContext, source._1());
+            sourceTables.add(
+                    new TableInfo(
+                            source._2(),
+                            
ReadonlyConfig.fromConfig(sourceConfig).get(PLUGIN_OUTPUT)));
+        }
+        return sourceTables;
+    }
+
+    private static List<TableInfo> validateTransforms(
+            List<TableInfo> upstreamTables,
+            List<? extends Config> transformConfigs,
+            JobContext jobContext) {
+        if (transformConfigs.isEmpty()) {
+            return upstreamTables;
+        }
+
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(TableTransformFactory.class, 
ADD_URL_TO_CLASSLOADER);
+        TableInfo defaultInput = upstreamTables.get(0);
+        Map<String, TableInfo> outputTables =
+                upstreamTables.stream()
+                        .collect(
+                                Collectors.toMap(
+                                        TableInfo::getTableName,
+                                        Function.identity(),
+                                        (left, right) -> right,
+                                        LinkedHashMap::new));
+
+        for (Config transformConfig : transformConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(transformConfig, PluginType.TRANSFORM);
+
+            TableInfo inputTable =
+                    resolveInputTable(
+                                    transformConfig,
+                                    new ArrayList<>(outputTables.values()),
+                                    "Multiple input tables are not supported 
in the current version")
+                            .orElse(defaultInput);
+
+            TableTransformFactory factory =
+                    (TableTransformFactory) 
factoryDiscovery.createPluginInstance(pluginIdentifier);
+            TableTransformFactoryContext context =
+                    new TableTransformFactoryContext(
+                            inputTable.getCatalogTables(),
+                            ReadonlyConfig.fromConfig(transformConfig),
+                            classLoader);
+            
org.apache.seatunnel.api.configuration.util.ConfigValidator.of(context.getOptions())
+                    .validate(factory.optionRule());
+            SeaTunnelTransform<?> transform = 
factory.createTransform(context).createTransform();
+            transform.setJobContext(jobContext);
+
+            String pluginOutputIdentifier =
+                    
ReadonlyConfig.fromConfig(transformConfig).get(PLUGIN_OUTPUT);
+            outputTables.put(
+                    pluginOutputIdentifier,
+                    new TableInfo(transform.getProducedCatalogTables(), 
pluginOutputIdentifier));
+        }
+        return new ArrayList<>(outputTables.values());
+    }
+
+    private static void validateSinks(
+            List<TableInfo> upstreamTables,
+            List<? extends Config> sinkConfigs,
+            JobContext jobContext) {
+        if (sinkConfigs.isEmpty()) {
+            throw new ConfigCheckException("At least one sink plugin must be 
configured.");
+        }
+
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(TableSinkFactory.class, 
ADD_URL_TO_CLASSLOADER);
+        SeaTunnelSinkPluginDiscovery sinkPluginDiscovery =
+                new SeaTunnelSinkPluginDiscovery(ADD_URL_TO_CLASSLOADER);
+        Function<PluginIdentifier, SeaTunnelSink> fallbackCreateSink =
+                sinkPluginDiscovery::createPluginInstance;
+
+        TableInfo defaultInput = upstreamTables.get(upstreamTables.size() - 1);
+        for (Config sinkConfig : sinkConfigs) {
+            PluginIdentifier pluginIdentifier = 
getPluginIdentifier(sinkConfig, PluginType.SINK);
+            TableInfo inputTable =
+                    resolveInputTable(
+                                    sinkConfig,
+                                    upstreamTables,
+                                    "Multiple input tables are not supported 
in the current version")
+                            .orElse(defaultInput);
+
+            Map<TablePath, SeaTunnelSink> sinks = new LinkedHashMap<>();
+            for (CatalogTable catalogTable : inputTable.getCatalogTables()) {
+                SeaTunnelSink sink =
+                        FactoryUtil.createAndPrepareSink(
+                                catalogTable,
+                                ReadonlyConfig.fromConfig(sinkConfig),
+                                classLoader,
+                                pluginIdentifier.getPluginName(),
+                                fallbackCreateSink,
+                                (TableSinkFactory)
+                                        factoryDiscovery
+                                                
.createOptionalPluginInstance(pluginIdentifier)
+                                                .orElse(null));

Review Comment:
   Fixed in 3646fc626. Sink preparation now reuses the discovered sink factory 
classloader, so FactoryUtil fallback detection, placeholder replacement, and 
option validation all run with the same loader.



##########
seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/utils/ConfigValidationUtils.java:
##########
@@ -0,0 +1,339 @@
+/*
+ * 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.core.starter.utils;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.common.PluginIdentifier;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.SeaTunnelSink;
+import org.apache.seatunnel.api.sink.SupportMultiTableSink;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.factory.FactoryUtil;
+import org.apache.seatunnel.api.table.factory.TableSinkFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactory;
+import org.apache.seatunnel.api.table.factory.TableTransformFactoryContext;
+import org.apache.seatunnel.api.transform.SeaTunnelTransform;
+import org.apache.seatunnel.common.Constants;
+import org.apache.seatunnel.common.config.CheckResult;
+import org.apache.seatunnel.common.config.TypesafeConfigUtils;
+import org.apache.seatunnel.common.constants.EngineType;
+import org.apache.seatunnel.common.constants.PluginType;
+import org.apache.seatunnel.common.utils.ExceptionUtils;
+import org.apache.seatunnel.common.utils.ReflectionUtils;
+import org.apache.seatunnel.core.starter.exception.ConfigCheckException;
+import org.apache.seatunnel.core.starter.execution.RuntimeEnvironment;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelFactoryDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSinkPluginDiscovery;
+import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelSourcePluginDiscovery;
+
+import scala.Tuple2;
+
+import java.io.Serializable;
+import java.lang.reflect.Method;
+import java.net.URL;
+import java.net.URLClassLoader;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.function.BiConsumer;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_INPUT;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME;
+import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_OUTPUT;
+import static 
org.apache.seatunnel.api.table.factory.FactoryUtil.ensureJobModeMatch;
+
+/** Utility methods for validating SeaTunnel job configuration without 
executing a job. */
+@SuppressWarnings({"rawtypes", "unchecked"})
+public final class ConfigValidationUtils {
+
+    private static final BiConsumer<ClassLoader, List<URL>> 
ADD_URL_TO_CLASSLOADER =
+            (classLoader, urls) -> {
+                if (classLoader instanceof URLClassLoader) {
+                    urls.forEach(url -> ReflectionUtils.invoke(classLoader, 
"addURL", url));
+                } else {
+                    try {
+                        Optional<Method> method =
+                                ReflectionUtils.getDeclaredMethod(
+                                        URLClassLoader.class, "addURL", 
URL.class);
+                        if (!method.isPresent()) {
+                            throw new IllegalStateException(
+                                    "Unable to find addURL method from 
URLClassLoader");
+                        }
+                        method.get().setAccessible(true);
+                        for (URL url : urls) {
+                            method.get().invoke(classLoader, url);
+                        }
+                    } catch (Exception e) {
+                        throw new RuntimeException(
+                                "Unsupported classloader: " + 
classLoader.getClass().getName(), e);
+                    }
+                }
+            };
+
+    private ConfigValidationUtils() {}
+
+    public static void validate(Config config) {
+        validate(config, CheckResult.success());
+    }
+
+    public static void validate(Config config, CheckResult checkResult) {
+        if (!checkResult.isSuccess()) {
+            throw new ConfigCheckException(checkResult.getMsg());
+        }
+
+        try {
+            JobContext jobContext = new JobContext();
+            jobContext.setJobMode(RuntimeEnvironment.getJobMode(config));
+            
jobContext.setEnableCheckpoint(RuntimeEnvironment.getEnableCheckpoint(config));
+
+            List<TableInfo> sourceTables =
+                    validateSources(
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.SOURCE, 
Collections.emptyList()),
+                            jobContext);
+            if (sourceTables.isEmpty()) {
+                throw new ConfigCheckException("At least one source plugin 
must be configured.");
+            }
+
+            List<TableInfo> outputTables =
+                    validateTransforms(
+                            sourceTables,
+                            TypesafeConfigUtils.getConfigList(
+                                    config, Constants.TRANSFORM, 
Collections.emptyList()),
+                            jobContext);
+
+            validateSinks(
+                    outputTables,
+                    TypesafeConfigUtils.getConfigList(
+                            config, Constants.SINK, Collections.emptyList()),
+                    jobContext);
+        } catch (ConfigCheckException e) {
+            throw e;
+        } catch (Exception e) {
+            Throwable rootException = ExceptionUtils.getRootException(e);
+            String message = rootException.getMessage();
+            throw new ConfigCheckException(
+                    message == null || message.isEmpty() ? e.getMessage() : 
message, e);
+        }
+    }
+
+    private static List<TableInfo> validateSources(
+            List<? extends Config> sourceConfigs, JobContext jobContext) {
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(
+                        
org.apache.seatunnel.api.table.factory.TableSourceFactory.class,
+                        ADD_URL_TO_CLASSLOADER);
+        SeaTunnelSourcePluginDiscovery sourcePluginDiscovery =
+                new SeaTunnelSourcePluginDiscovery(ADD_URL_TO_CLASSLOADER);
+        Function<PluginIdentifier, SeaTunnelSource> fallbackCreateSource =
+                sourcePluginDiscovery::createPluginInstance;
+
+        List<TableInfo> sourceTables = new ArrayList<>();
+        for (Config sourceConfig : sourceConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(sourceConfig, PluginType.SOURCE);
+            Tuple2<SeaTunnelSource<Object, SourceSplit, Serializable>, 
List<CatalogTable>> source =
+                    FactoryUtil.createAndPrepareSource(
+                            ReadonlyConfig.fromConfig(sourceConfig),
+                            classLoader,
+                            pluginIdentifier.getPluginName(),
+                            fallbackCreateSource,
+                            
(org.apache.seatunnel.api.table.factory.TableSourceFactory)
+                                    factoryDiscovery
+                                            
.createOptionalPluginInstance(pluginIdentifier)
+                                            .orElse(null));
+
+            source._1().setJobContext(jobContext);
+            ensureJobModeMatch(jobContext, source._1());
+            sourceTables.add(
+                    new TableInfo(
+                            source._2(),
+                            
ReadonlyConfig.fromConfig(sourceConfig).get(PLUGIN_OUTPUT)));
+        }
+        return sourceTables;
+    }
+
+    private static List<TableInfo> validateTransforms(
+            List<TableInfo> upstreamTables,
+            List<? extends Config> transformConfigs,
+            JobContext jobContext) {
+        if (transformConfigs.isEmpty()) {
+            return upstreamTables;
+        }
+
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(TableTransformFactory.class, 
ADD_URL_TO_CLASSLOADER);
+        TableInfo defaultInput = upstreamTables.get(0);
+        Map<String, TableInfo> outputTables =
+                upstreamTables.stream()
+                        .collect(
+                                Collectors.toMap(
+                                        TableInfo::getTableName,
+                                        Function.identity(),
+                                        (left, right) -> right,
+                                        LinkedHashMap::new));
+
+        for (Config transformConfig : transformConfigs) {
+            PluginIdentifier pluginIdentifier =
+                    getPluginIdentifier(transformConfig, PluginType.TRANSFORM);
+
+            TableInfo inputTable =
+                    resolveInputTable(
+                                    transformConfig,
+                                    new ArrayList<>(outputTables.values()),
+                                    "Multiple input tables are not supported 
in the current version")
+                            .orElse(defaultInput);
+
+            TableTransformFactory factory =
+                    (TableTransformFactory) 
factoryDiscovery.createPluginInstance(pluginIdentifier);
+            TableTransformFactoryContext context =
+                    new TableTransformFactoryContext(
+                            inputTable.getCatalogTables(),
+                            ReadonlyConfig.fromConfig(transformConfig),
+                            classLoader);
+            
org.apache.seatunnel.api.configuration.util.ConfigValidator.of(context.getOptions())
+                    .validate(factory.optionRule());
+            SeaTunnelTransform<?> transform = 
factory.createTransform(context).createTransform();
+            transform.setJobContext(jobContext);
+
+            String pluginOutputIdentifier =
+                    
ReadonlyConfig.fromConfig(transformConfig).get(PLUGIN_OUTPUT);
+            outputTables.put(
+                    pluginOutputIdentifier,
+                    new TableInfo(transform.getProducedCatalogTables(), 
pluginOutputIdentifier));
+        }
+        return new ArrayList<>(outputTables.values());
+    }
+
+    private static void validateSinks(
+            List<TableInfo> upstreamTables,
+            List<? extends Config> sinkConfigs,
+            JobContext jobContext) {
+        if (sinkConfigs.isEmpty()) {
+            throw new ConfigCheckException("At least one sink plugin must be 
configured.");
+        }
+
+        ClassLoader classLoader = 
Thread.currentThread().getContextClassLoader();
+        SeaTunnelFactoryDiscovery factoryDiscovery =
+                new SeaTunnelFactoryDiscovery(TableSinkFactory.class, 
ADD_URL_TO_CLASSLOADER);
+        SeaTunnelSinkPluginDiscovery sinkPluginDiscovery =
+                new SeaTunnelSinkPluginDiscovery(ADD_URL_TO_CLASSLOADER);
+        Function<PluginIdentifier, SeaTunnelSink> fallbackCreateSink =
+                sinkPluginDiscovery::createPluginInstance;
+
+        TableInfo defaultInput = upstreamTables.get(upstreamTables.size() - 1);
+        for (Config sinkConfig : sinkConfigs) {
+            PluginIdentifier pluginIdentifier = 
getPluginIdentifier(sinkConfig, PluginType.SINK);
+            TableInfo inputTable =
+                    resolveInputTable(
+                                    sinkConfig,
+                                    upstreamTables,
+                                    "Multiple input tables are not supported 
in the current version")
+                            .orElse(defaultInput);
+
+            Map<TablePath, SeaTunnelSink> sinks = new LinkedHashMap<>();
+            for (CatalogTable catalogTable : inputTable.getCatalogTables()) {
+                SeaTunnelSink sink =
+                        FactoryUtil.createAndPrepareSink(
+                                catalogTable,
+                                ReadonlyConfig.fromConfig(sinkConfig),
+                                classLoader,
+                                pluginIdentifier.getPluginName(),
+                                fallbackCreateSink,
+                                (TableSinkFactory)
+                                        factoryDiscovery
+                                                
.createOptionalPluginInstance(pluginIdentifier)
+                                                .orElse(null));
+                sink.setJobContext(jobContext);
+                sinks.put(catalogTable.getTableId().toTablePath(), sink);
+            }
+
+            if (!sinks.isEmpty()
+                    && 
sinks.values().stream().allMatch(SupportMultiTableSink.class::isInstance)) {
+                FactoryUtil.createMultiTableSink(
+                        sinks, ReadonlyConfig.fromConfig(sinkConfig), 
classLoader);
+            }

Review Comment:
   Fixed in 3646fc626. Multi-table sink creation now receives the same 
classloader selected for the sink factory, keeping the entire sink validation 
path consistent.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


Reply via email to