diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java b/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java index bdaa49c5..39cd6ad0 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-flink/datavines-engine-flink-config/src/main/java/io/datavines/engine/flink/config/BaseFlinkConfigurationBuilder.java @@ -91,8 +91,17 @@ protected List getSourceConfigs() throws DataVinesException { if (tableArray.length == 1) { metricInputParameter.put(TABLE, tableArray[0]); } else { - metricInputParameter.put(DATABASE, tableArray[0]); - metricInputParameter.put(TABLE, tableArray[1]); + String connectorSchema = connectorParameter.getParameters().get(SCHEMA) != null ? + (String)connectorParameter.getParameters().get(SCHEMA) : null; + if (connectorSchema != null && connectorSchema.equals(tableArray[0])) { + // When the table name from SQL already includes the schema prefix + // (e.g., "schema.table"), and it matches the configured schema, + // treat it as schema.table, not database.table + metricInputParameter.put(TABLE, tableArray[1]); + } else { + metricInputParameter.put(DATABASE, tableArray[0]); + metricInputParameter.put(TABLE, tableArray[1]); + } } SourceConfig sourceConfig = new SourceConfig(); diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java b/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java index 930b0192..4752b07d 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-local/datavines-engine-local-config/src/main/java/io/datavines/engine/local/config/BaseLocalConfigurationBuilder.java @@ -82,8 +82,17 @@ protected List getSourceConfigs() throws DataVinesException { if (tableArray.length == 1) { metricInputParameter.put(TABLE, tableArray[0]); } else { - metricInputParameter.put(DATABASE, tableArray[0]); - metricInputParameter.put(TABLE, tableArray[1]); + String connectorSchema = connectorParameter.getParameters().get(SCHEMA) != null ? + (String)connectorParameter.getParameters().get(SCHEMA) : null; + if (connectorSchema != null && connectorSchema.equals(tableArray[0])) { + // When the table name from SQL already includes the schema prefix + // (e.g., "schema.table"), and it matches the configured schema, + // treat it as schema.table, not database.table + metricInputParameter.put(TABLE, tableArray[1]); + } else { + metricInputParameter.put(DATABASE, tableArray[0]); + metricInputParameter.put(TABLE, tableArray[1]); + } } Map connectorParameterMap = new HashMap<>(connectorParameter.getParameters()); diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java b/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java index eafb45a3..eff538d7 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-config/src/main/java/io/datavines/engine/spark/config/BaseSparkConfigurationBuilder.java @@ -117,8 +117,17 @@ protected List getSourceConfigs() throws DataVinesException { if (tableArray.length == 1) { metricInputParameter.put(TABLE, tableArray[0]); } else { - metricInputParameter.put(DATABASE, tableArray[0]); - metricInputParameter.put(TABLE, tableArray[1]); + String connectorSchema = connectorParameter.getParameters().get(SCHEMA) != null ? + (String)connectorParameter.getParameters().get(SCHEMA) : null; + if (connectorSchema != null && connectorSchema.equals(tableArray[0])) { + // When the table name from SQL already includes the schema prefix + // (e.g., "schema.table"), and it matches the configured schema, + // treat it as schema.table, not database.table + metricInputParameter.put(TABLE, tableArray[1]); + } else { + metricInputParameter.put(DATABASE, tableArray[0]); + metricInputParameter.put(TABLE, tableArray[1]); + } } SourceConfig sourceConfig = new SourceConfig(); diff --git a/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-core/src/main/java/io/datavines/engine/spark/jdbc/source/JdbcSource.java b/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-core/src/main/java/io/datavines/engine/spark/jdbc/source/JdbcSource.java index 926ef5f2..e9f19258 100644 --- a/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-core/src/main/java/io/datavines/engine/spark/jdbc/source/JdbcSource.java +++ b/datavines-engine/datavines-engine-plugins/datavines-engine-spark/datavines-engine-spark-core/src/main/java/io/datavines/engine/spark/jdbc/source/JdbcSource.java @@ -101,7 +101,12 @@ public Dataset getData(SparkRuntimeEnvironment env) { } DataFrameReader reader = new DataFrameReader(env.sparkSession()); - return reader.jdbc(config.getString(URL), config.getString(TABLE), properties); + String table = config.getString(TABLE); + String schema = config.getString(SCHEMA); + if (!StringUtils.isEmptyOrNullStr(schema)) { + table = schema + "." + table; + } + return reader.jdbc(config.getString(URL), table, properties); } private Dataset hiveSourceData(SparkRuntimeEnvironment env) {