{"record":{"id":"229270fb329bd6ad","repo":"apache/seatunnel","slug":"multiple-input-tables-are-not-supported-in-flink-p","errorCode":null,"errorMessage":"Multiple input tables are not supported in flink plugin","messagePattern":"Multiple input tables are not supported in flink plugin","errorType":"validation","errorClass":"UnsupportedOperationException","httpStatus":null,"severity":"error","filePath":"seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/FlinkAbstractPluginExecuteProcessor.java","lineNumber":101,"sourceCode":"        this.pluginConfigs = pluginConfigs;\n        this.jobContext = jobContext;\n        this.plugins = initializePlugins(jarPaths, pluginConfigs);\n        this.envConfig = envConfig;\n    }\n\n    @Override\n    public void setRuntimeEnvironment(FlinkRuntimeEnvironment flinkRuntimeEnvironment) {\n        this.flinkRuntimeEnvironment = flinkRuntimeEnvironment;\n    }\n\n    protected Optional<DataStreamTableInfo> fromSourceTable(\n            Config pluginConfig, List<DataStreamTableInfo> upstreamDataStreams) {\n        ReadonlyConfig readonlyConfig = ReadonlyConfig.fromConfig(pluginConfig);\n\n        if (readonlyConfig.getOptional(PLUGIN_INPUT).isPresent()) {\n            List<String> pluginInputIdentifiers = readonlyConfig.get(PLUGIN_INPUT);\n            if (pluginInputIdentifiers.size() > 1) {\n                throw new UnsupportedOperationException(\n                        \"Multiple input tables are not supported in flink plugin\");\n            }\n\n            String tableName = pluginInputIdentifiers.get(0);\n            DataStreamTableInfo dataStreamTableInfo =\n                    upstreamDataStreams.stream()\n                            .filter(info -> tableName.equals(info.getTableName()))\n                            .findFirst()\n                            .orElseThrow(\n                                    () ->\n                                            new SeaTunnelException(\n                                                    String.format(\n                                                            \"table %s not found\", tableName)));\n            return Optional.of(\n                    new DataStreamTableInfo(\n                            dataStreamTableInfo.getDataStream(),\n                            dataStreamTableInfo.getCatalogTables(),\n                            tableName));","sourceCodeStart":83,"sourceCodeEnd":119,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/FlinkAbstractPluginExecuteProcessor.java#L83-L119","documentation":"FlinkAbstractPluginExecuteProcessor.fromSourceTable resolves which upstream DataStream(s) a plugin consumes via the 'plugin_input' option. The Flink starter only supports single-input plugins, so it explicitly rejects configurations listing more than one input table identifier with an UnsupportedOperationException.","triggerScenarios":"A job config specifies 'plugin_input' with more than one table identifier (or 'plugin_input = [\"t1\",\"t2\"]') for a Flink-executed source/transform/sink plugin.","commonSituations":"Users porting multi-input jobs from the Zeta engine to Flink; copy-pasting configs that feed two upstream transforms into one plugin; joining two streams in a Flink-launched SeaTunnel job.","solutions":["Reduce the job to a single upstream input per plugin and add intermediate transform stages instead of multi-input joins","Run the job on the Zeta (SeaTunnel) engine, which supports multiple input tables","Combine upstream streams before the plugin so only one DataStream reaches it","Remove duplicate plugin_input entries if the multiple values were unintentional"],"exampleFix":"// before\ntransform {\n  Sql = {\n    plugin_input = [\"source_a\", \"source_b\"]\n  }\n}\n// after\ntransform {\n  Sql = {\n    plugin_input = \"source_a\"\n  }\n}","handlingStrategy":"validation","validationCode":"List<String> inputs = readonlyConfig.get(PLUGIN_INPUT);\nif (inputs != null && inputs.size() > 1) {\n    throw new IllegalArgumentException(\"Flink starter supports at most one plugin_input, got: \" + inputs);\n}","typeGuard":"boolean supportsMultiInput(JobMode mode) {\n    return mode == JobMode.ZETA; // only Zeta supports multiple input tables\n}","tryCatchPattern":"try {\n    processor.execute(...);\n} catch (UnsupportedOperationException e) {\n    // fall back to single-input config or Zeta engine\n}","preventionTips":["Keep one plugin_input per plugin in Flink jobs","Use the Zeta engine when a job needs multi-input/join semantics","Lint job configs to reject plugin_input lists longer than 1 for Flink execution"],"tags":["flink","multi-table","unsupported-operation"],"backgroundTag":"unsupported-operation","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}