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
12 changes: 1 addition & 11 deletions backend/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -348,15 +348,10 @@
<artifactId>jooq</artifactId>
<version>3.20.2</version>
</dependency>
<dependency>
<groupId>org.jooq</groupId>
<artifactId>jooq-postgres-extensions</artifactId>
<version>3.20.2</version>
</dependency>
<dependency>
<groupId>com.zaxxer</groupId>
<artifactId>HikariCP</artifactId>
<version>5.1.0</version>
<version>7.1.0</version>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
Expand All @@ -368,11 +363,6 @@
<artifactId>testcontainers-junit-jupiter</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>testcontainers-postgresql</artifactId>
<scope>test</scope>
</dependency>
<dependency>
<groupId>org.testcontainers</groupId>
<artifactId>testcontainers-clickhouse</artifactId>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand All @@ -100,10 +97,6 @@ public void run(Manager manager) throws InterruptedException {

environment.lifecycle().manage(this);

loadNamespaces();

loadMetaStorage();

Comment thread
thoniTUB marked this conversation as resolved.
// Create AdminServlet first to make it available to the realms
admin = new AdminServlet(this);

Expand Down Expand Up @@ -159,7 +152,6 @@ public void loadNamespaces() {
});
}


loaders.shutdown();
while (!loaders.awaitTermination(1, TimeUnit.MINUTES)) {
final int countLoaded = registry.getNamespaces().size();
Expand Down Expand Up @@ -201,6 +193,9 @@ private void registerTasks(Manager manager, Environment environment, ConqueryCon

@Override
public void start() throws Exception {
loadNamespaces();
loadMetaStorage();

manager.start();
}

Expand All @@ -214,7 +209,6 @@ public void stop() throws Exception {
catch (Exception e) {
log.error("{} could not be closed", provider, e);
}

}

try {
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -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;
Expand All @@ -24,6 +21,8 @@
import lombok.extern.slf4j.Slf4j;
import org.jooq.DSLContext;

import java.time.Clock;

@RequiredArgsConstructor
@Slf4j
public class LocalNamespaceHandler implements NamespaceHandler<LocalNamespace> {
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Die 1000 ist ein Timeout oder?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

genau, der geht hier aber nur auf die connection. Denke, wenn die Verbindung nach 1sek nicht als valide aufgebaut werden kann ist da was kaputt

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);
}
}

Expand Down
Original file line number Diff line number Diff line change
@@ -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
Comment thread
thoniTUB marked this conversation as resolved.
@EqualsAndHashCode(callSuper = false)
public class UpdateMatchingStatsSqlJob extends Job {

@ToString.Exclude
private final List<Concept<?>> concepts;
private final Dataset dataset;
@ToString.Exclude
private final List<Concept<?>> 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<ListenableFuture<?>> jobs = new ArrayList<>();
Map<ConceptId, ListenableFuture<?>> jobsByConcept = new HashedMap<>();
Collection<ListenableFuture<?>> 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<List<@Nullable Object>> 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);
Comment thread
awildturtok marked this conversation as resolved.
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) {
Comment thread
awildturtok marked this conversation as resolved.
// 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());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Comment thread
thoniTUB marked this conversation as resolved.

/**
* 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.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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
Expand Down
Loading
Loading