Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
## SQRL
##############################
plan-output/
sqrl_iceberg_data/

##############################
## Java
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ public enum EngineType {
SERVER,
LOG,
QUERY,
SHALLOW_QUERY,
EXPORT;

public boolean isWrite() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,11 @@ public enum JdbcDialect {
SQLServer,
H2,
SQLite,
Iceberg,
Snowflake,
DuckDB;
Iceberg;

private final String[] synonyms;

private JdbcDialect(String... synonyms) {
JdbcDialect(String... synonyms) {
this.synonyms = synonyms;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,14 +36,19 @@ public List<ObjectNode> convertConfigsToJson() {

for (var engine : enabledQueryEngines) {
var queryEngine = (QueryEngine) engine;
var engineConf = packageJson.getEngines().getEngineConfig(queryEngine.getName()).get();
var engineConf = packageJson.getEngines().getEngineConfig(queryEngine.getName());
if (engineConf.isEmpty()) {
continue;
}

if (engineConf instanceof EngineConfigImpl impl) {
if (engineConf.get() instanceof EngineConfigImpl impl) {
var engineConfigMap = impl.sqrlConfig.toMap();

var rootNode = JsonUtils.MAPPER.createObjectNode();
var configNode = JsonUtils.MAPPER.valueToTree(engineConfigMap);
rootNode.set(queryEngine.serverConfigName(), configNode);
queryEngine
.serverConfigName()
.ifPresent(engineConfigName -> rootNode.set(engineConfigName, configNode));

convertedConfigs.add(rootNode);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
import com.datasqrl.engine.EnginePhysicalPlan;
import com.datasqrl.engine.ExecutionEngine;
import com.datasqrl.planner.dag.plan.MaterializationStagePlan;
import java.util.Optional;

/**
* A {@link QueryEngine} executes queries against a {@link DatabaseEngine} that supports the query
Expand All @@ -28,5 +29,5 @@ public interface QueryEngine extends ExecutionEngine {

EnginePhysicalPlan plan(MaterializationStagePlan stagePlan);

String serverConfigName();
Optional<String> serverConfigName();
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
/*
* Copyright © 2021 DataSQRL (contact@datasqrl.com)
*
* Licensed 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 com.datasqrl.engine.database.relational;

import static com.datasqrl.engine.EngineFeature.STANDARD_QUERY;

import com.datasqrl.config.ConnectorFactoryFactory;
import com.datasqrl.config.EngineType;
import com.datasqrl.config.PackageJson.EngineConfig;
import com.datasqrl.engine.database.QueryEngine;
import com.datasqrl.planner.tables.FlinkTableBuilder;
import com.datasqrl.server.jdbc.DatabaseType;
import java.util.Optional;
import lombok.NonNull;

/** Abstract implementation of a relational {@link QueryEngine}. */
public abstract class AbstractJDBCShallowQueryEngine extends AbstractJDBCEngine
implements QueryEngine {

protected AbstractJDBCShallowQueryEngine(
String name, @NonNull EngineConfig engineConfig, ConnectorFactoryFactory connectorFactory) {
super(name, EngineType.SHALLOW_QUERY, STANDARD_QUERY, engineConfig, connectorFactory);
}

@Override
public final Optional<String> serverConfigName() {
return Optional.empty();
}

@Override
protected final DatabaseType getDatabaseType() {
return DatabaseType.NONE;
}

@Override
protected String getConnectorTableName(FlinkTableBuilder tableBuilder) {
throw new UnsupportedOperationException();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import com.datasqrl.config.PackageJson;
import com.datasqrl.server.jdbc.DatabaseType;
import jakarta.inject.Inject;
import java.util.Optional;
import lombok.NonNull;

public class DuckDBEngine extends AbstractJDBCQueryEngine {
Expand All @@ -33,8 +34,8 @@ public DuckDBEngine(@NonNull PackageJson json, ConnectorFactoryFactory connector
}

@Override
public String serverConfigName() {
return "duckDbConfig";
public Optional<String> serverConfigName() {
return Optional.of("duckDbConfig");
}

@Override
Expand All @@ -49,6 +50,6 @@ protected DatabaseType getDatabaseType() {

@Override
public JdbcStatementFactory getStatementFactory() {
return new DuckDbStatementFactory(engineConfig);
return new DuckDbStatementFactory();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,6 @@
import com.datasqrl.calcite.Dialect;
import com.datasqrl.calcite.dialect.DuckDbSqlDialect;
import com.datasqrl.calcite.type.TypeFactory;
import com.datasqrl.config.JdbcDialect;
import com.datasqrl.config.PackageJson.EngineConfig;
import com.datasqrl.engine.database.relational.ddl.GenericCreateTableDdlFactory;
import com.datasqrl.plan.global.IndexDefinition;
import com.datasqrl.planner.dag.plan.MaterializationStagePlan.Query;
Expand All @@ -51,18 +49,10 @@

public class DuckDbStatementFactory extends AbstractJdbcStatementFactory {

private final EngineConfig engineConfig;

public DuckDbStatementFactory(EngineConfig engineConfig) {
public DuckDbStatementFactory() {
super(
Dialect.DUCKDB,
new GenericCreateTableDdlFactory()); // Iceberg creates the tables, DuckDB only queries
this.engineConfig = engineConfig;
}

@Override
public JdbcDialect getDialect() {
return JdbcDialect.DuckDB;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
import com.datasqrl.calcite.Dialect;
import com.datasqrl.calcite.OperatorRuleTransformer;
import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect;
import com.datasqrl.config.JdbcDialect;
import com.datasqrl.engine.database.relational.ddl.IcebergCreateTableDdlFactory;
import com.datasqrl.plan.global.IndexDefinition;
import com.datasqrl.planner.hint.DataTypeHint;
Expand All @@ -41,11 +40,6 @@ protected SqlDataTypeSpec getSqlType(RelDataType type, Optional<DataTypeHint> hi
return ExtendedPostgresSqlDialect.DEFAULT.getCastSpec(type);
}

@Override
public JdbcDialect getDialect() {
return JdbcDialect.Postgres;
}

@Override
public boolean supportsQueries() {
return false;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,6 @@
*/
package com.datasqrl.engine.database.relational;

import com.datasqrl.config.JdbcDialect;
import com.datasqrl.plan.global.IndexDefinition;
import com.datasqrl.planner.dag.plan.MaterializationStagePlan.Query;
import java.util.Collection;
Expand All @@ -24,8 +23,6 @@

public interface JdbcStatementFactory {

JdbcDialect getDialect();

JdbcStatement createTable(JdbcEngineCreateTable createTable);

default List<JdbcStatement> applyTableExtensions(Collection<CreateTableJdbcStatement> tables) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@

import com.datasqrl.calcite.Dialect;
import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect;
import com.datasqrl.config.JdbcDialect;
import com.datasqrl.config.PackageJson.EngineConfig;
import com.datasqrl.deployment.model.JdbcStatementModel.PartitionType;
import com.datasqrl.deployment.model.JdbcStatementModel.Type;
Expand Down Expand Up @@ -64,11 +63,6 @@ public PostgresStatementFactory(int partitionTtlDivisor) {
this.partitionTtlDivisor = partitionTtlDivisor;
}

@Override
public JdbcDialect getDialect() {
return JdbcDialect.Postgres;
}

@Override
protected SqlDataTypeSpec getSqlType(RelDataType type, Optional<DataTypeHint> hint) {
SqlDataTypeSpec spec = ExtendedPostgresSqlDialect.DEFAULT.getCastSpec(type);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,10 @@
import com.datasqrl.config.ConnectorFactoryFactory;
import com.datasqrl.config.JdbcDialect;
import com.datasqrl.config.PackageJson;
import com.datasqrl.server.jdbc.DatabaseType;
import jakarta.inject.Inject;
import lombok.NonNull;

public class SnowflakeEngine extends AbstractJDBCQueryEngine {
public class SnowflakeEngine extends AbstractJDBCShallowQueryEngine {

@Inject
public SnowflakeEngine(@NonNull PackageJson json, ConnectorFactoryFactory connectorFactory) {
Expand All @@ -32,19 +31,9 @@ public SnowflakeEngine(@NonNull PackageJson json, ConnectorFactoryFactory connec
connectorFactory);
}

@Override
public String serverConfigName() {
return "snowflakeConfig";
}

@Override
protected JdbcDialect getDialect() {
return JdbcDialect.Snowflake;
}

@Override
protected DatabaseType getDatabaseType() {
return DatabaseType.SNOWFLAKE;
return JdbcDialect.Iceberg;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
import com.datasqrl.calcite.Dialect;
import com.datasqrl.calcite.dialect.ExtendedPostgresSqlDialect;
import com.datasqrl.calcite.dialect.snowflake.SqlCreateIcebergTableFromObjectStorage;
import com.datasqrl.config.JdbcDialect;
import com.datasqrl.config.PackageJson.EngineConfig;
import com.datasqrl.deployment.model.JdbcStatementModel.Type;
import com.datasqrl.engine.database.relational.ddl.GenericCreateTableDdlFactory;
Expand All @@ -40,11 +39,6 @@ public SnowflakeStatementFactory(EngineConfig engineConfig) {
this.engineConfig = engineConfig;
}

@Override
public JdbcDialect getDialect() {
return JdbcDialect.Snowflake;
}

@Override
public JdbcStatement createTable(JdbcEngineCreateTable createTable) {
var tableName = createTable.tableName();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;

/**
* A simple pipeline that has a single stream, log, and server engine with support for multiple
Expand All @@ -38,8 +39,8 @@ public record SimplePipeline(
HashMultimap<ExecutionStage, ExecutionStage> downstream)
implements ExecutionPipeline {

private static final List<String> AVAILABLE_QUERY_ENGINES =
EngineUtil.getAvailableQueryEngineNames();
private static final String QUERY_ENGINE_NAMES =
EngineUtil.formatEngineNames(EngineUtil.getAvailableQueryEngines());

public static SimplePipeline of(Map<String, ExecutionEngine> engines, ErrorCollector errors) {
var upstream = HashMultimap.<ExecutionStage, ExecutionStage>create();
Expand Down Expand Up @@ -76,8 +77,28 @@ public static SimplePipeline of(Map<String, ExecutionEngine> engines, ErrorColle
streamStage.ifPresent(ss -> upstream.put(dbStage, ss));
serverStage.ifPresent(vs -> downstream.put(dbStage, vs));

// Make sure if server is present, then a non-view query engine is also present
if (serverStage.isPresent() && dbStage.engine() instanceof AbstractJDBCTableFormatEngine) {
validatePipelineForQueryEngine(dbStage.name(), engines, errors);
var queryStages = getStage(EngineType.QUERY, engines);
var shallowQueryStages = getStage(EngineType.SHALLOW_QUERY, engines);

if (queryStages.isEmpty() && !shallowQueryStages.isEmpty()) {
var shallowQueryEngines =
shallowQueryStages.stream()
.map(EngineStage::name)
.map(s -> '\'' + s + '\'')
.collect(Collectors.joining(", "));

errors.fatal(
"When '%s' is enabled as a server, '%s' cannot use shallow query engines (%s) to process server queries because they are not integrated at the database level. Available query engines: %s",
serverStage.get().name(), dbStage.name(), shallowQueryEngines, QUERY_ENGINE_NAMES);
}

if (queryStages.isEmpty()) {
errors.fatal(
"When '%s' is enabled as a server, '%s' requires a query engine to process server queries, but none are listed under 'enabled-engines'. Available query engines: %s",
serverStage.get().name(), dbStage.name(), QUERY_ENGINE_NAMES);
}
}
}

Expand Down Expand Up @@ -129,18 +150,6 @@ private static Optional<EngineStage> getSingleStage(
"Expected a single %s engine but found multiple: %s".formatted(engineType, engineList));
}

private static void validatePipelineForQueryEngine(
String tableFormatEngineName, Map<String, ExecutionEngine> engines, ErrorCollector errors) {
var queryStages = getStage(EngineType.QUERY, engines);
if (!queryStages.isEmpty()) {
return;
}

errors.fatal(
"Engine '%s' requires a query engine, but none are listed under 'enabled-engines'. Available options: %s",
tableFormatEngineName, AVAILABLE_QUERY_ENGINES);
}

@Override
public Set<ExecutionStage> getUpStreamFrom(ExecutionStage stage) {
Preconditions.checkArgument(upstream.containsKey(stage), "Invalid stage: %s", stage);
Expand Down
22 changes: 9 additions & 13 deletions sqrl-planner/src/main/java/com/datasqrl/util/EngineUtil.java
Original file line number Diff line number Diff line change
Expand Up @@ -17,28 +17,24 @@

import com.datasqrl.config.EngineFactory;
import com.datasqrl.engine.IExecutionEngine;
import com.datasqrl.engine.database.DatabaseEngine;
import com.datasqrl.engine.database.QueryEngine;
import java.util.List;
import java.util.Set;
import com.datasqrl.engine.database.relational.AbstractJDBCQueryEngine;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import lombok.AccessLevel;
import lombok.NoArgsConstructor;

@NoArgsConstructor(access = AccessLevel.PRIVATE)
public final class EngineUtil {

public static List<String> getAvailableDatabaseEngineNames(String... exclusions) {
var exclusionSet = Set.of(exclusions);

return getChildEngineFactories(DatabaseEngine.class)
.map(EngineFactory::getEngineName)
.filter(name -> !exclusionSet.contains(name))
.toList();
public static Stream<EngineFactory> getAvailableQueryEngines() {
return getChildEngineFactories(AbstractJDBCQueryEngine.class);
}

public static List<String> getAvailableQueryEngineNames() {
return getChildEngineFactories(QueryEngine.class).map(EngineFactory::getEngineName).toList();
public static String formatEngineNames(Stream<EngineFactory> engines) {
return engines
.map(EngineFactory::getEngineName)
.map(name -> '\'' + name + '\'')
.collect(Collectors.joining(", "));
}

public static Stream<EngineFactory> getChildEngineFactories(
Expand Down
Loading