From be1aa2bb6189d5f61719c877692fc04143672b11 Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Mon, 29 Jun 2026 16:57:52 +0200 Subject: [PATCH 1/7] Adds logging for SqlTableValidator.java for presumably missing tables --- .../com/bakdata/conquery/util/validation/SqlTableValidator.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/backend/src/main/java/com/bakdata/conquery/util/validation/SqlTableValidator.java b/backend/src/main/java/com/bakdata/conquery/util/validation/SqlTableValidator.java index 6c4a4c9e67..6616702229 100644 --- a/backend/src/main/java/com/bakdata/conquery/util/validation/SqlTableValidator.java +++ b/backend/src/main/java/com/bakdata/conquery/util/validation/SqlTableValidator.java @@ -41,6 +41,8 @@ public boolean isValid(Table value, ConstraintValidatorContext context) { .fetch(); } catch (DataAccessException e) { + log.trace("Failed to test for table {}", value.getName(), e); + context.buildConstraintViolationWithTemplate("SQL table %s does not exist".formatted(value.getName())) .addPropertyNode("name") .addConstraintViolation(); From d679125f18a26db027d8d15a9b22542bb739313e Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Thu, 9 Jul 2026 11:29:43 +0200 Subject: [PATCH 2/7] Bugfixes from Matching Stats testing --- backend/pom.xml | 12 +- .../mode/local/LocalNamespaceHandler.java | 3 +- .../mode/local/UpdateMatchingStatsSqlJob.java | 99 +-- .../sql/conquery/SqlMatchingStats.java | 655 +++++++++--------- .../sql/conversion/NodeConversions.java | 31 +- .../conquery/util/TablePrimaryColumnUtil.java | 6 + .../integration/sql/CsvTableImporter.java | 10 +- 7 files changed, 403 insertions(+), 413 deletions(-) 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/mode/local/LocalNamespaceHandler.java b/backend/src/main/java/com/bakdata/conquery/mode/local/LocalNamespaceHandler.java index fced1ca2da..8a3b55c3a2 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 @@ -13,7 +13,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; @@ -49,7 +48,7 @@ public LocalNamespace createNamespace( ResultSetProcessor resultSetProcessor = dialectBundle.getResultSetProcessor(config); SqlExecutionService sqlExecutionService = new SqlExecutionService(dslContext, resultSetProcessor); - NodeConversions nodeConversions = new NodeConversions(config.getIdColumns(), dialectBundle, dslContext, sqlExecutionService, clock); + 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); 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..73c362ecf6 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,89 @@ package com.bakdata.conquery.mode.local; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.*; - 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.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; + +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; @Slf4j @Data 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; + + + @Override + public void execute() throws Exception { - @ToString.Exclude - private final SqlMatchingStats matchingStats; + log.info("BEGIN collecting SQL matching stats for {}", dataset); + Stopwatch stopwatch = Stopwatch.createStarted(); - @Override - public void execute() throws Exception { - log.info("BEGIN collecting SQL matching stats for {}", dataset); - Stopwatch stopwatch = Stopwatch.createStarted(); + ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(5)); - ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newVirtualThreadPerTaskExecutor()); + Map> jobsByConcept = new HashedMap<>(); + Collection> jobs = jobsByConcept.values(); - List> jobs = new ArrayList<>(); + for (Concept concept : concepts) { + if (!(concept instanceof TreeConcept)) { + continue; + } + jobsByConcept.put(concept.getId(), matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, 5)); + } - for (Concept concept : concepts) { - if (!(concept instanceof TreeConcept)) { - continue; - } - jobs.add(matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService)); - } + log.info("WAITING for {} jobs in {}", jobs.size(), dataset); - ListenableFuture> all = Futures.allAsList(jobs); + while (jobs.stream().anyMatch(job -> job.state().equals(Future.State.RUNNING))) { + for (ListenableFuture someJob : jobs) { + if (someJob.isDone()) { + continue; + } - while (!all.isDone()) { - if (isCancelled()) { - all.cancel(true); - log.debug("CANCELLED update matching stats for {}", getDataset(), all.exceptionNow()); - return; - } + try { + someJob.get(5, TimeUnit.SECONDS); + } catch (Exception e) { + // intentionally left blank + } - all.get(5, TimeUnit.SECONDS); - log.trace("WAITING for matching stats to finish {}", getDataset()); + log.debug("Still waiting for {} matching stats to finish.", jobs.stream().filter(job -> job.state().equals(Future.State.RUNNING)).count()); + } + } - if (all.state().equals(Future.State.FAILED)) { - log.error("FAILED update matching stats for {}", getDataset(), all.exceptionNow()); - return; - } - } + for (Map.Entry> conceptState : jobsByConcept.entrySet()) { + if (conceptState.getValue().state().equals(Future.State.FAILED)) { + log.warn("Failed to collect SQL matching stats for {}", conceptState.getKey(), conceptState.getValue().exceptionNow()); + } + } - log.debug("DONE collecting SQL matching stats for {} within {}", dataset, stopwatch); - } + 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()); - } + @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/sql/conquery/SqlMatchingStats.java b/backend/src/main/java/com/bakdata/conquery/sql/conquery/SqlMatchingStats.java index a98646b38f..f6aa01db44 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,396 @@ 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); - } - + 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)); - return unioned; - } + 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) { + 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> collectAllFields(List conceptConditions) { + List> fields = conceptConditions.stream().flatMap(e -> e.conditions().keySet().stream()).distinct().toList(); + return fields; + } - if (field.getDataType().isString()) { - return inline(null, String.class); - } + private static Select unionSelects(List> connectorTableSelects) { + Select unioned = null; - throw new IllegalStateException("Fields of type %s are not expected".formatted(field.getDataType())); - } + for (Select connectorTable : connectorTableSelects) { + if (unioned == null) { + unioned = (Select) connectorTable; + continue; + } - /** - * Assembles the join table and inserts it into the database. - * - * @param concept - */ - public void createConceptIdJoinTable(TreeConcept concept) { - CTConditionContext context = CTConditionContext.forJoinTables(functionProvider); + unioned = unioned.unionAll(connectorTable); + } - List conceptConditions = collectAllExpressions(concept, null, context); - List> allFields = collectAllFields(conceptConditions); + return unioned; + } - List rows = expressionsToRows(conceptConditions, allFields); + private static Param defaultValue(Field field) { + if (field.getDataType().isBoolean()) { + return inline(false); + } - Name tableName = idsTableName(concept.getName()); + if (field.getDataType().isString()) { + return inline(null, String.class); + } - // allFields are the statements to extract values from the underlying tables, we use them to generate the field names - List> fields = new ArrayList<>(); + throw new IllegalStateException("Fields of type %s are not expected".formatted(field.getDataType())); + } - fields.addAll(allFields); - fields.addFirst(CONCEPT_ID_FIELD); + /** + * Assembles the join table and inserts it into the database. + * + * @param concept + */ + public void createConceptIdJoinTable(TreeConcept concept) { + CTConditionContext context = CTConditionContext.forJoinTables(functionProvider); - // Make sure there's no table present. - deleteConceptIdJoinTable(concept.getId()); - createConceptIdsTable(tableName, fields); + List conceptConditions = collectAllExpressions(concept, null, context); - insertConceptIdMappings(tableName, fields, rows, dslContext); - } + List> allFields = collectAllFields(conceptConditions); - @NotNull - private Field[] collectValidityDateFields(Connector connector) { - List> validityDates = new ArrayList<>(); + List rows = expressionsToRows(conceptConditions, 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)); - } + Name tableName = idsTableName(concept.getName()); - } - return (Field[]) validityDates.toArray(Field[]::new); - } + // allFields are the statements to extract values from the underlying tables, we use them to generate the field names + List> fields = new ArrayList<>(); - private void assignStats(Map, MatchingStats.Entry> matchingStats) { - for (Map.Entry, MatchingStats.Entry> entry : matchingStats.entrySet()) { - ConceptElementId conceptElementId = entry.getKey(); + fields.addAll(allFields); + fields.addFirst(CONCEPT_ID_FIELD); - MatchingStats stats = new MatchingStats(); - stats.putEntry("sql", entry.getValue()); - conceptElementId.resolve().setMatchingStats(stats); - } - } + // Make sure there's no table present. + deleteConceptIdJoinTable(concept.getId()); + createConceptIdsTable(tableName, fields); - @NotNull - private Map, MatchingStats.Entry> readStats(TreeConcept concept, SelectJoinStep selectJoinStep) { - Map, MatchingStats.Entry> matchingStats = new HashMap<>(); + insertConceptIdMappings(tableName, fields, rows, dslContext); + } - Stopwatch stopwatch = Stopwatch.createStarted(); + @NotNull + private Field[] collectValidityDateFields(Connector connector) { + List> validityDates = new ArrayList<>(); - log.info("BEGIN fetching matching stats for {}", concept.getId()); - log.trace("{}", selectJoinStep); + 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)); + } - try (Cursor cursor = selectJoinStep.fetchSize(fetchBatchSize).fetchLazy()) { + } + return (Field[]) validityDates.toArray(Field[]::new); + } - for (Record record : cursor) { + private void assignStats(Map, MatchingStats.Entry> matchingStats) { + for (Map.Entry, MatchingStats.Entry> entry : matchingStats.entrySet()) { + ConceptElementId conceptElementId = entry.getKey(); - Integer rawId = record.get(CONCEPT_ID_FIELD); - ConceptElement resolvedId; - if (rawId == null) { - resolvedId = concept; - } else { - resolvedId = concept.getElementByLocalId(rawId); - } + MatchingStats stats = new MatchingStats(); + stats.putEntry("sql", entry.getValue()); + conceptElementId.resolve().setMatchingStats(stats); + } + } - String entity = record.get(PID_FIELD); - Date min = record.get(LB_FIELD); - Date max = record.get(UB_FIELD); + @NotNull + private Map, MatchingStats.Entry> readStats(TreeConcept concept, SelectJoinStep selectJoinStep) { + Map, MatchingStats.Entry> matchingStats = new HashMap<>(); - CDateRange span = CDateRange.of(min != null ? min.toLocalDate() : null, max != null ? max.toLocalDate() : null); + Stopwatch stopwatch = Stopwatch.createStarted(); - assignStatsToPath(resolvedId, matchingStats, entity, span); - } - } + log.info("BEGIN fetching matching stats for {}", concept.getId()); + log.trace("{}", selectJoinStep); - log.debug("DONE fetching matching stats for {} within {}", concept.getId(), stopwatch); + try (Cursor cursor = selectJoinStep.fetchSize(fetchBatchSize).fetchLazy()) { + for (Record record : cursor) { - return matchingStats; - } + Integer rawId = record.get(CONCEPT_ID_FIELD); + ConceptElement resolvedId; + if (rawId == null) { + resolvedId = concept; + } else { + resolvedId = concept.getElementByLocalId(rawId); + } - @NotNull - private Name idsTableName(@NotBlank String name) { - return name("%s_ids".formatted(name)); - } + String entity = record.get(PID_FIELD); + Date min = record.get(LB_FIELD); + Date max = record.get(UB_FIELD); - private void insertConceptIdMappings(Name tableName, List> fieldNames, List rows, DSLContext dsl) { - log.info("BEGIN inserting {} rows into {}", rows.size(), tableName); + CDateRange span = CDateRange.of(min != null ? min.toLocalDate() : null, max != null ? max.toLocalDate() : null); - // 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()); + assignStatsToPath(resolvedId, matchingStats, entity, span); + } + } - for (RowN row : rows) { - inserts.add(dsl.insertInto(table(tableName)).columns(fieldNames).values(row)); - } + log.debug("DONE fetching matching stats for {} within {}", concept.getId(), stopwatch); - dsl.batch(inserts).execute(); + return matchingStats; + } - log.trace("DONE inserting into {}", tableName); - } + @NotNull + private Name idsTableName(@NotBlank String name) { + return name("%s_ids".formatted(name)); + } - /** - * Create table and fields. Assumes, table has been dropped already. - */ - private void createConceptIdsTable(Name tableName, List> fields) { + private void insertConceptIdMappings(Name tableName, List> fieldNames, List rows, DSLContext dsl) { + log.info("BEGIN inserting {} rows into {}", rows.size(), tableName); + Stopwatch stopwatch = Stopwatch.createStarted(); - log.debug("Creating table {} with fields {}", tableName, fields); + // 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()); - CreateTableElementListStep createTable = dslContext.createTable(tableName).columns(fields); + for (RowN row : rows) { + inserts.add(dsl.insertInto(table(tableName)).columns(fieldNames).values(row)); + } - log.info("{}", createTable); + dsl.batch(inserts).execute(); - createTable.execute(); - } + log.debug("DONE inserting into {} within {}", tableName, stopwatch); + } - 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); - }); - }); + /** + * Create table and fields. Assumes, table has been dropped already. + */ + private void createConceptIdsTable(Name tableName, List> fields) { - } + log.debug("Creating table {} with fields {}", tableName, fields); - @NotNull - private SelectJoinStep createMatchingStatsStatement(TreeConcept concept) { + CreateTableElementListStep createTable = dslContext.createTable(tableName).columns(fields); - List> connectorTables = new ArrayList<>(); + log.trace("{}", createTable); - Field positiveInfinity = functionProvider.getMaxDateExpression(); - Field negativeInfinity = functionProvider.getMinDateExpression(); + createTable.execute(); + } - for (Connector connector : concept.getConnectors()) { + 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); - CTConditionContext context = CTConditionContext.forConnector(connector, functionProvider); + } catch (DataAccessException e) { + log.debug("Failed to connect to database for concept {}. Retrying.", concept.getId(), (Exception) (log.isTraceEnabled() || tries == 0 ? e : null)); - Field[] validityDates = collectValidityDateFields(connector); + if (tries > 0) { + collectMatchingStatsForConcept(concept, executorService, tries - 1); + } + } + }); + }); + } - 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()); + @NotNull + private SelectJoinStep createMatchingStatsStatement(TreeConcept concept) { - connectorTables.add(connectorTable); - } + List> connectorTables = new ArrayList<>(); - Name ct_name = name("connector_tables"); - CommonTableExpression unioned = ct_name.as(unionSelects(connectorTables)); + Field positiveInfinity = functionProvider.getMaxDateExpression(); + Field negativeInfinity = functionProvider.getMinDateExpression(); - 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); + for (Connector connector : concept.getConnectors()) { - return records; - } + CTConditionContext context = CTConditionContext.forConnector(connector, functionProvider); - 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); - } - } + Field[] validityDates = collectValidityDateFields(connector); - /** - * 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); + 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()); + connectorTables.add(connectorTable); + } - if (conceptConditions.isEmpty()) { - return context.getFunctionProvider().unconditionalJoinCondition(); - } - - Set conditions = new HashSet<>(); + Name ct_name = name("connector_tables"); + CommonTableExpression unioned = ct_name.as(unionSelects(connectorTables)); - 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); - } + 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); + 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); + } + } + + /** + * 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..a3fbc60985 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,14 @@ 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); }; } From b36ca6ae9930ba2add19bc6a538c63306132af61 Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Thu, 16 Jul 2026 10:15:25 +0200 Subject: [PATCH 3/7] fixes initializtaion order --- .../conquery/commands/ManagerNode.java | 14 ++---- .../mode/local/LocalNamespaceHandler.java | 49 ++++++++++--------- .../mode/local/ManagedConnection.java | 2 +- 3 files changed, 32 insertions(+), 33 deletions(-) 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 8a3b55c3a2..b57d045e2d 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; @@ -23,6 +21,8 @@ import lombok.extern.slf4j.Slf4j; import org.jooq.DSLContext; +import java.time.Clock; + @RequiredArgsConstructor @Slf4j public class LocalNamespaceHandler implements NamespaceHandler { @@ -42,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, 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); + 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.getDataset(), 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..4b5e6fefd5 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,7 +31,6 @@ 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)) { @@ -43,6 +42,7 @@ public void start() throws Exception { } catch (SQLException exception) { log.error("FAILED connecting to {}", connection.getJdbcConnectionUrl(), exception); + throw exception; } } From e75125ed6bfde039576f8e0d6d0785dbfc29e177 Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Thu, 23 Jul 2026 12:13:43 +0200 Subject: [PATCH 4/7] minor adjustment for fields --- .../sql/conquery/SqlMatchingStats.java | 24 +++++++++++-------- 1 file changed, 14 insertions(+), 10 deletions(-) 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 f6aa01db44..5e6af62cc5 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 @@ -25,6 +25,7 @@ import org.jooq.*; import org.jooq.Record; import org.jooq.exception.DataAccessException; +import org.jooq.impl.DSL; import java.sql.Date; import java.util.*; @@ -108,15 +109,9 @@ public void createConceptIdJoinTable(TreeConcept concept) { Name tableName = idsTableName(concept.getName()); - // allFields are the statements to extract values from the underlying tables, we use them to generate the field names - List> fields = new ArrayList<>(); - - fields.addAll(allFields); - fields.addFirst(CONCEPT_ID_FIELD); - // Make sure there's no table present. deleteConceptIdJoinTable(concept.getId()); - createConceptIdsTable(tableName, fields); + List> fields = createConceptIdsTable(tableName, allFields); insertConceptIdMappings(tableName, fields, rows, dslContext); } @@ -210,15 +205,24 @@ private void insertConceptIdMappings(Name tableName, List> fieldNames, /** * Create table and fields. Assumes, table has been dropped already. */ - private void createConceptIdsTable(Name tableName, List> fields) { + private List> createConceptIdsTable(Name tableName, List> keyFields) { + + List> fields = new ArrayList<>(); + + fields.addAll(keyFields); + fields.addFirst(CONCEPT_ID_FIELD); log.debug("Creating table {} with fields {}", tableName, fields); - CreateTableElementListStep createTable = dslContext.createTable(tableName).columns(fields); + //TODO Option to create primaryKeys and indices here, but Hana is a bit flaky with it, would need to differentiate the impls. - log.trace("{}", createTable); + CreateTableElementListStep createTable = + dslContext.createTable(tableName) + .columns(fields); createTable.execute(); + + return fields; } public ListenableFuture collectMatchingStatsForConcept(TreeConcept concept, ListeningExecutorService executorService, int tries) { From 3db299f516d1d1f94dde1a567d4f23157cef6c62 Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Thu, 23 Jul 2026 12:32:16 +0200 Subject: [PATCH 5/7] Wire through worker size --- .../mode/local/UpdateMatchingStatsSqlJob.java | 12 ++++-------- .../models/config/DatabaseConnectionConfig.java | 8 ++++++++ .../conquery/models/worker/LocalNamespace.java | 14 +++++++------- .../conquery/sql/conquery/SqlMatchingStats.java | 3 ++- 4 files changed, 21 insertions(+), 16 deletions(-) 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 73c362ecf6..b639b3265b 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 @@ -41,9 +41,7 @@ public void execute() throws Exception { Stopwatch stopwatch = Stopwatch.createStarted(); - - - ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(5)); + ListeningExecutorService executorService = MoreExecutors.listeningDecorator(Executors.newFixedThreadPool(getMatchingStats().getMatchingStatsWorkers())); Map> jobsByConcept = new HashedMap<>(); Collection> jobs = jobsByConcept.values(); @@ -55,8 +53,6 @@ public void execute() throws Exception { jobsByConcept.put(concept.getId(), matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, 5)); } - log.info("WAITING for {} jobs in {}", jobs.size(), dataset); - while (jobs.stream().anyMatch(job -> job.state().equals(Future.State.RUNNING))) { for (ListenableFuture someJob : jobs) { if (someJob.isDone()) { @@ -64,18 +60,18 @@ public void execute() throws Exception { } try { - someJob.get(5, TimeUnit.SECONDS); + someJob.get(30, TimeUnit.SECONDS); } catch (Exception e) { // intentionally left blank } - log.debug("Still waiting for {} matching stats to finish.", jobs.stream().filter(job -> job.state().equals(Future.State.RUNNING)).count()); + log.debug("WAITING for {} matching stats to finish.", jobs.stream().filter(job -> job.state().equals(Future.State.RUNNING)).count()); } } for (Map.Entry> conceptState : jobsByConcept.entrySet()) { if (conceptState.getValue().state().equals(Future.State.FAILED)) { - log.warn("Failed to collect SQL matching stats for {}", conceptState.getKey(), conceptState.getValue().exceptionNow()); + log.warn("FAILED to collect SQL matching stats for {}", conceptState.getKey(), conceptState.getValue().exceptionNow()); } } 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..78d3571a47 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,13 @@ 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; + /** * 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/worker/LocalNamespace.java b/backend/src/main/java/com/bakdata/conquery/models/worker/LocalNamespace.java index 6ba57c5127..64d720b995 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()); } @@ -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 5e6af62cc5..a45804a53d 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 @@ -46,6 +46,7 @@ public class SqlMatchingStats { private final SqlFunctionProvider functionProvider; private final String defaultPrimaryColumn; private final int fetchBatchSize = 100; + private final int matchingStatsWorkers; private static void assignStatsToPath(ConceptElement element, Map, MatchingStats.Entry> matchingStats, String entity, CDateRange span) { while (element != null) { @@ -176,7 +177,6 @@ private Map, MatchingStats.Entry> readStats(TreeConcept conc log.debug("DONE fetching matching stats for {} within {}", concept.getId(), stopwatch); - return matchingStats; } @@ -287,6 +287,7 @@ private SelectJoinStep createMatchingStatsStatement(TreeConcep 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) { From 1e2bc61905de543ed8fa29e7cf4b603afa4ec58d Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Tue, 4 Aug 2026 16:37:55 +0200 Subject: [PATCH 6/7] Cleanup --- .../mode/local/LocalNamespaceHandler.java | 4 ++-- .../mode/local/UpdateMatchingStatsSqlJob.java | 11 ++++++--- .../config/DatabaseConnectionConfig.java | 7 ++++++ .../specific/UpdateMatchingStatsMessage.java | 2 ++ .../models/worker/LocalNamespace.java | 2 +- .../sql/conquery/SqlMatchingStats.java | 23 +++++++++++-------- .../sql/conversion/NodeConversions.java | 4 +++- 7 files changed, 37 insertions(+), 16 deletions(-) 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 b57d045e2d..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 @@ -67,8 +67,8 @@ public LocalNamespace createNamespace( connection.getConnection() ); } catch (Exception e) { - log.error("Failed to load namespaceStorage for {}", namespaceStorage.getDataset(), e); - throw e; + log.error("Failed to load namespaceStorage for {}", namespaceStorage.getPathName(), e); + throw e; } } 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 b639b3265b..0a288411f5 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 @@ -47,13 +47,18 @@ public void execute() throws Exception { Collection> jobs = jobsByConcept.values(); for (Concept concept : concepts) { - if (!(concept instanceof TreeConcept)) { - continue; + if (concept instanceof TreeConcept) { + jobsByConcept.put(concept.getId(), matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, getMatchingStats().getMatchingStatsRetries())); } - jobsByConcept.put(concept.getId(), matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, 5)); } while (jobs.stream().anyMatch(job -> job.state().equals(Future.State.RUNNING))) { + if (isCancelled()) { + for (ListenableFuture job : jobs) { + job.cancel(true); + } + } + for (ListenableFuture someJob : jobs) { if (someJob.isDone()) { continue; 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 78d3571a47..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 @@ -38,6 +38,13 @@ public class DatabaseConnectionConfig { @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 64d720b995..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 @@ -45,7 +45,7 @@ public LocalNamespace( this.dslContext = dslContext; this.storageHandler = storageHandler; this.dialect = dialect; - this.matchingStats = new SqlMatchingStats(dslContext, dialect.getFunctionProvider(), databaseConfig.getPrimaryColumn(), databaseConfig.getMatchingStatsWorkers()); + this.matchingStats = new SqlMatchingStats(dslContext, dialect.getFunctionProvider(), databaseConfig.getPrimaryColumn(), databaseConfig.getMatchingStatsWorkers(), databaseConfig.getMatchingStatsRetries()); } 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 a45804a53d..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 @@ -25,7 +25,6 @@ import org.jooq.*; import org.jooq.Record; import org.jooq.exception.DataAccessException; -import org.jooq.impl.DSL; import java.sql.Date; import java.util.*; @@ -36,17 +35,23 @@ @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)); + /** + * 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) { @@ -61,7 +66,7 @@ private static void assignStatsToPath(ConceptElement element, Map> collectAllFields(List conceptConditions) { + private static List> collectReferencedFields(List conceptConditions) { List> fields = conceptConditions.stream().flatMap(e -> e.conditions().keySet().stream()).distinct().toList(); return fields; } @@ -104,7 +109,7 @@ public void createConceptIdJoinTable(TreeConcept concept) { List conceptConditions = collectAllExpressions(concept, null, context); - List> allFields = collectAllFields(conceptConditions); + List> allFields = collectReferencedFields(conceptConditions); List rows = expressionsToRows(conceptConditions, allFields); @@ -139,7 +144,7 @@ private void assignStats(Map, MatchingStats.Entry> matchingS ConceptElementId conceptElementId = entry.getKey(); MatchingStats stats = new MatchingStats(); - stats.putEntry("sql", entry.getValue()); + stats.putEntry(SQL_SOURCE_MATCHING_STATS_LABEL, entry.getValue()); conceptElementId.resolve().setMatchingStats(stats); } } 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 a3fbc60985..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 @@ -33,7 +33,9 @@ public NodeConversions( IdColumnConfig idColumns, DialectBundle dialectBundle, DSLContext dslContext, - SqlExecutionService executionService, Clock clock, String defaultPrimaryColumn + SqlExecutionService executionService, + Clock clock, + String defaultPrimaryColumn ) { super(dialectBundle.getNodeConverters(dslContext)); this.idColumns = idColumns; From 4efd8d0ede8fd0d11c6c7d1aa1a21da8d7e72119 Mon Sep 17 00:00:00 2001 From: Fabian Kovacs Date: Mon, 10 Aug 2026 12:07:53 +0200 Subject: [PATCH 7/7] Cleanup 2 --- .../mode/local/ManagedConnection.java | 11 +++---- .../mode/local/UpdateMatchingStatsSqlJob.java | 33 +++++++++++-------- 2 files changed, 23 insertions(+), 21 deletions(-) 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 4b5e6fefd5..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 @@ -33,16 +33,13 @@ 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); - throw 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 0a288411f5..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,5 +1,12 @@ package com.bakdata.conquery.mode.local; +import java.util.Collection; +import java.util.List; +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; @@ -11,19 +18,14 @@ 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.collections4.map.HashedMap; -import java.util.Collection; -import java.util.List; -import java.util.Map; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; -import java.util.concurrent.TimeUnit; - @Slf4j @Data +@EqualsAndHashCode(callSuper = false) public class UpdateMatchingStatsSqlJob extends Job { @ToString.Exclude @@ -48,7 +50,16 @@ public void execute() throws Exception { for (Concept concept : concepts) { if (concept instanceof TreeConcept) { - jobsByConcept.put(concept.getId(), matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, getMatchingStats().getMatchingStatsRetries())); + ListenableFuture job = matchingStats.collectMatchingStatsForConcept((TreeConcept) concept, executorService, getMatchingStats().getMatchingStatsRetries()); + + job.addListener( + () -> { + if (job.state().equals(Future.State.FAILED)) { + log.warn("FAILED to collect SQL matching stats for {}", concept, job.exceptionNow()); + } + }, MoreExecutors.directExecutor()); + + jobsByConcept.put(concept.getId(), job); } } @@ -74,12 +85,6 @@ public void execute() throws Exception { } } - for (Map.Entry> conceptState : jobsByConcept.entrySet()) { - if (conceptState.getValue().state().equals(Future.State.FAILED)) { - log.warn("FAILED to collect SQL matching stats for {}", conceptState.getKey(), conceptState.getValue().exceptionNow()); - } - } - log.debug("DONE collecting SQL matching stats for {} within {}", dataset, stopwatch); }