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}