From 28ef734b0d8fe48b705ba11dfccfdebba9151fc6 Mon Sep 17 00:00:00 2001 From: unknown Date: Mon, 31 Aug 2026 18:55:43 +0500 Subject: [PATCH] [JDBC Template] Use Spark distributed JDBC DataFrame reader in SqlToRddOperator (#761) --- pom.xml | 58 +++++++- wayang-commons/wayang-basic/pom.xml | 2 +- wayang-commons/wayang-core/pom.xml | 2 +- wayang-platforms/wayang-java/pom.xml | 4 +- .../jdbc/execution/DatabaseDescriptor.java | 16 +++ .../jdbc/operators/SqlToRddOperator.java | 126 +++++++++++++++--- .../jdbc/operators/SqlToRddOperatorTest.java | 21 +++ wayang-platforms/wayang-spark/pom.xml | 18 +-- 8 files changed, 213 insertions(+), 34 deletions(-) diff --git a/pom.xml b/pom.xml index 2eab5dd8e..4d08cc103 100644 --- a/pom.xml +++ b/pom.xml @@ -95,7 +95,9 @@ 2.12.17 2.12 - 3.4.4 + 3.5.7 + 4.9.3 + 1.13.1 1.6.1 1.20.0 1.39.0 @@ -392,7 +394,7 @@ 2.12.17 2.12 - 3.4.4 + 3.5.7 @@ -498,7 +500,7 @@ - 4.13.1 + 4.9.3 @@ -733,6 +735,56 @@ commons-cli 1.3.1 + + org.antlr + antlr4-runtime + ${antlr.version} + + + org.apache.spark + spark-core_${scala.mayor.version} + ${spark.version} + + + org.apache.spark + spark-graphx_${scala.mayor.version} + ${spark.version} + + + org.apache.spark + spark-sql_${scala.mayor.version} + ${spark.version} + + + org.apache.spark + spark-mllib_${scala.mayor.version} + ${spark.version} + + + org.apache.parquet + parquet-hadoop + ${parquet.version} + + + org.apache.parquet + parquet-column + ${parquet.version} + + + org.apache.parquet + parquet-avro + ${parquet.version} + + + org.apache.parquet + parquet-common + ${parquet.version} + + + org.apache.parquet + parquet-encoding + ${parquet.version} + diff --git a/wayang-commons/wayang-basic/pom.xml b/wayang-commons/wayang-basic/pom.xml index 342bb2b40..ec5737ad1 100644 --- a/wayang-commons/wayang-basic/pom.xml +++ b/wayang-commons/wayang-basic/pom.xml @@ -87,7 +87,7 @@ org.apache.parquet parquet-hadoop - 1.12.3 + ${parquet.version} org.apache.commons diff --git a/wayang-commons/wayang-core/pom.xml b/wayang-commons/wayang-core/pom.xml index 07fa0f549..704d9af1e 100644 --- a/wayang-commons/wayang-core/pom.xml +++ b/wayang-commons/wayang-core/pom.xml @@ -83,7 +83,7 @@ org.antlr antlr4-runtime - 4.13.1 + ${antlr.version} org.apache.logging.log4j diff --git a/wayang-platforms/wayang-java/pom.xml b/wayang-platforms/wayang-java/pom.xml index 9cd901d26..d07ef27d5 100644 --- a/wayang-platforms/wayang-java/pom.xml +++ b/wayang-platforms/wayang-java/pom.xml @@ -52,12 +52,12 @@ org.apache.parquet parquet-avro - 1.15.2 + ${parquet.version} org.apache.parquet parquet-hadoop - 1.15.2 + ${parquet.version} org.apache.avro diff --git a/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/DatabaseDescriptor.java b/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/DatabaseDescriptor.java index 0aeb47511..27b19b11c 100644 --- a/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/DatabaseDescriptor.java +++ b/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/execution/DatabaseDescriptor.java @@ -46,6 +46,22 @@ public DatabaseDescriptor(String jdbcUrl, String user, String password, String j this.jdbcDriverClassName = jdbcDriverClassName; } + public String getJdbcUrl() { + return this.jdbcUrl; + } + + public String getUser() { + return this.user; + } + + public String getPassword() { + return this.password; + } + + public String getJdbcDriverClassName() { + return this.jdbcDriverClassName; + } + /** * Creates a {@link Connection} to the database described by this instance. * diff --git a/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/operators/SqlToRddOperator.java b/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/operators/SqlToRddOperator.java index 48cff7194..265ce3ea0 100644 --- a/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/operators/SqlToRddOperator.java +++ b/wayang-platforms/wayang-jdbc-template/src/main/java/org/apache/wayang/jdbc/operators/SqlToRddOperator.java @@ -19,26 +19,32 @@ package org.apache.wayang.jdbc.operators; import org.apache.spark.api.java.JavaRDD; +import org.apache.spark.sql.DataFrameReader; +import org.apache.spark.sql.Dataset; +import org.apache.spark.sql.Row; import org.apache.wayang.basic.data.Record; import org.apache.wayang.core.optimizer.OptimizationContext; +import org.apache.wayang.core.optimizer.costs.LoadProfileEstimators; import org.apache.wayang.core.plan.wayangplan.UnaryToUnaryOperator; import org.apache.wayang.core.platform.ChannelDescriptor; import org.apache.wayang.core.platform.ChannelInstance; import org.apache.wayang.core.platform.lineage.ExecutionLineageNode; import org.apache.wayang.core.types.DataSetType; import org.apache.wayang.core.util.JsonSerializable; +import org.apache.wayang.core.util.ReflectionUtils; import org.apache.wayang.core.util.Tuple; import org.apache.wayang.core.util.json.WayangJsonObj; import org.apache.wayang.jdbc.channels.SqlQueryChannel; +import org.apache.wayang.jdbc.execution.DatabaseDescriptor; import org.apache.wayang.jdbc.platform.JdbcPlatformTemplate; import org.apache.wayang.spark.channels.RddChannel; import org.apache.wayang.spark.execution.SparkExecutor; import org.apache.wayang.spark.operators.SparkExecutionOperator; -import java.sql.Connection; -import java.util.*; -import java.util.stream.Collectors; -import java.util.stream.StreamSupport; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.List; public class SqlToRddOperator extends UnaryToUnaryOperator implements SparkExecutionOperator, JsonSerializable { @@ -79,37 +85,121 @@ public Tuple, Collection> eval final RddChannel.Instance output = (RddChannel.Instance) outputs[0]; JdbcPlatformTemplate producerPlatform = (JdbcPlatformTemplate) input.getChannel().getProducer().getPlatform(); - final Connection connection = producerPlatform - .createDatabaseDescriptor(executor.getConfiguration()) - .createJdbcConnection(); - - Iterator resultSetIterator = new SqlToStreamOperator.ResultSetIterator(connection, input.getSqlQuery()); - Iterable resultSetIterable = () -> resultSetIterator; - - // Convert the ResultSet to a JavaRDD. - JavaRDD resultSetRDD = executor.sc.parallelize( - StreamSupport.stream(resultSetIterable.spliterator(), false).collect(Collectors.toList()), - executor.getNumDefaultPartitions() - ); + DatabaseDescriptor databaseDescriptor = producerPlatform.createDatabaseDescriptor(executor.getConfiguration()); + + String sqlQuery = cleanQuery(input.getSqlQuery()); + String dbtable = isTableName(sqlQuery) ? sqlQuery : "(" + sqlQuery + ") as wayang_subquery"; + + DataFrameReader reader = executor.ss.read() + .format("jdbc") + .option("url", databaseDescriptor.getJdbcUrl()) + .option("dbtable", dbtable) + .option("driver", databaseDescriptor.getJdbcDriverClassName()); + + if (databaseDescriptor.getUser() != null) { + reader.option("user", databaseDescriptor.getUser()); + } + if (databaseDescriptor.getPassword() != null) { + reader.option("password", databaseDescriptor.getPassword()); + } + + // Apply optional partition properties if configured + String partitionColumn = executor.getConfiguration().getStringProperty( + String.format("wayang.%s.jdbc.partitionColumn", producerPlatform.getPlatformId()), null); + if (partitionColumn != null) { + String lowerBound = executor.getConfiguration().getStringProperty( + String.format("wayang.%s.jdbc.lowerBound", producerPlatform.getPlatformId()), null); + String upperBound = executor.getConfiguration().getStringProperty( + String.format("wayang.%s.jdbc.upperBound", producerPlatform.getPlatformId()), null); + if (lowerBound == null || upperBound == null) { + throw new IllegalArgumentException( + "JDBC partitioning requires lowerBound and upperBound when partitionColumn is set."); + } + int numPartitions = executor.getConfiguration().getOptionalIntProperty( + String.format("wayang.%s.jdbc.numPartitions", producerPlatform.getPlatformId())) + .orElse(executor.getNumDefaultPartitions()); + reader.option("partitionColumn", partitionColumn) + .option("lowerBound", lowerBound) + .option("upperBound", upperBound) + .option("numPartitions", String.valueOf(numPartitions)); + } + + // Apply optional fetchsize if configured + String fetchSize = executor.getConfiguration().getStringProperty( + String.format("wayang.%s.jdbc.fetchsize", producerPlatform.getPlatformId()), null); + if (fetchSize != null) { + reader.option("fetchsize", fetchSize); + } + + Dataset df = reader.load(); + + // Convert the distributed DataFrame to JavaRDD lazily on executors + JavaRDD resultSetRDD = df.toJavaRDD().map(SqlToRddOperator::rowToRecord); output.accept(resultSetRDD, executor); - // TODO: Add load profile estimators ExecutionLineageNode queryLineageNode = new ExecutionLineageNode(operatorContext); + queryLineageNode.add(LoadProfileEstimators.createFromSpecification( + String.format("wayang.%s.sqltordd.load.query", this.jdbcPlatform.getPlatformId()), + executor.getConfiguration() + )); queryLineageNode.addPredecessor(input.getLineage()); ExecutionLineageNode outputLineageNode = new ExecutionLineageNode(operatorContext); + outputLineageNode.add(LoadProfileEstimators.createFromSpecification( + String.format("wayang.%s.sqltordd.load.output", this.jdbcPlatform.getPlatformId()), + executor.getConfiguration() + )); output.getLineage().addPredecessor(outputLineageNode); return queryLineageNode.collectAndMark(); } + public static Record rowToRecord(Row row) { + int length = row.size(); + Object[] fields = new Object[length]; + for (int i = 0; i < length; i++) { + fields[i] = row.get(i); + } + return new Record(fields); + } + + private static String cleanQuery(String query) { + if (query == null) { + return ""; + } + String trimmed = query.trim(); + while (trimmed.endsWith(";")) { + trimmed = trimmed.substring(0, trimmed.length() - 1).trim(); + } + return trimmed; + } + + private static boolean isTableName(String query) { + return !query.contains(" ") && !query.contains("\t") && !query.contains("\n") && !query.contains("("); + } + @Override public boolean containsAction() { return false; } + @Override + public Collection getLoadProfileEstimatorConfigurationKeys() { + return Arrays.asList( + String.format("wayang.%s.sqltordd.load.query", this.jdbcPlatform.getPlatformId()), + String.format("wayang.%s.sqltordd.load.output", this.jdbcPlatform.getPlatformId()) + ); + } + @Override public WayangJsonObj toJson() { - return null; + return new WayangJsonObj().put("platform", this.jdbcPlatform.getClass().getCanonicalName()); + } + + @SuppressWarnings("unused") + public static SqlToRddOperator fromJson(WayangJsonObj wayangJsonObj) { + final String platformClassName = wayangJsonObj.getString("platform"); + JdbcPlatformTemplate jdbcPlatform = ReflectionUtils.evaluate(platformClassName + ".getInstance()"); + return new SqlToRddOperator(jdbcPlatform); } } diff --git a/wayang-platforms/wayang-jdbc-template/src/test/java/org/apache/wayang/jdbc/operators/SqlToRddOperatorTest.java b/wayang-platforms/wayang-jdbc-template/src/test/java/org/apache/wayang/jdbc/operators/SqlToRddOperatorTest.java index 1acfd426d..f1f50d297 100644 --- a/wayang-platforms/wayang-jdbc-template/src/test/java/org/apache/wayang/jdbc/operators/SqlToRddOperatorTest.java +++ b/wayang-platforms/wayang-jdbc-template/src/test/java/org/apache/wayang/jdbc/operators/SqlToRddOperatorTest.java @@ -165,4 +165,25 @@ void testWithEmptyHsqldb() throws SQLException { assertTrue(output.isEmpty()); } + @Test + void testRowToRecord() { + org.apache.spark.sql.Row row = org.apache.spark.sql.RowFactory.create(1, "test", 42.0); + Record record = SqlToRddOperator.rowToRecord(row); + assertEquals(3, record.size()); + assertEquals(1, record.getField(0)); + assertEquals("test", record.getField(1)); + assertEquals(42.0, record.getField(2)); + } + + @Test + void testJsonSerialization() { + SqlToRddOperator operator = new SqlToRddOperator(HsqldbPlatform.getInstance()); + org.apache.wayang.core.util.json.WayangJsonObj json = operator.toJson(); + assertEquals(HsqldbPlatform.class.getCanonicalName(), json.getString("platform")); + + SqlToRddOperator deserialized = SqlToRddOperator.fromJson(json); + assertEquals(operator.getInputType(), deserialized.getInputType()); + assertEquals(operator.getOutputType(), deserialized.getOutputType()); + } + } diff --git a/wayang-platforms/wayang-spark/pom.xml b/wayang-platforms/wayang-spark/pom.xml index 34b0cbf29..2460c1630 100644 --- a/wayang-platforms/wayang-spark/pom.xml +++ b/wayang-platforms/wayang-spark/pom.xml @@ -77,8 +77,8 @@ org.apache.spark - spark-core_2.12 - 3.5.7 + spark-core_${scala.mayor.version} + ${spark.version} org.xerial.snappy @@ -88,18 +88,18 @@ org.apache.spark - spark-graphx_2.12 - 3.4.4 + spark-graphx_${scala.mayor.version} + ${spark.version} org.apache.spark - spark-sql_2.12 - 3.4.4 + spark-sql_${scala.mayor.version} + ${spark.version} org.apache.spark - spark-mllib_2.12 - 3.4.4 + spark-mllib_${scala.mayor.version} + ${spark.version} @@ -117,7 +117,7 @@ org.antlr antlr4-runtime - 4.8 + ${antlr.version}