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