diff --git a/backend/pom.xml b/backend/pom.xml index 8ccf72de52..c766af7f29 100644 --- a/backend/pom.xml +++ b/backend/pom.xml @@ -348,15 +348,10 @@ jooq 3.20.2 - - org.jooq - jooq-postgres-extensions - 3.20.2 - com.zaxxer HikariCP - 5.1.0 + 7.1.0 org.testcontainers @@ -368,11 +363,6 @@ testcontainers-junit-jupiter test - - org.testcontainers - testcontainers-postgresql - test - org.testcontainers testcontainers-clickhouse diff --git a/backend/src/main/java/com/bakdata/conquery/commands/ManagerNode.java b/backend/src/main/java/com/bakdata/conquery/commands/ManagerNode.java index 94ac5437b8..3b093cf73d 100644 --- a/backend/src/main/java/com/bakdata/conquery/commands/ManagerNode.java +++ b/backend/src/main/java/com/bakdata/conquery/commands/ManagerNode.java @@ -73,22 +73,19 @@ public ManagerNode(@NonNull String name) { public void run(Manager manager) throws InterruptedException { Environment environment = manager.getEnvironment(); ConqueryConfig config = manager.getConfig(); - validator = environment.getValidator(); + validator = environment.getValidator(); this.manager = manager; final ObjectMapper apiObjectMapper = environment.getObjectMapper(); getInternalMapperFactory().customizeApiObjectMapper(apiObjectMapper, getDatasetRegistry(), getMetaStorage()); - // FormScanner needs to be instantiated before plugins are initialized formScanner = new FormScanner(config); - // Init all plugins config.getPlugins().forEach(pluginConfig -> pluginConfig.initialize(this)); - // Initialization of internationalization I18n.init(); @@ -100,10 +97,6 @@ public void run(Manager manager) throws InterruptedException { environment.lifecycle().manage(this); - loadNamespaces(); - - loadMetaStorage(); - // Create AdminServlet first to make it available to the realms admin = new AdminServlet(this); @@ -159,7 +152,6 @@ public void loadNamespaces() { }); } - loaders.shutdown(); while (!loaders.awaitTermination(1, TimeUnit.MINUTES)) { final int countLoaded = registry.getNamespaces().size(); @@ -201,6 +193,9 @@ private void registerTasks(Manager manager, Environment environment, ConqueryCon @Override public void start() throws Exception { + loadNamespaces(); + loadMetaStorage(); + manager.start(); } @@ -214,7 +209,6 @@ public void stop() throws Exception { catch (Exception e) { log.error("{} could not be closed", provider, e); } - } try { diff --git a/backend/src/main/java/com/bakdata/conquery/mode/local/LocalNamespaceHandler.java b/backend/src/main/java/com/bakdata/conquery/mode/local/LocalNamespaceHandler.java index fced1ca2da..8a4a292b61 100644 --- a/backend/src/main/java/com/bakdata/conquery/mode/local/LocalNamespaceHandler.java +++ b/backend/src/main/java/com/bakdata/conquery/mode/local/LocalNamespaceHandler.java @@ -1,7 +1,5 @@ package com.bakdata.conquery.mode.local; -import java.time.Clock; - import com.bakdata.conquery.io.storage.MetaStorage; import com.bakdata.conquery.io.storage.NamespaceStorage; import com.bakdata.conquery.mode.NamespaceHandler; @@ -13,7 +11,6 @@ import com.bakdata.conquery.models.worker.DatasetRegistry; import com.bakdata.conquery.models.worker.LocalNamespace; import com.bakdata.conquery.sql.conquery.SqlExecutionManager; -import com.bakdata.conquery.sql.conquery.SqlMatchingStats; import com.bakdata.conquery.sql.conversion.NodeConversions; import com.bakdata.conquery.sql.conversion.SqlConverter; import com.bakdata.conquery.sql.conversion.dialect.DialectBundle; @@ -24,6 +21,8 @@ import lombok.extern.slf4j.Slf4j; import org.jooq.DSLContext; +import java.time.Clock; + @RequiredArgsConstructor @Slf4j public class LocalNamespaceHandler implements NamespaceHandler { @@ -43,29 +42,34 @@ public LocalNamespace createNamespace( NamespaceSetupData namespaceData = NamespaceHandler.createNamespaceSetup(namespaceStorage, config, internalMapperFactory, datasetRegistry, environment); ManagedConnection connection = connectionManager.getConnection(namespaceStorage.getDataset()); - DSLContext dslContext = connection.connect(); - DialectBundle dialectBundle = connection.getConnection().getDialect().getDialectBundle(); + try { + DSLContext dslContext = connection.connect(); + DialectBundle dialectBundle = connection.getConnection().getDialect().getDialectBundle(); - ResultSetProcessor resultSetProcessor = dialectBundle.getResultSetProcessor(config); - SqlExecutionService sqlExecutionService = new SqlExecutionService(dslContext, resultSetProcessor); + ResultSetProcessor resultSetProcessor = dialectBundle.getResultSetProcessor(config); + SqlExecutionService sqlExecutionService = new SqlExecutionService(dslContext, resultSetProcessor); - NodeConversions nodeConversions = new NodeConversions(config.getIdColumns(), dialectBundle, dslContext, sqlExecutionService, clock); - SqlConverter sqlConverter = new SqlConverter(nodeConversions, config); - ExecutionManager executionManager = new SqlExecutionManager(sqlConverter, sqlExecutionService, metaStorage, datasetRegistry, config); - SqlStorageHandler sqlStorageHandler = new SqlStorageHandler(sqlExecutionService); - SqlEntityResolver sqlEntityResolver = new SqlEntityResolver(config.getIdColumns(), dslContext, dialectBundle, sqlExecutionService); + NodeConversions nodeConversions = new NodeConversions(config.getIdColumns(), dialectBundle, dslContext, sqlExecutionService, clock, connection.getConnection().getPrimaryColumn()); + SqlConverter sqlConverter = new SqlConverter(nodeConversions, config); + ExecutionManager executionManager = new SqlExecutionManager(sqlConverter, sqlExecutionService, metaStorage, datasetRegistry, config); + SqlStorageHandler sqlStorageHandler = new SqlStorageHandler(sqlExecutionService); + SqlEntityResolver sqlEntityResolver = new SqlEntityResolver(config.getIdColumns(), dslContext, dialectBundle, sqlExecutionService); - return new LocalNamespace( - dialectBundle, - namespaceData.preprocessMapper(), - namespaceStorage, - executionManager, - dslContext, sqlStorageHandler, - namespaceData.jobManager(), - namespaceData.filterSearch(), - sqlEntityResolver, - connection.getConnection() - ); + return new LocalNamespace( + dialectBundle, + namespaceData.preprocessMapper(), + namespaceStorage, + executionManager, + dslContext, sqlStorageHandler, + namespaceData.jobManager(), + namespaceData.filterSearch(), + sqlEntityResolver, + connection.getConnection() + ); + } catch (Exception e) { + log.error("Failed to load namespaceStorage for {}", namespaceStorage.getPathName(), e); + throw e; + } } @Override diff --git a/backend/src/main/java/com/bakdata/conquery/mode/local/ManagedConnection.java b/backend/src/main/java/com/bakdata/conquery/mode/local/ManagedConnection.java index e0fe124c8e..0e1bf735f7 100644 --- a/backend/src/main/java/com/bakdata/conquery/mode/local/ManagedConnection.java +++ b/backend/src/main/java/com/bakdata/conquery/mode/local/ManagedConnection.java @@ -31,18 +31,15 @@ public class ManagedConnection implements Managed { @Override public void start() throws Exception { dataSource = connection.createDataSource(healthCheckRegistry); - try { log.debug("TEST connecting to {}", connection.getJdbcConnectionUrl()); - if (dataSource.getConnection().isValid(100)) { + if (dataSource.getConnection().isValid(1000)) { log.info("SUCCESS connecting to {}", connection.getJdbcConnectionUrl()); - } - else { + } else { log.error("FAILED connecting to {}. Connection did not become valid.", connection.getJdbcConnectionUrl()); } - } - catch (SQLException exception) { - log.error("FAILED connecting to {}", connection.getJdbcConnectionUrl(), exception); + } catch (SQLException exception) { + throw new RuntimeException("FAILED connecting to %s".formatted(connection.getJdbcConnectionUrl()), exception); } } diff --git a/backend/src/main/java/com/bakdata/conquery/mode/local/UpdateMatchingStatsSqlJob.java b/backend/src/main/java/com/bakdata/conquery/mode/local/UpdateMatchingStatsSqlJob.java index 28d603d1de..945e1171ac 100644 --- a/backend/src/main/java/com/bakdata/conquery/mode/local/UpdateMatchingStatsSqlJob.java +++ b/backend/src/main/java/com/bakdata/conquery/mode/local/UpdateMatchingStatsSqlJob.java @@ -1,78 +1,95 @@ package com.bakdata.conquery.mode.local; -import java.util.ArrayList; +import java.util.Collection; import java.util.List; -import java.util.concurrent.*; +import java.util.Map; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import com.bakdata.conquery.models.datasets.Dataset; import com.bakdata.conquery.models.datasets.concepts.Concept; import com.bakdata.conquery.models.datasets.concepts.tree.TreeConcept; +import com.bakdata.conquery.models.identifiable.ids.specific.ConceptId; import com.bakdata.conquery.models.jobs.Job; import com.bakdata.conquery.sql.conquery.SqlMatchingStats; import com.google.common.base.Stopwatch; -import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.common.util.concurrent.ListeningExecutorService; import com.google.common.util.concurrent.MoreExecutors; import lombok.Data; +import lombok.EqualsAndHashCode; import lombok.ToString; import lombok.extern.slf4j.Slf4j; -import org.apache.commons.lang3.builder.ToStringExclude; -import org.checkerframework.checker.nullness.qual.Nullable; +import org.apache.commons.collections4.map.HashedMap; @Slf4j @Data +@EqualsAndHashCode(callSuper = false) public class UpdateMatchingStatsSqlJob extends Job { - @ToString.Exclude - private final List> concepts; - private final Dataset dataset; + @ToString.Exclude + private final List> concepts; + private final Dataset dataset; - @ToString.Exclude - private final SqlMatchingStats matchingStats; + @ToString.Exclude + private final SqlMatchingStats matchingStats; - @Override - public void execute() throws Exception { + @Override + public void execute() throws Exception { - log.info("BEGIN collecting SQL matching stats for {}", dataset); + log.info("BEGIN collecting SQL matching stats for {}", dataset); - Stopwatch stopwatch = Stopwatch.createStarted(); + Stopwatch stopwatch = Stopwatch.createStarted(); - ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newVirtualThreadPerTaskExecutor()); + ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(getMatchingStats().getMatchingStatsWorkers())); - List> jobs = new ArrayList<>(); + Map> jobsByConcept = new HashedMap<>(); + Collection> jobs = jobsByConcept.values(); - for (Concept concept : concepts) { - if (!(concept instanceof TreeConcept)) { - continue; - } - jobs.add(matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService)); - } + for (Concept concept : concepts) { + if (concept instanceof TreeConcept) { + ListenableFuture job = matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, getMatchingStats().getMatchingStatsRetries()); - ListenableFuture> all = Futures.allAsList(jobs); + job.addListener( + () -> { + if (job.state().equals(Future.State.FAILED)) { + log.warn("FAILED to collect SQL matching stats for {}", concept, job.exceptionNow()); + } + }, MoreExecutors.directExecutor()); - while (!all.isDone()) { - if (isCancelled()) { - all.cancel(true); - log.debug("CANCELLED update matching stats for {}", getDataset(), all.exceptionNow()); - return; - } + jobsByConcept.put(concept.getId(), job); + } + } - all.get(5, TimeUnit.SECONDS); - log.trace("WAITING for matching stats to finish {}", getDataset()); + while (jobs.stream().anyMatch(job -> job.state().equals(Future.State.RUNNING))) { + if (isCancelled()) { + for (ListenableFuture job : jobs) { + job.cancel(true); + } + } - if (all.state().equals(Future.State.FAILED)) { - log.error("FAILED update matching stats for {}", getDataset(), all.exceptionNow()); - return; - } - } + for (ListenableFuture someJob : jobs) { + if (someJob.isDone()) { + continue; + } - log.debug("DONE collecting SQL matching stats for {} within {}", dataset, stopwatch); - } + try { + someJob.get(30, TimeUnit.SECONDS); + } catch (Exception e) { + // intentionally left blank + } - @Override - public String getLabel() { - return "Collect matching stats for %s (%s concepts)".formatted(dataset.getName(), concepts.size()); - } + log.debug("WAITING for {} matching stats to finish.", jobs.stream().filter(job -> job.state().equals(Future.State.RUNNING)).count()); + } + } + + log.debug("DONE collecting SQL matching stats for {} within {}", dataset, stopwatch); + } + + @Override + public String getLabel() { + return "Collect matching stats for %s (%s concepts)".formatted(dataset.getName(), concepts.size()); + } } diff --git a/backend/src/main/java/com/bakdata/conquery/models/config/DatabaseConnectionConfig.java b/backend/src/main/java/com/bakdata/conquery/models/config/DatabaseConnectionConfig.java index 61a152bb3d..2c5cfb3fe6 100644 --- a/backend/src/main/java/com/bakdata/conquery/models/config/DatabaseConnectionConfig.java +++ b/backend/src/main/java/com/bakdata/conquery/models/config/DatabaseConnectionConfig.java @@ -4,6 +4,7 @@ import com.zaxxer.hikari.HikariConfig; import com.zaxxer.hikari.HikariDataSource; import io.dropwizard.util.Duration; +import jakarta.validation.constraints.Min; import jakarta.validation.constraints.NotNull; import lombok.*; import lombok.extern.jackson.Jacksonized; @@ -30,6 +31,20 @@ public class DatabaseConnectionConfig { @NotNull private Dialect dialect; + /** + * Maximum workers to schedule for matching stats. This depends on your HikariCP configuration as well. + */ + @Min(1) + @Builder.Default + private int matchingStatsWorkers = 5; + + /** + * Retries for matching stats. This is a bit of a workaround because some DBMS seem to struggle to properly communicate workload to HikariCP. + */ + @Min(1) + @Builder.Default + private int matchingStatsRetries = 3; + /** * Name of the column which is shared among the table and all aggregations are grouped by. */ diff --git a/backend/src/main/java/com/bakdata/conquery/models/messages/namespaces/specific/UpdateMatchingStatsMessage.java b/backend/src/main/java/com/bakdata/conquery/models/messages/namespaces/specific/UpdateMatchingStatsMessage.java index 7a6927360a..8b721e7b0c 100644 --- a/backend/src/main/java/com/bakdata/conquery/models/messages/namespaces/specific/UpdateMatchingStatsMessage.java +++ b/backend/src/main/java/com/bakdata/conquery/models/messages/namespaces/specific/UpdateMatchingStatsMessage.java @@ -33,6 +33,7 @@ import com.bakdata.conquery.util.progressreporter.ProgressReporter; import com.fasterxml.jackson.annotation.JsonCreator; import com.google.common.base.Functions; +import lombok.EqualsAndHashCode; import lombok.Getter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; @@ -45,6 +46,7 @@ @CPSType(id = "UPDATE_MATCHING_STATS", base = NamespacedMessage.class) @Slf4j @RequiredArgsConstructor(onConstructor_ = {@JsonCreator}) +@EqualsAndHashCode(callSuper = false) public class UpdateMatchingStatsMessage extends WorkerMessage { @Getter diff --git a/backend/src/main/java/com/bakdata/conquery/models/worker/LocalNamespace.java b/backend/src/main/java/com/bakdata/conquery/models/worker/LocalNamespace.java index 6ba57c5127..30d79c6da5 100644 --- a/backend/src/main/java/com/bakdata/conquery/models/worker/LocalNamespace.java +++ b/backend/src/main/java/com/bakdata/conquery/models/worker/LocalNamespace.java @@ -1,9 +1,5 @@ package com.bakdata.conquery.models.worker; -import java.util.Set; -import java.util.stream.Collectors; -import java.util.stream.Stream; - import com.bakdata.conquery.io.storage.NamespaceStorage; import com.bakdata.conquery.mode.local.SqlEntityResolver; import com.bakdata.conquery.mode.local.SqlStorageHandler; @@ -20,6 +16,10 @@ import lombok.extern.slf4j.Slf4j; import org.jooq.DSLContext; +import java.util.Set; +import java.util.stream.Collectors; +import java.util.stream.Stream; + @Getter @Slf4j public class LocalNamespace extends Namespace { @@ -41,10 +41,11 @@ public LocalNamespace( SqlEntityResolver sqlEntityResolver, DatabaseConnectionConfig databaseConfig ) { super(preprocessMapper, storage, executionManager, jobManager, filterSearch, sqlEntityResolver); + this.dslContext = dslContext; this.storageHandler = storageHandler; this.dialect = dialect; - matchingStats = new SqlMatchingStats(dslContext, dialect.getFunctionProvider(), databaseConfig.getPrimaryColumn()); + this.matchingStats = new SqlMatchingStats(dslContext, dialect.getFunctionProvider(), databaseConfig.getPrimaryColumn(), databaseConfig.getMatchingStatsWorkers(), databaseConfig.getMatchingStatsRetries()); } @@ -61,8 +62,7 @@ void registerColumnValuesInSearch(Set columns) { try { final Stream stringStream = storageHandler.lookupColumnValues(getStorage(), column); getFilterSearch().registerValues(column, stringStream.collect(Collectors.toSet())); - } - catch (Exception e) { + } catch (Exception e) { log.error("Problem collecting column values for {}", column, e); } } diff --git a/backend/src/main/java/com/bakdata/conquery/sql/conquery/SqlMatchingStats.java b/backend/src/main/java/com/bakdata/conquery/sql/conquery/SqlMatchingStats.java index a98646b38f..e35d7babd7 100644 --- a/backend/src/main/java/com/bakdata/conquery/sql/conquery/SqlMatchingStats.java +++ b/backend/src/main/java/com/bakdata/conquery/sql/conquery/SqlMatchingStats.java @@ -1,23 +1,5 @@ package com.bakdata.conquery.sql.conquery; -import static org.jooq.impl.DSL.*; - -import java.sql.Date; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.CompletionStage; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Future; - -import com.google.common.util.concurrent.ListenableFuture; -import com.google.common.util.concurrent.ListeningExecutorService; -import jakarta.validation.constraints.NotBlank; - import com.bakdata.conquery.models.common.daterange.CDateRange; import com.bakdata.conquery.models.datasets.Column; import com.bakdata.conquery.models.datasets.concepts.ConceptElement; @@ -27,7 +9,6 @@ import com.bakdata.conquery.models.datasets.concepts.conditions.CTCondition; import com.bakdata.conquery.models.datasets.concepts.tree.ConceptTreeChild; import com.bakdata.conquery.models.datasets.concepts.tree.TreeConcept; -import com.bakdata.conquery.models.events.MajorTypeId; import com.bakdata.conquery.models.identifiable.ids.specific.ConceptElementId; import com.bakdata.conquery.models.identifiable.ids.specific.ConceptId; import com.bakdata.conquery.sql.conversion.cqelement.concept.CTConditionContext; @@ -35,392 +16,406 @@ import com.bakdata.conquery.util.TablePrimaryColumnUtil; import com.google.common.base.Stopwatch; import com.google.common.collect.Sets; +import com.google.common.util.concurrent.ListenableFuture; +import com.google.common.util.concurrent.ListeningExecutorService; +import jakarta.validation.constraints.NotBlank; import lombok.Data; import lombok.extern.slf4j.Slf4j; import org.jetbrains.annotations.NotNull; -import org.jooq.CommonTableExpression; -import org.jooq.Condition; -import org.jooq.CreateTableElementListStep; -import org.jooq.Cursor; -import org.jooq.DSLContext; -import org.jooq.Field; -import org.jooq.InsertValuesStepN; -import org.jooq.Name; -import org.jooq.Param; +import org.jooq.*; import org.jooq.Record; -import org.jooq.Record4; -import org.jooq.RowN; -import org.jooq.Select; -import org.jooq.SelectConditionStep; -import org.jooq.SelectJoinStep; import org.jooq.exception.DataAccessException; +import java.sql.Date; +import java.util.*; + +import static org.jooq.impl.DSL.*; + @Slf4j @Data public class SqlMatchingStats { - private final Field PID_FIELD = field(name("pid"), String.class); - private final Field LB_FIELD = field(name("lower_bound"), Date.class); - private final Field UB_FIELD = field(name("upper_bound"), Date.class); - private final Field CONCEPT_ID_FIELD = field(name("resolved_id"), Integer.class); - private final Set> NULL_PARAMS = Collections.singleton(inline(null, String.class)); - - private final DSLContext dslContext; - private final SqlFunctionProvider functionProvider; - private final String defaultPrimaryColumn; - private final int fetchBatchSize = 100; - - private static void assignStatsToPath(ConceptElement element, Map, MatchingStats.Entry> matchingStats, String entity, CDateRange span) { - ConceptElementId id = element.getId(); - - while (element != null) { - matchingStats.computeIfAbsent(id, (ignored) -> new MatchingStats.Entry()).addEvents(entity, 1, span); - element = element.getParent(); - } - } - - /** - * collect unique fields used/defined in the expressions. - */ - private static List> collectAllFields(List conceptConditions) { - List> fields = conceptConditions.stream().flatMap(e -> e.conditions().keySet().stream()).distinct().toList(); - return fields; - } - - private static Select unionSelects(List> connectorTableSelects) { - Select unioned = null; - - for (Select connectorTable : connectorTableSelects) { - if (unioned == null) { - unioned = (Select) connectorTable; - continue; - } - - unioned = unioned.unionAll(connectorTable); - } - - - return unioned; - } + /** + * Legacy backend implementation requires separation by source (in that case different shards). For sql it's a constant. + */ + private static final String SQL_SOURCE_MATCHING_STATS_LABEL = "sql"; + + private static final Field PID_FIELD = field(name("pid"), String.class); + private static final Field LB_FIELD = field(name("lower_bound"), Date.class); + private static final Field UB_FIELD = field(name("upper_bound"), Date.class); + private static final Field CONCEPT_ID_FIELD = field(name("resolved_id"), Integer.class); + private static final Set> NULL_PARAMS = Collections.singleton(inline(null, String.class)); + + private final DSLContext dslContext; + private final SqlFunctionProvider functionProvider; + private final String defaultPrimaryColumn; + private final int fetchBatchSize = 100; + private final int matchingStatsWorkers; + private final int matchingStatsRetries; + + private static void assignStatsToPath(ConceptElement element, Map, MatchingStats.Entry> matchingStats, String entity, CDateRange span) { + while (element != null) { + ConceptElementId id = element.getId(); + + matchingStats.computeIfAbsent(id, (ignored) -> new MatchingStats.Entry()) + .addEvents(entity, 1, span); + element = element.getParent(); + } + } - private static Param defaultValue(Field field) { - if (field.getDataType().isBoolean()) { - return inline(false); - } + /** + * collect unique fields used/defined in the expressions. + */ + private static List> collectReferencedFields(List conceptConditions) { + List> fields = conceptConditions.stream().flatMap(e -> e.conditions().keySet().stream()).distinct().toList(); + return fields; + } + + private static Select unionSelects(List> connectorTableSelects) { + Select unioned = null; + + for (Select connectorTable : connectorTableSelects) { + if (unioned == null) { + unioned = (Select) connectorTable; + continue; + } - if (field.getDataType().isString()) { - return inline(null, String.class); - } + unioned = unioned.unionAll(connectorTable); + } - throw new IllegalStateException("Fields of type %s are not expected".formatted(field.getDataType())); - } - /** - * Assembles the join table and inserts it into the database. - * - * @param concept - */ - public void createConceptIdJoinTable(TreeConcept concept) { - CTConditionContext context = CTConditionContext.forJoinTables(functionProvider); + return unioned; + } - List conceptConditions = collectAllExpressions(concept, null, context); + private static Param defaultValue(Field field) { + if (field.getDataType().isBoolean()) { + return inline(false); + } - List> allFields = collectAllFields(conceptConditions); + if (field.getDataType().isString()) { + return inline(null, String.class); + } - List rows = expressionsToRows(conceptConditions, allFields); + throw new IllegalStateException("Fields of type %s are not expected".formatted(field.getDataType())); + } - Name tableName = idsTableName(concept.getName()); + /** + * Assembles the join table and inserts it into the database. + * + * @param concept + */ + public void createConceptIdJoinTable(TreeConcept concept) { + CTConditionContext context = CTConditionContext.forJoinTables(functionProvider); - // allFields are the statements to extract values from the underlying tables, we use them to generate the field names - List> fields = new ArrayList<>(); + List conceptConditions = collectAllExpressions(concept, null, context); - fields.addAll(allFields); - fields.addFirst(CONCEPT_ID_FIELD); + List> allFields = collectReferencedFields(conceptConditions); - // Make sure there's no table present. - deleteConceptIdJoinTable(concept.getId()); - createConceptIdsTable(tableName, fields); + List rows = expressionsToRows(conceptConditions, allFields); - insertConceptIdMappings(tableName, fields, rows, dslContext); - } + Name tableName = idsTableName(concept.getName()); - @NotNull - private Field[] collectValidityDateFields(Connector connector) { - List> validityDates = new ArrayList<>(); + // Make sure there's no table present. + deleteConceptIdJoinTable(concept.getId()); + List> fields = createConceptIdsTable(tableName, allFields); - for (ValidityDate validityDate : connector.getValidityDates()) { - if (validityDate.isSingleColumnDaterange()) { - Column column = validityDate.getColumn().get(); - validityDates.add(field(name(column.getName()), Date.class)); - } else { - validityDates.add(field(name(validityDate.getStartColumn().getColumn()), Date.class)); - validityDates.add(field(name(validityDate.getEndColumn().getColumn()), Date.class)); - } + insertConceptIdMappings(tableName, fields, rows, dslContext); + } - } - return (Field[]) validityDates.toArray(Field[]::new); - } + @NotNull + private Field[] collectValidityDateFields(Connector connector) { + List> validityDates = new ArrayList<>(); - private void assignStats(Map, MatchingStats.Entry> matchingStats) { - for (Map.Entry, MatchingStats.Entry> entry : matchingStats.entrySet()) { - ConceptElementId conceptElementId = entry.getKey(); + for (ValidityDate validityDate : connector.getValidityDates()) { + if (validityDate.isSingleColumnDaterange()) { + Column column = validityDate.getColumn().get(); + validityDates.add(field(name(column.getName()), Date.class)); + } else { + validityDates.add(field(name(validityDate.getStartColumn().getColumn()), Date.class)); + validityDates.add(field(name(validityDate.getEndColumn().getColumn()), Date.class)); + } - MatchingStats stats = new MatchingStats(); - stats.putEntry("sql", entry.getValue()); - conceptElementId.resolve().setMatchingStats(stats); - } - } + } + return (Field[]) validityDates.toArray(Field[]::new); + } - @NotNull - private Map, MatchingStats.Entry> readStats(TreeConcept concept, SelectJoinStep selectJoinStep) { - Map, MatchingStats.Entry> matchingStats = new HashMap<>(); + private void assignStats(Map, MatchingStats.Entry> matchingStats) { + for (Map.Entry, MatchingStats.Entry> entry : matchingStats.entrySet()) { + ConceptElementId conceptElementId = entry.getKey(); - Stopwatch stopwatch = Stopwatch.createStarted(); + MatchingStats stats = new MatchingStats(); + stats.putEntry(SQL_SOURCE_MATCHING_STATS_LABEL, entry.getValue()); + conceptElementId.resolve().setMatchingStats(stats); + } + } - log.info("BEGIN fetching matching stats for {}", concept.getId()); - log.trace("{}", selectJoinStep); + @NotNull + private Map, MatchingStats.Entry> readStats(TreeConcept concept, SelectJoinStep selectJoinStep) { + Map, MatchingStats.Entry> matchingStats = new HashMap<>(); - try (Cursor cursor = selectJoinStep.fetchSize(fetchBatchSize).fetchLazy()) { + Stopwatch stopwatch = Stopwatch.createStarted(); - for (Record record : cursor) { + log.info("BEGIN fetching matching stats for {}", concept.getId()); + log.trace("{}", selectJoinStep); - Integer rawId = record.get(CONCEPT_ID_FIELD); - ConceptElement resolvedId; - if (rawId == null) { - resolvedId = concept; - } else { - resolvedId = concept.getElementByLocalId(rawId); - } + try (Cursor cursor = selectJoinStep.fetchSize(fetchBatchSize).fetchLazy()) { - String entity = record.get(PID_FIELD); - Date min = record.get(LB_FIELD); - Date max = record.get(UB_FIELD); + for (Record record : cursor) { - CDateRange span = CDateRange.of(min != null ? min.toLocalDate() : null, max != null ? max.toLocalDate() : null); + Integer rawId = record.get(CONCEPT_ID_FIELD); + ConceptElement resolvedId; + if (rawId == null) { + resolvedId = concept; + } else { + resolvedId = concept.getElementByLocalId(rawId); + } - assignStatsToPath(resolvedId, matchingStats, entity, span); - } - } + String entity = record.get(PID_FIELD); + Date min = record.get(LB_FIELD); + Date max = record.get(UB_FIELD); - log.debug("DONE fetching matching stats for {} within {}", concept.getId(), stopwatch); + CDateRange span = CDateRange.of(min != null ? min.toLocalDate() : null, max != null ? max.toLocalDate() : null); + assignStatsToPath(resolvedId, matchingStats, entity, span); + } + } - return matchingStats; - } + log.debug("DONE fetching matching stats for {} within {}", concept.getId(), stopwatch); - @NotNull - private Name idsTableName(@NotBlank String name) { - return name("%s_ids".formatted(name)); - } + return matchingStats; + } - private void insertConceptIdMappings(Name tableName, List> fieldNames, List rows, DSLContext dsl) { - log.info("BEGIN inserting {} rows into {}", rows.size(), tableName); + @NotNull + private Name idsTableName(@NotBlank String name) { + return name("%s_ids".formatted(name)); + } - // We're using batching here because some DBMS don't allow mass inserts. - // There's a chance, we rework this to use a prepared statement with lots of bindings under the hood. But that needs to rework the entire stream of rows. - List> inserts = new ArrayList<>(rows.size()); + private void insertConceptIdMappings(Name tableName, List> fieldNames, List rows, DSLContext dsl) { + log.info("BEGIN inserting {} rows into {}", rows.size(), tableName); + Stopwatch stopwatch = Stopwatch.createStarted(); - for (RowN row : rows) { - inserts.add(dsl.insertInto(table(tableName)).columns(fieldNames).values(row)); - } + // We're using batching here because some DBMS don't allow mass inserts. + // There's a chance, we rework this to use a prepared statement with lots of bindings under the hood. But that needs to rework the entire stream of rows. + List> inserts = new ArrayList<>(rows.size()); - dsl.batch(inserts).execute(); + for (RowN row : rows) { + inserts.add(dsl.insertInto(table(tableName)).columns(fieldNames).values(row)); + } + dsl.batch(inserts).execute(); - log.trace("DONE inserting into {}", tableName); - } + log.debug("DONE inserting into {} within {}", tableName, stopwatch); + } - /** - * Create table and fields. Assumes, table has been dropped already. - */ - private void createConceptIdsTable(Name tableName, List> fields) { + /** + * Create table and fields. Assumes, table has been dropped already. + */ + private List> createConceptIdsTable(Name tableName, List> keyFields) { - log.debug("Creating table {} with fields {}", tableName, fields); + List> fields = new ArrayList<>(); - CreateTableElementListStep createTable = dslContext.createTable(tableName).columns(fields); + fields.addAll(keyFields); + fields.addFirst(CONCEPT_ID_FIELD); - log.info("{}", createTable); + log.debug("Creating table {} with fields {}", tableName, fields); - createTable.execute(); - } + //TODO Option to create primaryKeys and indices here, but Hana is a bit flaky with it, would need to differentiate the impls. - public ListenableFuture collectMatchingStatsForConcept(TreeConcept concept, ListeningExecutorService executorService) { - return executorService.submit(() -> { - dslContext.connection(cfg -> { - SelectJoinStep matchingStatsStatement = createMatchingStatsStatement(concept); - Map, MatchingStats.Entry> matchingStats = readStats(concept, matchingStatsStatement); - assignStats(matchingStats); - }); - }); + CreateTableElementListStep createTable = + dslContext.createTable(tableName) + .columns(fields); - } + createTable.execute(); - @NotNull - private SelectJoinStep createMatchingStatsStatement(TreeConcept concept) { + return fields; + } - List> connectorTables = new ArrayList<>(); + public ListenableFuture collectMatchingStatsForConcept(TreeConcept concept, ListeningExecutorService executorService, int tries) { + return executorService.submit(() -> { + dslContext.connection(cfg -> { + try { + SelectJoinStep matchingStatsStatement = createMatchingStatsStatement(concept); + Map, MatchingStats.Entry> matchingStats = readStats(concept, matchingStatsStatement); + assignStats(matchingStats); - Field positiveInfinity = functionProvider.getMaxDateExpression(); - Field negativeInfinity = functionProvider.getMinDateExpression(); + } catch (DataAccessException e) { + log.debug("Failed to connect to database for concept {}. Retrying.", concept.getId(), (Exception) (log.isTraceEnabled() || tries == 0 ? e : null)); - for (Connector connector : concept.getConnectors()) { + if (tries > 0) { + collectMatchingStatsForConcept(concept, executorService, tries - 1); + } + } + }); + }); + } - CTConditionContext context = CTConditionContext.forConnector(connector, functionProvider); + @NotNull + private SelectJoinStep createMatchingStatsStatement(TreeConcept concept) { - Field[] validityDates = collectValidityDateFields(connector); + List> connectorTables = new ArrayList<>(); - SelectConditionStep connectorTable = - dslContext.select( - TablePrimaryColumnUtil.findPrimaryColumn(connector.getResolvedTable(), defaultPrimaryColumn).as(PID_FIELD), - // The infinities are intentionally swapped - least(positiveInfinity, validityDates).as(LB_FIELD), - greatest(negativeInfinity, validityDates).as(UB_FIELD), - CONCEPT_ID_FIELD) - .from(table(name(connector.getResolvedTable().getName()))) - .leftJoin(idsTableName(concept.getName())) - // join onto the concept-ids table to assign the most specific id. - .on(getJoinConditions(concept, context)) - .where(connector.getCondition() != null ? connector.getCondition().convertToSqlCondition(context).condition() : noCondition()); + Field positiveInfinity = functionProvider.getMaxDateExpression(); + Field negativeInfinity = functionProvider.getMinDateExpression(); - connectorTables.add(connectorTable); - } + for (Connector connector : concept.getConnectors()) { - Name ct_name = name("connector_tables"); - CommonTableExpression unioned = ct_name.as(unionSelects(connectorTables)); + CTConditionContext context = CTConditionContext.forConnector(connector, functionProvider); - SelectJoinStep> records = dslContext.with(unioned).select(unioned.field(CONCEPT_ID_FIELD), PID_FIELD, - // The infinities are intentionally swapped - nullif(unioned.field(LB_FIELD), positiveInfinity).as(LB_FIELD), nullif(unioned.field(UB_FIELD), negativeInfinity).as(UB_FIELD)).from(ct_name); + Field[] validityDates = collectValidityDateFields(connector); - return records; - } + SelectConditionStep connectorTable = + dslContext.select( + TablePrimaryColumnUtil.findPrimaryColumn(connector.getResolvedTable(), defaultPrimaryColumn).as(PID_FIELD), + // The infinities are intentionally swapped + least(positiveInfinity, validityDates).as(LB_FIELD), + greatest(negativeInfinity, validityDates).as(UB_FIELD), + CONCEPT_ID_FIELD) + .from(table(name(connector.getResolvedTable().getName()))) + .leftJoin(idsTableName(concept.getName())) + // join onto the concept-ids table to assign the most specific id. + .on(getJoinConditions(concept, context)) + .where(connector.getCondition() != null ? connector.getCondition().convertToSqlCondition(context).condition() : noCondition()); - public void deleteConceptIdJoinTable(ConceptId concept) { - Name tableName = idsTableName(concept.getName()); - log.debug("Trying to delete id-table {}", tableName); - try { - dslContext.dropTable(tableName).execute(); - } catch (DataAccessException exception) { - // Likely it doesn't exist. Some DBMS just don't support drop-IfExists so this is the next best thing :^) - log.trace("Failed to drop table {}", tableName, exception); - } - } + connectorTables.add(connectorTable); + } - /** - * Using the expressions of a concept, build a Condition that descibes the left-join onto the ids table, from any connector-table. - */ - private Condition getJoinConditions(TreeConcept concept, CTConditionContext context) { - List conceptConditions = collectAllExpressions(concept, null, context); + Name ct_name = name("connector_tables"); + CommonTableExpression unioned = ct_name.as(unionSelects(connectorTables)); + SelectJoinStep> records = dslContext.with(unioned).select(unioned.field(CONCEPT_ID_FIELD), PID_FIELD, + // The infinities are intentionally swapped + nullif(unioned.field(LB_FIELD), positiveInfinity).as(LB_FIELD), nullif(unioned.field(UB_FIELD), negativeInfinity).as(UB_FIELD)).from(ct_name); + + return records; + } + + public void deleteConceptIdJoinTable(ConceptId concept) { + Name tableName = idsTableName(concept.getName()); + log.debug("Trying to delete id-table {}", tableName); - if (conceptConditions.isEmpty()) { - return context.getFunctionProvider().unconditionalJoinCondition(); - } - - Set conditions = new HashSet<>(); + try { + dslContext.dropTable(tableName).execute(); + } catch (DataAccessException exception) { + // Likely it doesn't exist. Some DBMS just don't support drop-IfExists so this is the next best thing :^) + log.trace("Failed to drop table {}", tableName, exception); + } + } - for (CTCondition.ConceptConditions conceptCondition : conceptConditions) { - for (Map.Entry, CTCondition.FieldCondition> entry : conceptCondition.conditions().entrySet()) { - conditions.add(entry.getKey().eq((Field) entry.getValue().extractor())); - } - } - - Condition reduced = conditions.stream().reduce(noCondition(), Condition::and); - - if (reduced.equals(noCondition())) { - return context.getFunctionProvider().unconditionalJoinCondition(); - } - - return reduced; - } - - private List expressionsToRows(List conceptConditions, List> allFields) { - Map>, ConceptElement> byDepth = new HashMap<>(); - - for (CTCondition.ConceptConditions conceptCondition : conceptConditions) { - ConceptElement elt = conceptCondition.conceptElement(); - Map, CTCondition.FieldCondition> conditions = conceptCondition.conditions(); - - List>> rowValues = new ArrayList<>(); - for (Field field : allFields) { - if (conditions.containsKey(field)) { - rowValues.add(conditions.get(field).params()); - } else { - rowValues.add(Set.of(defaultValue(field))); - } - } - - Set>> flattened = Sets.cartesianProduct(rowValues); - - // Group by params, find deepest params. This ensures we map to the most-specific element. - for (List> params : flattened) { - byDepth.compute(params, (__, prior) -> { - if (prior == null || prior.getDepth() < elt.getDepth()) { - return elt; - } - if (prior.getDepth() == elt.getDepth() && !prior.equals(elt)) { - log.warn("Nodes {} and {} are mapped by the same params {}", prior.getId(), elt.getId(), params); - } - return prior; - }); - } - } - - List rows = new ArrayList<>(); - - for (Map.Entry>, ConceptElement> entry : byDepth.entrySet()) { - List> params = new ArrayList<>(entry.getKey().size() + 1); - - params.addFirst(val(entry.getValue().getLocalId())); - params.addAll(entry.getKey()); - - rows.add(row(params)); - } - return rows; - } - - /** - * Collect all mappings from values to conceptElement for the entire concept. This means the column-value and the auxiliary columns. - * We use them to construct a table building an injective mapping from values to concept element that can be used for performant joins instead of resolving the concept every time. - */ - private List collectAllExpressions(ConceptElement current, CTCondition.ConceptConditions parentConceptCondition, CTConditionContext context) { - - final CTCondition.ConceptConditions forCurrent = switch (current) { - case TreeConcept concept -> new CTCondition.ConceptConditions(concept, Collections.emptyMap()); - // concept elements implicitly inherit the conditions of its parents - case ConceptTreeChild child -> - child.getCondition().buildExpression(context, current).and(parentConceptCondition); - case null, default -> throw new IllegalStateException(); - }; - - final List out = new ArrayList<>(); - - out.add(forCurrent); - - for (ConceptTreeChild child : current.getChildren()) { - out.addAll(collectAllExpressions(child, forCurrent, context)); - } - - return out; - } - - /** - * recursively build just a single expression - * - * @param current - * @param context TODO use this to implement joining in queries - */ - private CTCondition.ConceptConditions collectExpressionsForSingleNode(ConceptElement current, CTConditionContext context) { - - if (current instanceof TreeConcept concept) { - return new CTCondition.ConceptConditions(concept, Collections.emptyMap()); - } - - CTCondition.ConceptConditions parentConceptCondition = collectExpressionsForSingleNode(current.getParent(), context); - CTCondition.ConceptConditions currentConceptCondition = ((ConceptTreeChild) current).getCondition().buildExpression(context, current); - - return currentConceptCondition.and(parentConceptCondition); - } + /** + * Using the expressions of a concept, build a Condition that descibes the left-join onto the ids table, from any connector-table. + */ + private Condition getJoinConditions(TreeConcept concept, CTConditionContext context) { + List conceptConditions = collectAllExpressions(concept, null, context); + + + if (conceptConditions.isEmpty()) { + return context.getFunctionProvider().unconditionalJoinCondition(); + } + + Set conditions = new HashSet<>(); + + for (CTCondition.ConceptConditions conceptCondition : conceptConditions) { + for (Map.Entry, CTCondition.FieldCondition> entry : conceptCondition.conditions().entrySet()) { + conditions.add(entry.getKey().eq((Field) entry.getValue().extractor())); + } + } + + Condition reduced = conditions.stream().reduce(noCondition(), Condition::and); + + if (reduced.equals(noCondition())) { + return context.getFunctionProvider().unconditionalJoinCondition(); + } + + return reduced; + } + + private List expressionsToRows(List conceptConditions, List> allFields) { + Map>, ConceptElement> byDepth = new HashMap<>(); + + for (CTCondition.ConceptConditions conceptCondition : conceptConditions) { + ConceptElement elt = conceptCondition.conceptElement(); + Map, CTCondition.FieldCondition> conditions = conceptCondition.conditions(); + + List>> rowValues = new ArrayList<>(); + for (Field field : allFields) { + if (conditions.containsKey(field)) { + rowValues.add(conditions.get(field).params()); + } else { + rowValues.add(Set.of(defaultValue(field))); + } + } + + Set>> flattened = Sets.cartesianProduct(rowValues); + + // Group by params, find deepest params. This ensures we map to the most-specific element. + for (List> params : flattened) { + byDepth.compute(params, (__, prior) -> { + if (prior == null || prior.getDepth() < elt.getDepth()) { + return elt; + } + if (prior.getDepth() == elt.getDepth() && !prior.equals(elt)) { + log.warn("Nodes {} and {} are mapped by the same params {}", prior.getId(), elt.getId(), params); + } + return prior; + }); + } + } + + List rows = new ArrayList<>(); + + for (Map.Entry>, ConceptElement> entry : byDepth.entrySet()) { + List> params = new ArrayList<>(entry.getKey().size() + 1); + + params.addFirst(val(entry.getValue().getLocalId())); + params.addAll(entry.getKey()); + + rows.add(row(params)); + } + return rows; + } + + /** + * Collect all mappings from values to conceptElement for the entire concept. This means the column-value and the auxiliary columns. + * We use them to construct a table building an injective mapping from values to concept element that can be used for performant joins instead of resolving the concept every time. + */ + private List collectAllExpressions(ConceptElement current, CTCondition.ConceptConditions parentConceptCondition, CTConditionContext context) { + + final CTCondition.ConceptConditions forCurrent = switch (current) { + case TreeConcept concept -> new CTCondition.ConceptConditions(concept, Collections.emptyMap()); + // concept elements implicitly inherit the conditions of its parents + case ConceptTreeChild child -> + child.getCondition().buildExpression(context, current).and(parentConceptCondition); + case null, default -> throw new IllegalStateException(); + }; + + final List out = new ArrayList<>(); + + out.add(forCurrent); + + for (ConceptTreeChild child : current.getChildren()) { + out.addAll(collectAllExpressions(child, forCurrent, context)); + } + + return out; + } + + /** + * recursively build just a single expression + * + * @param current + * @param context TODO use this to implement joining in queries + */ + private CTCondition.ConceptConditions collectExpressionsForSingleNode(ConceptElement current, CTConditionContext context) { + + if (current instanceof TreeConcept concept) { + return new CTCondition.ConceptConditions(concept, Collections.emptyMap()); + } + + CTCondition.ConceptConditions parentConceptCondition = collectExpressionsForSingleNode(current.getParent(), context); + CTCondition.ConceptConditions currentConceptCondition = ((ConceptTreeChild) current).getCondition().buildExpression(context, current); + + return currentConceptCondition.and(parentConceptCondition); + } } diff --git a/backend/src/main/java/com/bakdata/conquery/sql/conversion/NodeConversions.java b/backend/src/main/java/com/bakdata/conquery/sql/conversion/NodeConversions.java index 38a6eaaf33..6b3a3fa08a 100644 --- a/backend/src/main/java/com/bakdata/conquery/sql/conversion/NodeConversions.java +++ b/backend/src/main/java/com/bakdata/conquery/sql/conversion/NodeConversions.java @@ -1,8 +1,5 @@ package com.bakdata.conquery.sql.conversion; -import java.time.Clock; -import java.util.Locale; - import com.bakdata.conquery.apiv1.query.QueryDescription; import com.bakdata.conquery.models.config.ConqueryConfig; import com.bakdata.conquery.models.config.IdColumnConfig; @@ -13,8 +10,12 @@ import com.bakdata.conquery.sql.conversion.dialect.DialectBundle; import com.bakdata.conquery.sql.conversion.model.NameGenerator; import com.bakdata.conquery.sql.execution.SqlExecutionService; +import lombok.NonNull; import org.jooq.DSLContext; +import java.time.Clock; +import java.util.Locale; + /** * Entry point for converting {@link QueryDescription} to an SQL query. */ @@ -25,12 +26,16 @@ public class NodeConversions extends Conversions findPrimaryColumn(Table table, String defaultPrimaryColumn) { @@ -15,6 +17,10 @@ public static Field findPrimaryColumn(Table table, String defaultPrimary primaryColumnName = table.getPrimaryColumn().getName(); } + if (primaryColumnName == null) { + throw new IllegalArgumentException("Unable to determine primary column for table " + table.getId()); + } + return field(name(table.getName(), primaryColumnName), String.class); } diff --git a/backend/src/test/java/com/bakdata/conquery/integration/sql/CsvTableImporter.java b/backend/src/test/java/com/bakdata/conquery/integration/sql/CsvTableImporter.java index 8838c5ecf6..a0543cdeb3 100644 --- a/backend/src/test/java/com/bakdata/conquery/integration/sql/CsvTableImporter.java +++ b/backend/src/test/java/com/bakdata/conquery/integration/sql/CsvTableImporter.java @@ -39,7 +39,6 @@ import org.jooq.impl.BuiltInDataType; import org.jooq.impl.DSL; import org.jooq.impl.SQLDataType; -import org.jooq.postgres.extensions.types.DateRange; @Slf4j public class CsvTableImporter { @@ -153,7 +152,7 @@ private Field createField(RequiredColumn requiredColumn) { // TODO (ja) how do we handle REAL and DECIMAL properly? case REAL, DECIMAL, MONEY -> SQLDataType.DECIMAL(10, 2); case DATE -> SQLDataType.DATE.nullable(true); - case DATE_RANGE -> new BuiltInDataType<>(DateRange.class, "daterange"); + default -> throw new IllegalArgumentException("Unsupported data type: " + requiredColumn.getType()); }; // Set all columns except 'pid' to nullable, important for ClickHouse compatibility @@ -236,12 +235,7 @@ private Object readAccordingToColumnType(com.univocity.parsers.common.record.Rec case REAL -> record.getDouble(column); case DECIMAL, MONEY -> record.getBigDecimal(column); case DATE -> dateReader.parseToLocalDate(record.getString(column)); - case DATE_RANGE -> { - CDateRange dateRange = dateReader.parseToCDateRange(record.getString(column)); - yield DateRange.dateRange(dateRange.getMin() != null ? Date.valueOf(dateRange.getMin()) : null, true, - dateRange.getMax() != null ? Date.valueOf(dateRange.getMax()) : null, true - ); - } + default -> throw new IllegalArgumentException("Unsupported data type: " + type); }; }