diff --git a/deployment/helm/ditto/Chart.yaml b/deployment/helm/ditto/Chart.yaml index 140f74d1e4..4b24d7f017 100644 --- a/deployment/helm/ditto/Chart.yaml +++ b/deployment/helm/ditto/Chart.yaml @@ -16,7 +16,7 @@ description: | A digital twin is a virtual, cloud based, representation of his real world counterpart (real world “Things”, e.g. devices like sensors, smart heating, connected cars, smart grids, EV charging stations etc). type: application -version: 4.6.0 # chart version is effectively set by release-job +version: 4.6.1 # chart version is effectively set by release-job appVersion: 3.9.6 keywords: - iot-chart diff --git a/deployment/helm/ditto/service-config/search-extension.conf.tpl b/deployment/helm/ditto/service-config/search-extension.conf.tpl index c276512b2e..c23a253023 100644 --- a/deployment/helm/ditto/service-config/search-extension.conf.tpl +++ b/deployment/helm/ditto/service-config/search-extension.conf.tpl @@ -146,6 +146,28 @@ ditto { } {{- end }} } + + {{- with .Values.thingsSearch.config.operatorMetrics.customMetricsPersistence }} + {{- if .readPreference }} + custom-metrics-persistence { + readPreference = "{{ .readPreference }}" + {{- if .readConcern }} + readConcern = "{{ .readConcern }}" + {{- end }} + } + {{- end }} + {{- end }} + + {{- with .Values.thingsSearch.config.operatorMetrics.customAggregationMetricsPersistence }} + {{- if .readPreference }} + custom-aggregation-metrics-persistence { + readPreference = "{{ .readPreference }}" + {{- if .readConcern }} + readConcern = "{{ .readConcern }}" + {{- end }} + } + {{- end }} + {{- end }} } } } diff --git a/deployment/helm/ditto/values.yaml b/deployment/helm/ditto/values.yaml index c6946f4e92..59eeba9e89 100644 --- a/deployment/helm/ditto/values.yaml +++ b/deployment/helm/ditto/values.yaml @@ -1729,6 +1729,19 @@ thingsSearch: # # or as index key spec: # # indexHint: # # "t.attributes.location": 1 + # customMetricsPersistence optionally overrides the MongoDB read preference / read concern used for the + # count based "customMetrics" queries above. When unset (the default), the general query.persistence + # settings are used - so this can offload the periodic operator metric counts to a secondary node without + # affecting user-facing search/count requests. + # customMetricsPersistence: + # readPreference: secondaryPreferred + # readConcern: default + # customAggregationMetricsPersistence optionally overrides the MongoDB read preference / read concern used + # for the "customAggregationMetrics" ($group) queries above. When unset (the default), the general + # query.persistence settings are used. Can be configured independently from customMetricsPersistence. + # customAggregationMetricsPersistence: + # readPreference: secondaryPreferred + # readConcern: default # dispatchers contains tuning for the things-search service's Pekko dispatchers. # Defaults below match the in-jar HOCON defaults and only need overriding under sustained load. dispatchers: diff --git a/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadConcern.java b/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadConcern.java index 98586ac91e..48f6ff2aee 100644 --- a/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadConcern.java +++ b/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadConcern.java @@ -42,6 +42,16 @@ public com.mongodb.ReadConcern getMongoReadConcern() { return mongoReadConcern; } + /** + * Returns the config string value of this read concern (e.g. {@code "local"}). + * + * @return the read concern config string. + * @since 3.9.7 + */ + public String getName() { + return name; + } + /** * Tries to create a ReadConcern from the given read concern string. * diff --git a/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadPreference.java b/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadPreference.java index b80dae2b4b..a06556df3d 100644 --- a/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadPreference.java +++ b/internal/utils/persistence/src/main/java/org/eclipse/ditto/internal/utils/persistence/mongo/config/ReadPreference.java @@ -40,6 +40,16 @@ public com.mongodb.ReadPreference getMongoReadPreference() { return mongoReadPreference; } + /** + * Returns the config string value of this read preference (e.g. {@code "secondaryPreferred"}). + * + * @return the read preference config string. + * @since 3.9.7 + */ + public String getName() { + return name; + } + /** * Tries to create a ReadPreference from the given read preference string. * diff --git a/thingsearch/api/src/main/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThings.java b/thingsearch/api/src/main/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThings.java index 225b80ae36..0e11deda9b 100644 --- a/thingsearch/api/src/main/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThings.java +++ b/thingsearch/api/src/main/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThings.java @@ -73,6 +73,14 @@ public final class SudoCountThings extends AbstractCommand JsonFactory.newJsonValueFieldDefinition("indexHint", FieldType.REGULAR, JsonSchemaVersion.V_2); + static final JsonFieldDefinition JSON_READ_PREFERENCE = + JsonFactory.newStringFieldDefinition("readPreference", FieldType.REGULAR, + JsonSchemaVersion.V_2); + + static final JsonFieldDefinition JSON_READ_CONCERN = + JsonFactory.newStringFieldDefinition("readConcern", FieldType.REGULAR, + JsonSchemaVersion.V_2); + @Nullable private final String filter; @@ -82,8 +90,15 @@ public final class SudoCountThings extends AbstractCommand @Nullable private final JsonValue indexHint; + @Nullable + private final String readPreference; + + @Nullable + private final String readConcern; + private SudoCountThings(final DittoHeaders dittoHeaders, @Nullable final String filter, - @Nullable final Collection namespaces, @Nullable final JsonValue indexHint) { + @Nullable final Collection namespaces, @Nullable final JsonValue indexHint, + @Nullable final String readPreference, @Nullable final String readConcern) { super(TYPE, dittoHeaders); this.filter = filter; if (namespaces != null) { @@ -92,6 +107,8 @@ private SudoCountThings(final DittoHeaders dittoHeaders, @Nullable final String this.namespaces = null; } this.indexHint = indexHint; + this.readPreference = readPreference; + this.readConcern = readConcern; } /** @@ -103,7 +120,7 @@ private SudoCountThings(final DittoHeaders dittoHeaders, @Nullable final String * @throws NullPointerException if {@code dittoHeaders} is {@code null}. */ public static SudoCountThings of(@Nullable final String filter, final DittoHeaders dittoHeaders) { - return new SudoCountThings(dittoHeaders, filter, null, null); + return new SudoCountThings(dittoHeaders, filter, null, null, null, null); } /** @@ -117,7 +134,7 @@ public static SudoCountThings of(@Nullable final String filter, final DittoHeade */ public static SudoCountThings of(@Nullable final String filter, @Nullable final Collection namespaces, final DittoHeaders dittoHeaders) { - return new SudoCountThings(dittoHeaders, filter, namespaces, null); + return new SudoCountThings(dittoHeaders, filter, namespaces, null, null, null); } /** @@ -132,7 +149,47 @@ public static SudoCountThings of(@Nullable final String filter, @Nullable final */ public static SudoCountThings of(@Nullable final String filter, @Nullable final Collection namespaces, @Nullable final JsonValue indexHint, final DittoHeaders dittoHeaders) { - return new SudoCountThings(dittoHeaders, filter, namespaces, indexHint); + return new SudoCountThings(dittoHeaders, filter, namespaces, indexHint, null, null); + } + + /** + * Returns a new instance of {@code SudoCountThings}. + * + * @param filter the optional filter string. + * @param namespaces the namespaces to perform the count in. + * @param indexHint the optional index hint for the MongoDB query. + * @param readPreference the optional MongoDB read preference to use for this count query (e.g. + * {@code "secondaryPreferred"}); when {@code null} the persistence default read preference is used. + * @param dittoHeaders the headers of the command. + * @return a new command for counting Things. + * @throws NullPointerException if {@code dittoHeaders} is {@code null}. + * @since 3.9.7 + */ + public static SudoCountThings of(@Nullable final String filter, @Nullable final Collection namespaces, + @Nullable final JsonValue indexHint, @Nullable final String readPreference, + final DittoHeaders dittoHeaders) { + return new SudoCountThings(dittoHeaders, filter, namespaces, indexHint, readPreference, null); + } + + /** + * Returns a new instance of {@code SudoCountThings}. + * + * @param filter the optional filter string. + * @param namespaces the namespaces to perform the count in. + * @param indexHint the optional index hint for the MongoDB query. + * @param readPreference the optional MongoDB read preference to use for this count query (e.g. + * {@code "secondaryPreferred"}); when {@code null} the persistence default read preference is used. + * @param readConcern the optional MongoDB read concern to use for this count query (e.g. {@code "local"}); + * when {@code null} the persistence default read concern is used. + * @param dittoHeaders the headers of the command. + * @return a new command for counting Things. + * @throws NullPointerException if {@code dittoHeaders} is {@code null}. + * @since 3.9.7 + */ + public static SudoCountThings of(@Nullable final String filter, @Nullable final Collection namespaces, + @Nullable final JsonValue indexHint, @Nullable final String readPreference, + @Nullable final String readConcern, final DittoHeaders dittoHeaders) { + return new SudoCountThings(dittoHeaders, filter, namespaces, indexHint, readPreference, readConcern); } /** @@ -143,7 +200,7 @@ public static SudoCountThings of(@Nullable final String filter, @Nullable final * @throws NullPointerException if any argument is {@code null}. */ public static SudoCountThings of(final DittoHeaders dittoHeaders) { - return new SudoCountThings(dittoHeaders, null, null, null); + return new SudoCountThings(dittoHeaders, null, null, null, null, null); } /** @@ -186,7 +243,12 @@ public static SudoCountThings fromJson(final JsonObject jsonObject, final DittoH final JsonValue extractedIndexHint = jsonObject.getValue(JSON_INDEX_HINT).orElse(null); - return new SudoCountThings(dittoHeaders, extractedFilter, extractedNamespaces, extractedIndexHint); + final String extractedReadPreference = jsonObject.getValue(JSON_READ_PREFERENCE).orElse(null); + + final String extractedReadConcern = jsonObject.getValue(JSON_READ_CONCERN).orElse(null); + + return new SudoCountThings(dittoHeaders, extractedFilter, extractedNamespaces, extractedIndexHint, + extractedReadPreference, extractedReadConcern); }); } @@ -217,6 +279,26 @@ public Optional getIndexHint() { return Optional.ofNullable(indexHint); } + /** + * Get the optional MongoDB read preference override for this count query. + * + * @return the optional read preference (e.g. {@code "secondaryPreferred"}). + * @since 3.9.7 + */ + public Optional getReadPreference() { + return Optional.ofNullable(readPreference); + } + + /** + * Get the optional MongoDB read concern override for this count query. + * + * @return the optional read concern (e.g. {@code "local"}). + * @since 3.9.7 + */ + public Optional getReadConcern() { + return Optional.ofNullable(readConcern); + } + @Override protected void appendPayload(final JsonObjectBuilder jsonObjectBuilder, final JsonSchemaVersion schemaVersion, final Predicate thePredicate) { @@ -227,6 +309,8 @@ protected void appendPayload(final JsonObjectBuilder jsonObjectBuilder, final Js .map(JsonValue::of) .collect(JsonCollectors.valuesToArray()), predicate)); getIndexHint().ifPresent(hint -> jsonObjectBuilder.set(JSON_INDEX_HINT, hint, predicate)); + getReadPreference().ifPresent(rp -> jsonObjectBuilder.set(JSON_READ_PREFERENCE, rp, predicate)); + getReadConcern().ifPresent(rc -> jsonObjectBuilder.set(JSON_READ_CONCERN, rc, predicate)); } @Override @@ -236,7 +320,7 @@ public Category getCategory() { @Override public SudoCountThings setDittoHeaders(final DittoHeaders dittoHeaders) { - return new SudoCountThings(dittoHeaders, filter, namespaces, indexHint); + return new SudoCountThings(dittoHeaders, filter, namespaces, indexHint, readPreference, readConcern); } @Override @@ -250,12 +334,14 @@ public boolean equals(@Nullable final Object o) { final SudoCountThings that = (SudoCountThings) o; return Objects.equals(filter, that.filter) && Objects.equals(namespaces, that.namespaces) && - Objects.equals(indexHint, that.indexHint); + Objects.equals(indexHint, that.indexHint) && + Objects.equals(readPreference, that.readPreference) && + Objects.equals(readConcern, that.readConcern); } @Override public int hashCode() { - return Objects.hash(super.hashCode(), filter, namespaces, indexHint); + return Objects.hash(super.hashCode(), filter, namespaces, indexHint, readPreference, readConcern); } @Override @@ -264,6 +350,8 @@ public String toString() { "filter='" + filter + "'" + ", namespaces=" + namespaces + ", indexHint=" + indexHint + + ", readPreference=" + readPreference + + ", readConcern=" + readConcern + "]"; } diff --git a/thingsearch/api/src/test/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThingsTest.java b/thingsearch/api/src/test/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThingsTest.java index b080b8447b..c3d8f57202 100644 --- a/thingsearch/api/src/test/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThingsTest.java +++ b/thingsearch/api/src/test/java/org/eclipse/ditto/thingsearch/api/commands/sudo/SudoCountThingsTest.java @@ -135,4 +135,42 @@ public void setDittoHeadersPreservesIndexHint() { assertThat(withNewHeaders.getIndexHint()).contains(JsonValue.of("my_index")); } + @Test + public void toJsonWithReadPreferenceAndReadConcern() { + final SudoCountThings command = SudoCountThings.of( + KNOWN_FILTER_STR, List.of("ns1"), JsonValue.of("my_index"), "secondaryPreferred", "local", + DittoHeaders.empty()); + + final String json = command.toJsonString(); + final SudoCountThings deserialized = SudoCountThings.fromJson(json, DittoHeaders.empty()); + + assertThat(deserialized.getFilter()).contains(KNOWN_FILTER_STR); + assertThat(deserialized.getIndexHint()).contains(JsonValue.of("my_index")); + assertThat(deserialized.getReadPreference()).contains("secondaryPreferred"); + assertThat(deserialized.getReadConcern()).contains("local"); + } + + @Test + public void toJsonWithoutReadPreferenceOrReadConcern() { + final SudoCountThings command = SudoCountThings.of(KNOWN_FILTER_STR, DittoHeaders.empty()); + + final String json = command.toJsonString(); + final SudoCountThings deserialized = SudoCountThings.fromJson(json, DittoHeaders.empty()); + + assertThat(deserialized.getReadPreference()).isEmpty(); + assertThat(deserialized.getReadConcern()).isEmpty(); + } + + @Test + public void setDittoHeadersPreservesReadPreferenceAndReadConcern() { + final SudoCountThings command = SudoCountThings.of( + KNOWN_FILTER_STR, null, null, "secondaryPreferred", "local", DittoHeaders.empty()); + + final SudoCountThings withNewHeaders = command.setDittoHeaders( + DittoHeaders.newBuilder().correlationId("test").build()); + + assertThat(withNewHeaders.getReadPreference()).contains("secondaryPreferred"); + assertThat(withNewHeaders.getReadConcern()).contains("local"); + } + } diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfig.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfig.java index a8bdfe8502..63f505863d 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfig.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfig.java @@ -17,6 +17,7 @@ import java.util.LinkedHashMap; import java.util.Map; import java.util.Objects; +import java.util.Optional; import java.util.Set; import java.util.function.BiConsumer; import java.util.function.BinaryOperator; @@ -26,6 +27,7 @@ import java.util.stream.Collectors; import java.util.stream.Stream; +import javax.annotation.Nullable; import javax.annotation.concurrent.Immutable; import org.eclipse.ditto.internal.utils.config.ConfigWithFallback; @@ -47,10 +49,22 @@ public final class DefaultOperatorMetricsConfig implements OperatorMetricsConfig */ static final String CONFIG_PATH = "operator-metrics"; + /** + * Path of the optional persistence config block used for the count based {@code custom-metrics} queries. + */ + static final String CUSTOM_METRICS_PERSISTENCE_PATH = "custom-metrics-persistence"; + + /** + * Path of the optional persistence config block used for the {@code custom-aggregation-metrics} queries. + */ + static final String CUSTOM_AGGREGATION_METRICS_PERSISTENCE_PATH = "custom-aggregation-metrics-persistence"; + private final boolean enabled; private final Duration scrapeInterval; private final Map customMetricConfigurations; private final Map customAggregationMetricConfigs; + @Nullable private final SearchPersistenceConfig customMetricsPersistenceConfig; + @Nullable private final SearchPersistenceConfig customAggregationMetricsPersistenceConfig; private DefaultOperatorMetricsConfig(final ConfigWithFallback updaterScopedConfig) { enabled = updaterScopedConfig.getBoolean(OperatorMetricsConfigValue.ENABLED.getConfigPath()); @@ -59,6 +73,14 @@ private DefaultOperatorMetricsConfig(final ConfigWithFallback updaterScopedConfi OperatorMetricsConfigValue.CUSTOM_METRICS); customAggregationMetricConfigs = loadCustomAggregatedMetricConfigurations(updaterScopedConfig, OperatorMetricsConfigValue.CUSTOM_AGGREGATION_METRIC); + customMetricsPersistenceConfig = updaterScopedConfig.hasPath(CUSTOM_METRICS_PERSISTENCE_PATH) + ? DefaultSearchPersistenceConfig.of(updaterScopedConfig, CUSTOM_METRICS_PERSISTENCE_PATH) + : null; + customAggregationMetricsPersistenceConfig = + updaterScopedConfig.hasPath(CUSTOM_AGGREGATION_METRICS_PERSISTENCE_PATH) + ? DefaultSearchPersistenceConfig.of(updaterScopedConfig, + CUSTOM_AGGREGATION_METRICS_PERSISTENCE_PATH) + : null; } /** @@ -99,12 +121,18 @@ public boolean equals(final Object o) { } final DefaultOperatorMetricsConfig that = (DefaultOperatorMetricsConfig) o; return enabled == that.enabled && - Objects.equals(scrapeInterval, that.scrapeInterval); + Objects.equals(scrapeInterval, that.scrapeInterval) && + Objects.equals(customMetricConfigurations, that.customMetricConfigurations) && + Objects.equals(customAggregationMetricConfigs, that.customAggregationMetricConfigs) && + Objects.equals(customMetricsPersistenceConfig, that.customMetricsPersistenceConfig) && + Objects.equals(customAggregationMetricsPersistenceConfig, + that.customAggregationMetricsPersistenceConfig); } @Override public int hashCode() { - return Objects.hash(enabled, scrapeInterval, customMetricConfigurations); + return Objects.hash(enabled, scrapeInterval, customMetricConfigurations, customAggregationMetricConfigs, + customMetricsPersistenceConfig, customAggregationMetricsPersistenceConfig); } @Override @@ -113,6 +141,9 @@ public String toString() { "enabled=" + enabled + ", scrapeInterval=" + scrapeInterval + ", customMetricConfigurations=" + customMetricConfigurations + + ", customAggregationMetricConfigs=" + customAggregationMetricConfigs + + ", customMetricsPersistenceConfig=" + customMetricsPersistenceConfig + + ", customAggregationMetricsPersistenceConfig=" + customAggregationMetricsPersistenceConfig + "]"; } @@ -136,6 +167,16 @@ public Map getCustomAggregationMetricConf return customAggregationMetricConfigs; } + @Override + public Optional getCustomMetricsPersistenceConfig() { + return Optional.ofNullable(customMetricsPersistenceConfig); + } + + @Override + public Optional getCustomAggregationMetricsPersistenceConfig() { + return Optional.ofNullable(customAggregationMetricsPersistenceConfig); + } + private static class CustomMetricConfigCollector implements Collector, Map, Map> { diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultSearchPersistenceConfig.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultSearchPersistenceConfig.java index 66d937b935..c41c59f296 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultSearchPersistenceConfig.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultSearchPersistenceConfig.java @@ -69,7 +69,22 @@ private DefaultSearchPersistenceConfig(final ConfigWithFallback config) { * @throws DittoConfigError if {@code config} is invalid. */ public static DefaultSearchPersistenceConfig of(final Config config) { - return new DefaultSearchPersistenceConfig(ConfigWithFallback.newInstance(config, CONFIG_PATH, ConfigValue.values())); + return of(config, CONFIG_PATH); + } + + /** + * Returns an instance of DefaultSearchPersistenceConfig based on the settings of the specified Config at the + * given {@code configPath}. + * + * @param config the Config which provides the persistence settings at {@code configPath}. + * @param configPath the path under {@code config} at which the persistence settings are expected. + * @return the instance. + * @throws DittoConfigError if {@code config} is invalid. + * @since 3.9.7 + */ + public static DefaultSearchPersistenceConfig of(final Config config, final String configPath) { + return new DefaultSearchPersistenceConfig( + ConfigWithFallback.newInstance(config, configPath, ConfigValue.values())); } diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/OperatorMetricsConfig.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/OperatorMetricsConfig.java index 49b1d6f726..ac4fc0ac27 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/OperatorMetricsConfig.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/common/config/OperatorMetricsConfig.java @@ -15,6 +15,7 @@ import java.time.Duration; import java.util.Collections; import java.util.Map; +import java.util.Optional; import javax.annotation.concurrent.Immutable; @@ -54,6 +55,25 @@ public interface OperatorMetricsConfig { */ Map getCustomAggregationMetricConfigs(); + /** + * Returns the optional persistence (read preference / read concern) configuration to use for the count based + * {@code custom-metrics} queries. When empty, the general {@code query.persistence} configuration is used. + * + * @return the optional persistence configuration for count based custom metrics. + * @since 3.9.7 + */ + Optional getCustomMetricsPersistenceConfig(); + + /** + * Returns the optional persistence (read preference / read concern) configuration to use for the + * {@code custom-aggregation-metrics} ({@code $group}) queries. When empty, the general {@code query.persistence} + * configuration is used. + * + * @return the optional persistence configuration for custom aggregation metrics. + * @since 3.9.7 + */ + Optional getCustomAggregationMetricsPersistenceConfig(); + /** * An enumeration of the known config path expressions and their associated default values for * OperatorMetricsConfig. diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsAggregationPersistence.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsAggregationPersistence.java index 62b77b6d2e..373535e129 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsAggregationPersistence.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsAggregationPersistence.java @@ -89,8 +89,13 @@ private MongoThingsAggregationPersistence(final DittoMongoClient mongoClient, public static ThingsAggregationPersistence of(final DittoMongoClient mongoClient, final SearchConfig searchConfig, final LoggingAdapter log) { + // use the dedicated operator-metrics aggregation persistence config (read preference / read concern) if + // configured, otherwise fall back to the general query persistence config: + final SearchPersistenceConfig persistenceConfig = searchConfig.getOperatorMetricsConfig() + .getCustomAggregationMetricsPersistenceConfig() + .orElseGet(searchConfig::getQueryPersistenceConfig); return new MongoThingsAggregationPersistence(mongoClient, searchConfig.getMongoHintsByNamespace(), - searchConfig.getSimpleFieldMappings(), searchConfig.getQueryPersistenceConfig(), log); + searchConfig.getSimpleFieldMappings(), persistenceConfig, log); } @Override diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsSearchPersistence.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsSearchPersistence.java index d4fbb2c682..bbcb617660 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsSearchPersistence.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/MongoThingsSearchPersistence.java @@ -48,6 +48,8 @@ import org.eclipse.ditto.internal.utils.persistence.mongo.DittoBsonJson; import org.eclipse.ditto.internal.utils.persistence.mongo.DittoMongoClient; import org.eclipse.ditto.internal.utils.persistence.mongo.config.IndexInitializationConfig; +import org.eclipse.ditto.internal.utils.persistence.mongo.config.ReadConcern; +import org.eclipse.ditto.internal.utils.persistence.mongo.config.ReadPreference; import org.eclipse.ditto.internal.utils.persistence.mongo.indices.Index; import org.eclipse.ditto.internal.utils.persistence.mongo.indices.IndexInitializer; import org.eclipse.ditto.json.JsonArray; @@ -215,6 +217,13 @@ public Source sudoCount(final Query query, final DittoHeaders dit @Override public Source sudoCount(final Query query, final DittoHeaders dittoHeaders, @Nullable final JsonValue indexHint) { + return sudoCount(query, dittoHeaders, indexHint, null, null); + } + + @Override + public Source sudoCount(final Query query, final DittoHeaders dittoHeaders, + @Nullable final JsonValue indexHint, @Nullable final String readPreferenceOverride, + @Nullable final String readConcernOverride) { checkNotNull(query, "query"); @@ -227,11 +236,38 @@ public Source sudoCount(final Query query, final DittoHeaders dit .maxTime(maxQueryTime.getSeconds(), TimeUnit.SECONDS); applyHint(countOptions, indexHint); - return Source.fromPublisher(collection.countDocuments(queryFilter, countOptions)) + return Source.fromPublisher(withReadOverrides(readPreferenceOverride, readConcernOverride, dittoHeaders) + .countDocuments(queryFilter, countOptions)) .mapError(handleMongoExecutionTimeExceededException(dittoHeaders)) .log("sudoCount"); } + private MongoCollection withReadOverrides(@Nullable final String readPreferenceOverride, + @Nullable final String readConcernOverride, final DittoHeaders dittoHeaders) { + MongoCollection result = collection; + if (readPreferenceOverride != null) { + final Optional parsed = ReadPreference.ofReadPreference(readPreferenceOverride); + if (parsed.isPresent()) { + result = result.withReadPreference(parsed.get().getMongoReadPreference()); + } else { + LOGGER.withCorrelationId(dittoHeaders) + .warn("Ignoring unknown read preference override <{}> for count query.", + readPreferenceOverride); + } + } + if (readConcernOverride != null) { + final Optional parsed = ReadConcern.ofReadConcern(readConcernOverride); + if (parsed.isPresent()) { + result = result.withReadConcern(parsed.get().getMongoReadConcern()); + } else { + LOGGER.withCorrelationId(dittoHeaders) + .warn("Ignoring unknown read concern override <{}> for count query.", + readConcernOverride); + } + } + return result; + } + private void applyHint(final CountOptions countOptions, @Nullable final JsonValue indexHint) { if (indexHint != null) { if (indexHint.isString()) { diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/ThingsSearchPersistence.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/ThingsSearchPersistence.java index d9fee7184f..93b9e799b9 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/ThingsSearchPersistence.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/persistence/read/ThingsSearchPersistence.java @@ -90,6 +90,27 @@ default Source sudoCount(Query query, DittoHeaders dittoHeaders, return sudoCount(query, dittoHeaders); } + /** + * Returns the count of documents found by the given {@code query} regardless of visibility, + * using an optional per-query index hint and an optional per-query read preference override. + * + * @param query the query for matching. + * @param dittoHeaders the headers of the request. + * @param indexHint the optional index hint (string for index name, object for index key spec). + * @param readPreferenceOverride the optional MongoDB read preference to use for this query (e.g. + * {@code "secondaryPreferred"}); when {@code null} the persistence default read preference is used. + * @param readConcernOverride the optional MongoDB read concern to use for this query (e.g. {@code "local"}); + * when {@code null} the persistence default read concern is used. + * @return an {@link Source} which emits the count. + * @throws NullPointerException if {@code query} is {@code null}. + * @since 3.9.7 + */ + default Source sudoCount(Query query, DittoHeaders dittoHeaders, + @Nullable JsonValue indexHint, @Nullable String readPreferenceOverride, + @Nullable String readConcernOverride) { + return sudoCount(query, dittoHeaders, indexHint); + } + /** * Returns the IDs for all found documents. * diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/OperatorMetricsProviderActor.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/OperatorMetricsProviderActor.java index a14046f664..a8e26fccbc 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/OperatorMetricsProviderActor.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/OperatorMetricsProviderActor.java @@ -19,6 +19,8 @@ import java.util.UUID; import java.util.concurrent.ThreadLocalRandom; +import javax.annotation.Nullable; + import org.apache.pekko.actor.AbstractActorWithTimers; import org.apache.pekko.actor.ActorRef; import org.apache.pekko.actor.Props; @@ -57,12 +59,20 @@ public final class OperatorMetricsProviderActor extends AbstractActorWithTimers private final ActorRef searchActor; private final Map metricsGauges; + @Nullable private final String readPreferenceOverride; + @Nullable private final String readConcernOverride; @SuppressWarnings("unused") private OperatorMetricsProviderActor(final OperatorMetricsConfig operatorMetricsConfig, final ActorRef searchActor) { this.searchActor = searchActor; + readPreferenceOverride = operatorMetricsConfig.getCustomMetricsPersistenceConfig() + .map(persistenceConfig -> persistenceConfig.readPreference().getName()) + .orElse(null); + readConcernOverride = operatorMetricsConfig.getCustomMetricsPersistenceConfig() + .map(persistenceConfig -> persistenceConfig.readConcern().getName()) + .orElse(null); metricsGauges = new HashMap<>(); operatorMetricsConfig.getCustomMetricConfigurations().forEach((metricName, config) -> { if (config.isEnabled()) { @@ -131,7 +141,7 @@ private void handleGatheringMetrics(final GatherMetrics gatherMetrics) { .build(); final SudoCountThings sudoCountThings = SudoCountThings.of( filter.isEmpty() ? null : filter, namespaces.isEmpty() ? null : namespaces, - config.getIndexHint().orElse(null), dittoHeaders); + config.getIndexHint().orElse(null), readPreferenceOverride, readConcernOverride, dittoHeaders); final long startTs = System.nanoTime(); log.withCorrelationId(dittoHeaders) diff --git a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActor.java b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActor.java index 0b6f524546..a184155097 100644 --- a/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActor.java +++ b/thingsearch/service/src/main/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActor.java @@ -327,7 +327,9 @@ private > CompletionStage executeCount(final T coun (theQuery, headers) -> { if (isSudo && tracedCountCommand instanceof SudoCountThings sudoCmd) { return searchPersistence.sudoCount(theQuery, headers, - sudoCmd.getIndexHint().orElse(null)); + sudoCmd.getIndexHint().orElse(null), + sudoCmd.getReadPreference().orElse(null), + sudoCmd.getReadConcern().orElse(null)); } else if (isSudo) { return searchPersistence.sudoCount(theQuery, headers); } else { diff --git a/thingsearch/service/src/main/resources/search.conf b/thingsearch/service/src/main/resources/search.conf index cc898a5310..48fec8dcaa 100755 --- a/thingsearch/service/src/main/resources/search.conf +++ b/thingsearch/service/src/main/resources/search.conf @@ -402,6 +402,28 @@ ditto { # # index-hint { "t.attributes.location": 1 } # } } + + # Optional dedicated MongoDB read settings for the count based "custom-metrics" queries above. + # When this block is absent (the default), the general "query.persistence" read preference / read concern + # is used. Configure it to e.g. offload the periodic operator metric counts to a secondary node without + # affecting user-facing search/count requests: + # custom-metrics-persistence { + # readPreference = "secondaryPreferred" + # readPreference = ${?THINGS_SEARCH_OPERATOR_METRICS_CUSTOM_METRICS_READ_PREFERENCE} + # readConcern = "default" + # readConcern = ${?THINGS_SEARCH_OPERATOR_METRICS_CUSTOM_METRICS_READ_CONCERN} + # } + + # Optional dedicated MongoDB read settings for the "custom-aggregation-metrics" ($group) queries. + # When this block is absent (the default), the general "query.persistence" read preference / read concern + # is used. Configure it independently from the count based "custom-metrics-persistence" above, e.g. to run + # the (potentially expensive) aggregation scans on a secondary node: + # custom-aggregation-metrics-persistence { + # readPreference = "secondaryPreferred" + # readPreference = ${?THINGS_SEARCH_OPERATOR_METRICS_CUSTOM_AGGREGATION_METRICS_READ_PREFERENCE} + # readConcern = "default" + # readConcern = ${?THINGS_SEARCH_OPERATOR_METRICS_CUSTOM_AGGREGATION_METRICS_READ_CONCERN} + # } } } } diff --git a/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfigTest.java b/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfigTest.java new file mode 100644 index 0000000000..35cec93890 --- /dev/null +++ b/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/common/config/DefaultOperatorMetricsConfigTest.java @@ -0,0 +1,74 @@ +/* + * Copyright (c) 2026 Contributors to the Eclipse Foundation + * + * See the NOTICE file(s) distributed with this work for additional + * information regarding copyright ownership. + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0 + * + * SPDX-License-Identifier: EPL-2.0 + */ +package org.eclipse.ditto.thingsearch.service.common.config; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.eclipse.ditto.internal.utils.persistence.mongo.config.ReadConcern; +import org.eclipse.ditto.internal.utils.persistence.mongo.config.ReadPreference; +import org.junit.Test; + +import com.typesafe.config.Config; +import com.typesafe.config.ConfigFactory; + +/** + * Unit test for {@link DefaultOperatorMetricsConfig}, focusing on the optional per-metric-type persistence + * (read preference / read concern) configuration. + */ +public final class DefaultOperatorMetricsConfigTest { + + @Test + public void persistenceConfigsAreEmptyWhenNotConfigured() { + final Config config = ConfigFactory.parseString( + "operator-metrics {\n" + + " enabled = true\n" + + " scrape-interval = 15m\n" + + " custom-metrics {}\n" + + " custom-aggregation-metrics {}\n" + + "}"); + + final DefaultOperatorMetricsConfig underTest = DefaultOperatorMetricsConfig.of(config); + + assertThat(underTest.getCustomMetricsPersistenceConfig()).isEmpty(); + assertThat(underTest.getCustomAggregationMetricsPersistenceConfig()).isEmpty(); + } + + @Test + public void persistenceConfigsAreParsedWhenConfigured() { + final Config config = ConfigFactory.parseString( + "operator-metrics {\n" + + " enabled = true\n" + + " scrape-interval = 15m\n" + + " custom-metrics {}\n" + + " custom-aggregation-metrics {}\n" + + " custom-metrics-persistence {\n" + + " readPreference = \"nearest\"\n" + + " readConcern = \"local\"\n" + + " }\n" + + " custom-aggregation-metrics-persistence {\n" + + " readPreference = \"secondaryPreferred\"\n" + + " }\n" + + "}"); + + final DefaultOperatorMetricsConfig underTest = DefaultOperatorMetricsConfig.of(config); + + assertThat(underTest.getCustomMetricsPersistenceConfig()).hasValueSatisfying(persistenceConfig -> { + assertThat(persistenceConfig.readPreference()).isEqualTo(ReadPreference.NEAREST); + assertThat(persistenceConfig.readConcern()).isEqualTo(ReadConcern.LOCAL); + }); + assertThat(underTest.getCustomAggregationMetricsPersistenceConfig()).hasValueSatisfying(persistenceConfig -> + // read concern falls back to its own default when not explicitly configured: + assertThat(persistenceConfig.readPreference()).isEqualTo(ReadPreference.SECONDARY_PREFERRED)); + } + +} diff --git a/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActorTest.java b/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActorTest.java index 1364122a02..6308f3e072 100644 --- a/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActorTest.java +++ b/thingsearch/service/src/test/java/org/eclipse/ditto/thingsearch/service/starter/actors/SearchActorTest.java @@ -122,7 +122,7 @@ public void waitForQueries() { final var serviceRequestsDone = SearchActor.Control.SERVICE_REQUESTS_DONE; final var countActor = use(p -> p.count(any(), any(), any())); - final var sudoCountActor = use(p -> p.sudoCount(any(), any(), any())); + final var sudoCountActor = use(p -> p.sudoCount(any(), any(), any(), any(), any())); final var queryActor = use(p -> p.findAll(any(), any(), any(), any())); final var shutdownProbe = TestProbe.apply(actorSystemResource.getActorSystem());