This is an automated email from the ASF dual-hosted git repository.
dominikriemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git
The following commit(s) were added to refs/heads/dev by this push:
new 67015eac37 chore: Improve query performance of latest event queries
(#4685)
67015eac37 is described below
commit 67015eac3732f1e57a83316aaece076cc9e62d1f
Author: Dominik Riemer <[email protected]>
AuthorDate: Mon Jul 6 21:39:36 2026 +0200
chore: Improve query performance of latest event queries (#4685)
---
.../api/IDataExplorerQueryManagement.java | 31 ++
.../influx/DataExplorerInfluxQueryExecutor.java | 74 +++++
.../influx/DataExplorerQueryManagementInflux.java | 88 ++++++
.../DataExplorerInfluxQueryExecutorTest.java | 41 +++
.../influx/InfluxLatestTimestampBenchmarkTest.java | 331 +++++++++++++++++++++
.../dataexplorer/influx/SelectQueryParamsTest.java | 27 ++
.../rest/impl/datalake/DataLakeResource.java | 25 +-
.../DataLakeLatestEventsBenchmarkTest.java | 108 +++++++
.../rest/impl/datalake/DataLakeResourceTest.java | 84 ++++++
9 files changed, 785 insertions(+), 24 deletions(-)
diff --git
a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
index d114cb2886..a5f68204cd 100644
---
a/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
+++
b/streampipes-data-explorer-api/src/main/java/org/apache/streampipes/dataexplorer/api/IDataExplorerQueryManagement.java
@@ -25,8 +25,16 @@ import
org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
import java.io.IOException;
import java.io.OutputStream;
+import java.util.HashMap;
+import java.util.List;
import java.util.Map;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_END_DATE;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_LIMIT;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_MISSING_VALUE_BEHAVIOUR;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_ORDER;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_START_DATE;
+
public interface IDataExplorerQueryManagement {
SpQueryResult getData(
@@ -47,4 +55,27 @@ public interface IDataExplorerQueryManagement {
Map<String, Object> getTagValues(String measurementId,
String fields);
+
+ default Map<String, Long> getLatestTimestamps(List<String> measurementNames)
{
+ Map<String, Long> latestTimestamps = new HashMap<>();
+ measurementNames.forEach(measurementName ->
latestTimestamps.put(measurementName, getLatestTimestamp(
+ measurementName)));
+ return latestTimestamps;
+ }
+
+ private Long getLatestTimestamp(String measurementName) {
+ Map<String, String> queryParams = Map.of(
+ QP_START_DATE, "0",
+ QP_END_DATE, String.valueOf(System.currentTimeMillis()),
+ QP_LIMIT, "1",
+ QP_ORDER, "DESC",
+ QP_MISSING_VALUE_BEHAVIOUR, "empty"
+ );
+
+ try {
+ return getData(new ProvidedRestQueryParams(measurementName,
queryParams), true).getLastTimestamp();
+ } catch (RuntimeException e) {
+ return 0L;
+ }
+ }
}
diff --git
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
index fbb715be88..d7a00a00a7 100644
---
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
+++
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutor.java
@@ -37,6 +37,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.TreeMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
@@ -180,6 +181,79 @@ public class DataExplorerInfluxQueryExecutor extends
DataExplorerQueryExecutor<Q
}
}
+ public Map<String, Long> getLatestTimestamps(Map<String, String>
measurementFields) {
+ if (measurementFields.isEmpty()) {
+ return Map.of();
+ }
+
+ var query = makeLatestTimestampQuery(measurementFields);
+ try (final InfluxDB influxDB = InfluxClientProvider.getInfluxDBClient()) {
+ return parseLatestTimestampResult(influxDB.query(query,
TimeUnit.MILLISECONDS));
+ }
+ }
+
+ Query makeLatestTimestampQuery(Map<String, String> measurementFields) {
+ var query = measurementFields.entrySet()
+ .stream()
+ .collect(Collectors.groupingBy(Map.Entry::getValue, TreeMap::new,
+ Collectors.mapping(Map.Entry::getKey, Collectors.toList())))
+ .entrySet()
+ .stream()
+ .map(entry -> makeLastSelectorQuery(entry.getKey(), entry.getValue()))
+ .collect(Collectors.joining(";"));
+
+ return new Query(query, getDatabaseName());
+ }
+
+ private String makeLastSelectorQuery(String field,
+ List<String> measurements) {
+ return "SELECT LAST(\""
+ + field
+ + "\") FROM /"
+ + measurements.stream()
+ .map(this::escapeRegex)
+ .collect(Collectors.joining("|", "^(", ")$"))
+ + "/";
+ }
+
+ private String escapeRegex(String measurement) {
+ return measurement.replaceAll("([\\\\.\\[\\]{}()*+?^$|])", "\\\\$1");
+ }
+
+ private Map<String, Long> parseLatestTimestampResult(QueryResult
queryResult) {
+ Map<String, Long> latestTimestamps = new HashMap<>();
+ if (queryResult.getResults() != null) {
+ queryResult.getResults().forEach(result -> {
+ if (result.getSeries() != null) {
+ result.getSeries().forEach(series ->
parseLatestTimestampSeries(series, latestTimestamps));
+ }
+ });
+ }
+ return latestTimestamps;
+ }
+
+ private void parseLatestTimestampSeries(QueryResult.Series series,
+ Map<String, Long> latestTimestamps) {
+ var values = series.getValues();
+ if (values != null && !values.isEmpty() && !values.get(0).isEmpty()) {
+ latestTimestamps.put(series.getName(),
parseTimestamp(values.get(0).get(0)));
+ }
+ }
+
+ private Long parseTimestamp(Object timestamp) {
+ if (timestamp instanceof Number number) {
+ return number.longValue();
+ } else if (timestamp instanceof String timestampString) {
+ try {
+ return Long.parseLong(timestampString);
+ } catch (NumberFormatException e) {
+ return 0L;
+ }
+ } else {
+ return 0L;
+ }
+ }
+
@Override
public boolean deleteData(DataLakeMeasure measure) {
QueryResult queryResult = new DeleteDataQuery(measure).executeQuery();
diff --git
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
index eee432d3e2..ed2bea62fc 100644
---
a/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
+++
b/streampipes-data-explorer-influx/src/main/java/org/apache/streampipes/dataexplorer/influx/DataExplorerQueryManagementInflux.java
@@ -30,12 +30,25 @@ import
org.apache.streampipes.model.datalake.DataLakeMeasure;
import org.apache.streampipes.model.datalake.SpQueryResult;
import org.apache.streampipes.model.datalake.SpQueryStatus;
import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
+import org.apache.streampipes.model.schema.EventProperty;
+import org.apache.streampipes.model.schema.EventPropertyPrimitive;
import java.io.IOException;
import java.io.OutputStream;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
+import java.util.Objects;
+import java.util.Optional;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_END_DATE;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_LIMIT;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_MISSING_VALUE_BEHAVIOUR;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_ORDER;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_START_DATE;
public class DataExplorerQueryManagementInflux implements
IDataExplorerQueryManagement {
@@ -112,6 +125,81 @@ public class DataExplorerQueryManagementInflux implements
IDataExplorerQueryMana
return new DataExplorerInfluxQueryExecutor().getTagValues(measurementId,
fields);
}
+ @Override
+ public Map<String, Long> getLatestTimestamps(List<String> measurementNames) {
+ Map<String, Long> latestTimestamps = measurementNames.stream()
+ .collect(Collectors.toMap(
+ Function.identity(),
+ measurementName -> 0L,
+ (left, right) -> left
+ ));
+ var measurementFields = getLatestTimestampFields(measurementNames);
+
+ if (!measurementFields.isEmpty()) {
+ try {
+ var batchedLatestTimestamps = new
DataExplorerInfluxQueryExecutor().getLatestTimestamps(measurementFields);
+ latestTimestamps.putAll(batchedLatestTimestamps);
+ measurementFields.keySet()
+ .stream()
+ .filter(measurementName ->
!batchedLatestTimestamps.containsKey(measurementName))
+ .forEach(measurementName -> latestTimestamps.put(measurementName,
getLatestTimestampFallback(measurementName)));
+ } catch (RuntimeException e) {
+ measurementFields.keySet()
+ .forEach(measurementName ->
+ latestTimestamps.put(measurementName,
getLatestTimestampFallback(measurementName)));
+ }
+ }
+
+ measurementNames.stream()
+ .filter(measurementName ->
!measurementFields.containsKey(measurementName))
+ .forEach(measurementName -> latestTimestamps.put(measurementName,
getLatestTimestampFallback(measurementName)));
+
+ return latestTimestamps;
+ }
+
+ private Map<String, String> getLatestTimestampFields(List<String>
measurementNames) {
+ Map<String, String> measurementFields = new HashMap<>();
+ Map<String, DataLakeMeasure> measuresByName = getAllMeasurements()
+ .stream()
+ .collect(Collectors.toMap(DataLakeMeasure::getMeasureName,
Function.identity(), (left, right) -> left));
+
+ measurementNames.forEach(measurementName ->
findLatestTimestampField(measuresByName.get(measurementName))
+ .ifPresent(field -> measurementFields.put(measurementName, field)));
+
+ return measurementFields;
+ }
+
+ private Optional<String> findLatestTimestampField(DataLakeMeasure measure) {
+ if (measure == null || measure.getEventSchema() == null ||
measure.getEventSchema().getEventProperties() == null) {
+ return Optional.empty();
+ }
+
+ return measure.getEventSchema()
+ .getEventProperties()
+ .stream()
+ .filter(Objects::nonNull)
+ .filter(EventPropertyPrimitive.class::isInstance)
+ .map(EventProperty::getRuntimeName)
+ .filter(Objects::nonNull)
+ .findFirst();
+ }
+
+ private Long getLatestTimestampFallback(String measurementName) {
+ Map<String, String> queryParams = Map.of(
+ QP_START_DATE, "0",
+ QP_END_DATE, String.valueOf(System.currentTimeMillis()),
+ QP_LIMIT, "1",
+ QP_ORDER, "DESC",
+ QP_MISSING_VALUE_BEHAVIOUR, "empty"
+ );
+
+ try {
+ return getData(new ProvidedRestQueryParams(measurementName,
queryParams), true).getLastTimestamp();
+ } catch (RuntimeException e) {
+ return 0L;
+ }
+ }
+
private List<DataLakeMeasure> getAllMeasurements() {
return this.dataExplorerSchemaManagement.getAllMeasurements();
}
diff --git
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutorTest.java
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutorTest.java
new file mode 100644
index 0000000000..b8e5c3f8ef
--- /dev/null
+++
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/DataExplorerInfluxQueryExecutorTest.java
@@ -0,0 +1,41 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.dataexplorer.influx;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.LinkedHashMap;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+public class DataExplorerInfluxQueryExecutorTest {
+
+ @Test
+ public void makeLatestTimestampQueryGroupsMeasurementsByField() {
+ var measurementFields = new LinkedHashMap<String, String>();
+ measurementFields.put("measure-1", "value");
+ measurementFields.put("measure-2", "value");
+ measurementFields.put("measure-3", "temperature");
+
+ var query = new
DataExplorerInfluxQueryExecutor().makeLatestTimestampQuery(measurementFields);
+
+ assertEquals("SELECT LAST(\"temperature\") FROM /^(measure-3)$/;"
+ + "SELECT LAST(\"value\") FROM /^(measure-1|measure-2)$/",
query.getCommand());
+ }
+}
diff --git
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/InfluxLatestTimestampBenchmarkTest.java
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/InfluxLatestTimestampBenchmarkTest.java
new file mode 100644
index 0000000000..baedaf322d
--- /dev/null
+++
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/InfluxLatestTimestampBenchmarkTest.java
@@ -0,0 +1,331 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.dataexplorer.influx;
+
+import org.apache.streampipes.dataexplorer.influx.client.InfluxClientUtils;
+
+import org.influxdb.InfluxDB;
+import org.influxdb.InfluxDBFactory;
+import org.influxdb.dto.BatchPoints;
+import org.influxdb.dto.Point;
+import org.influxdb.dto.Query;
+import org.influxdb.dto.QueryResult;
+import org.junit.jupiter.api.Assumptions;
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Comparator;
+import java.util.List;
+import java.util.Locale;
+import java.util.concurrent.TimeUnit;
+import java.util.function.IntFunction;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+@Tag("benchmark")
+public class InfluxLatestTimestampBenchmarkTest {
+
+ private static final String ENABLED = "sp.influx.benchmark";
+ private static final String URL = "sp.influx.benchmark.url";
+ private static final String USERNAME = "sp.influx.benchmark.username";
+ private static final String PASSWORD = "sp.influx.benchmark.password";
+ private static final String TOKEN = "sp.influx.benchmark.token";
+ private static final String DATABASE = "sp.influx.benchmark.database";
+ private static final String RETENTION_POLICY =
"sp.influx.benchmark.retention-policy";
+ private static final String RECREATE_DATABASE =
"sp.influx.benchmark.recreate-database";
+ private static final String CLEANUP_MEASUREMENTS =
"sp.influx.benchmark.cleanup-measurements";
+ private static final String MEASUREMENTS =
"sp.influx.benchmark.measurements";
+ private static final String POINTS_PER_MEASUREMENT =
"sp.influx.benchmark.points-per-measurement";
+ private static final String FIELDS = "sp.influx.benchmark.fields";
+ private static final String BATCH_SIZE = "sp.influx.benchmark.batch-size";
+ private static final String WARMUP_RUNS = "sp.influx.benchmark.warmup-runs";
+ private static final String MEASUREMENT_RUNS =
"sp.influx.benchmark.measurement-runs";
+ private static final String ALLOW_NON_BENCHMARK_DATABASE =
"sp.influx.benchmark.allow-non-benchmark-database";
+
+ private static final long BASE_TIMESTAMP = 1700000000000L;
+
+ @Test
+ public void benchmarkLatestTimestampQueriesAgainstInfluxDb() {
+ Assumptions.assumeTrue(Boolean.getBoolean(ENABLED),
+ "Enable with -D" + ENABLED + "=true and provide a reachable InfluxDB
1.x instance");
+
+ var config = BenchmarkConfig.fromSystemProperties();
+ assertBenchmarkDatabase(config);
+ try (var influxDb = config.connect()) {
+ for (int measurementCount : config.measurementCounts()) {
+ if (config.recreateDatabase()) {
+ recreateDatabase(influxDb, config.database());
+ } else if (config.cleanupMeasurements()) {
+ influxDb.setDatabase(config.database());
+ dropBenchmarkMeasurements(influxDb, config.database());
+ } else {
+ influxDb.setDatabase(config.database());
+ }
+ loadData(influxDb, config, measurementCount);
+
+ printHeader();
+ runScenario(
+ influxDb,
+ config,
+ measurementCount,
+ "current_select_star",
+ index -> currentLatestTimestampQuery(index,
config.pointsPerMeasurement()),
+ measurementCount
+ );
+ runScenario(
+ influxDb,
+ config,
+ measurementCount,
+ "last_selector",
+ InfluxLatestTimestampBenchmarkTest::lastSelectorQuery,
+ measurementCount
+ );
+ runScenario(
+ influxDb,
+ config,
+ measurementCount,
+ "batched_last_selector",
+ ignored -> batchedLastSelectorQuery(),
+ 1
+ );
+ }
+ }
+ }
+
+ private static void printHeader() {
+
System.out.println("measurements,points_per_measurement,total_points,strategy,query_count,total_ms,mean_ms,p50_ms,"
+ + "p95_ms");
+ }
+
+ private static void runScenario(InfluxDB influxDb,
+ BenchmarkConfig config,
+ int measurementCount,
+ String strategy,
+ IntFunction<String> queryFactory,
+ int queryCount) {
+ for (int i = 0; i < config.warmupRuns(); i++) {
+ executeQueries(influxDb, config.database(), queryFactory, queryCount);
+ }
+
+ var durations = new ArrayList<Long>();
+ for (int i = 0; i < config.measurementRuns(); i++) {
+ durations.add(executeQueries(influxDb, config.database(), queryFactory,
queryCount));
+ }
+
+ durations.sort(Comparator.naturalOrder());
+ long totalNanos = durations.stream().mapToLong(Long::longValue).sum();
+ long totalPoints = (long) measurementCount * config.pointsPerMeasurement();
+ System.out.printf(Locale.ROOT,
+ "%d,%d,%d,%s,%d,%.3f,%.3f,%.3f,%.3f%n",
+ measurementCount,
+ config.pointsPerMeasurement(),
+ totalPoints,
+ strategy,
+ queryCount,
+ toMillis(totalNanos),
+ toMillis(totalNanos / config.measurementRuns()),
+ toMillis(percentile(durations, 0.50)),
+ toMillis(percentile(durations, 0.95)));
+ }
+
+ private static long executeQueries(InfluxDB influxDb,
+ String database,
+ IntFunction<String> queryFactory,
+ int queryCount) {
+ long startNanos = System.nanoTime();
+ for (int i = 0; i < queryCount; i++) {
+ var result = influxDb.query(new Query(queryFactory.apply(i), database),
TimeUnit.MILLISECONDS);
+ assertNoQueryError(result);
+ }
+ return System.nanoTime() - startNanos;
+ }
+
+ private static void assertNoQueryError(QueryResult result) {
+ assertFalse(result.hasError(), () -> "Influx query failed: " +
result.getError());
+ if (result.getResults() != null) {
+ result.getResults()
+ .forEach(r -> assertFalse(r.hasError(), () -> "Influx query failed:
" + r.getError()));
+ }
+ }
+
+ private static void assertBenchmarkDatabase(BenchmarkConfig config) {
+ assertTrue(
+ Boolean.getBoolean(ALLOW_NON_BENCHMARK_DATABASE) ||
config.database().toLowerCase(Locale.ROOT)
+ .contains("benchmark"),
+ "Refusing to drop and recreate database '"
+ + config.database()
+ + "'. Use a benchmark database name or set -D"
+ + ALLOW_NON_BENCHMARK_DATABASE
+ + "=true."
+ );
+ }
+
+ private static void recreateDatabase(InfluxDB influxDb,
+ String database) {
+ influxDb.query(new Query("DROP DATABASE \"" + database + "\"", ""));
+ influxDb.query(new Query("CREATE DATABASE \"" + database + "\"", ""));
+ influxDb.setDatabase(database);
+ }
+
+ private static void dropBenchmarkMeasurements(InfluxDB influxDb,
+ String database) {
+ assertNoQueryError(influxDb.query(new Query("DROP SERIES FROM
/^measure-[0-9]+$/", database)));
+ }
+
+ private static void loadData(InfluxDB influxDb,
+ BenchmarkConfig config,
+ int measurementCount) {
+ BatchPoints batch = newBatch(config.database(), config.retentionPolicy());
+ int pointsInBatch = 0;
+ for (int measurementIndex = 0; measurementIndex < measurementCount;
measurementIndex++) {
+ for (int pointIndex = 0; pointIndex < config.pointsPerMeasurement();
pointIndex++) {
+ batch.point(makePoint(measurementIndex, pointIndex, config.fields()));
+ pointsInBatch++;
+ if (pointsInBatch == config.batchSize()) {
+ influxDb.write(batch);
+ batch = newBatch(config.database(), config.retentionPolicy());
+ pointsInBatch = 0;
+ }
+ }
+ }
+
+ if (pointsInBatch > 0) {
+ influxDb.write(batch);
+ }
+ }
+
+ private static BatchPoints newBatch(String database,
+ String retentionPolicy) {
+ var builder = BatchPoints.database(database)
+ .consistency(InfluxDB.ConsistencyLevel.ALL);
+ if (retentionPolicy != null && !retentionPolicy.isBlank()) {
+ builder.retentionPolicy(retentionPolicy);
+ }
+ return builder.build();
+ }
+
+ private static Point makePoint(int measurementIndex,
+ int pointIndex,
+ int fields) {
+ var point = Point.measurement(measurementName(measurementIndex))
+ .time(BASE_TIMESTAMP + pointIndex * 1000L, TimeUnit.MILLISECONDS)
+ .tag("device", "device-" + measurementIndex)
+ .tag("site", "site-" + measurementIndex % 10)
+ .addField("value", pointIndex);
+
+ for (int fieldIndex = 1; fieldIndex < fields; fieldIndex++) {
+ point.addField("value_" + fieldIndex, pointIndex + fieldIndex);
+ }
+
+ return point.build();
+ }
+
+ private static String currentLatestTimestampQuery(int measurementIndex,
+ int pointsPerMeasurement) {
+ long endTimestamp = (BASE_TIMESTAMP + pointsPerMeasurement * 1000L +
1000L) * 1000000L;
+ return "SELECT * FROM \"" + measurementName(measurementIndex)
+ + "\" WHERE time > 0 AND time < "
+ + endTimestamp
+ + " ORDER BY time DESC LIMIT 1";
+ }
+
+ private static String lastSelectorQuery(int measurementIndex) {
+ return "SELECT LAST(\"value\") FROM \"" +
measurementName(measurementIndex) + "\"";
+ }
+
+ private static String batchedLastSelectorQuery() {
+ return "SELECT LAST(\"value\") FROM /^measure-[0-9]+$/";
+ }
+
+ private static String measurementName(int measurementIndex) {
+ return "measure-" + measurementIndex;
+ }
+
+ private static long percentile(List<Long> durations,
+ double percentile) {
+ int index = (int) Math.ceil(percentile * durations.size()) - 1;
+ return durations.get(Math.max(index, 0));
+ }
+
+ private static double toMillis(long durationNanos) {
+ return durationNanos / 1_000_000.0;
+ }
+
+ private record BenchmarkConfig(String url,
+ String username,
+ String password,
+ String token,
+ String database,
+ String retentionPolicy,
+ boolean recreateDatabase,
+ boolean cleanupMeasurements,
+ List<Integer> measurementCounts,
+ int pointsPerMeasurement,
+ int fields,
+ int batchSize,
+ int warmupRuns,
+ int measurementRuns) {
+
+ static BenchmarkConfig fromSystemProperties() {
+ return new BenchmarkConfig(
+ System.getProperty(URL, "http://localhost:8086"),
+ System.getProperty(USERNAME, ""),
+ System.getProperty(PASSWORD, ""),
+ System.getProperty(TOKEN, ""),
+ System.getProperty(DATABASE, "streampipes_benchmark"),
+ System.getProperty(RETENTION_POLICY, ""),
+ Boolean.parseBoolean(System.getProperty(RECREATE_DATABASE, "true")),
+ Boolean.parseBoolean(System.getProperty(CLEANUP_MEASUREMENTS,
"true")),
+ getIntList(MEASUREMENTS, "10,100"),
+ getInt(POINTS_PER_MEASUREMENT, 1000),
+ getInt(FIELDS, 3),
+ getInt(BATCH_SIZE, 5000),
+ getInt(WARMUP_RUNS, 2),
+ getInt(MEASUREMENT_RUNS, 5)
+ );
+ }
+
+ InfluxDB connect() {
+ if (token != null && !token.isBlank()) {
+ return InfluxDBFactory.connect(url,
InfluxClientUtils.getHttpClientBuilder(token));
+ } else if (username == null || username.isBlank()) {
+ return InfluxDBFactory.connect(url);
+ } else {
+ return InfluxDBFactory.connect(url, username, password);
+ }
+ }
+
+ private static int getInt(String property,
+ int defaultValue) {
+ return Integer.parseInt(System.getProperty(property,
String.valueOf(defaultValue)));
+ }
+
+ private static List<Integer> getIntList(String property,
+ String defaultValue) {
+ return Arrays.stream(System.getProperty(property,
defaultValue).split(","))
+ .map(String::trim)
+ .filter(value -> !value.isBlank())
+ .map(Integer::parseInt)
+ .toList();
+ }
+ }
+}
diff --git
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/SelectQueryParamsTest.java
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/SelectQueryParamsTest.java
index 6513f2a7c3..62188caac1 100644
---
a/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/SelectQueryParamsTest.java
+++
b/streampipes-data-explorer-influx/src/test/java/org/apache/streampipes/dataexplorer/influx/SelectQueryParamsTest.java
@@ -21,17 +21,44 @@ package org.apache.streampipes.dataexplorer.influx;
import
org.apache.streampipes.dataexplorer.influx.utils.ProvidedQueryParameterBuilder;
import
org.apache.streampipes.dataexplorer.param.ProvidedRestQueryParamConverter;
import org.apache.streampipes.dataexplorer.param.SelectQueryParams;
+import org.apache.streampipes.model.datalake.param.ProvidedRestQueryParams;
import org.junit.jupiter.api.Test;
import java.util.Arrays;
import java.util.List;
+import java.util.Map;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_END_DATE;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_LIMIT;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_MISSING_VALUE_BEHAVIOUR;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_ORDER;
+import static
org.apache.streampipes.model.datalake.param.SupportedRestQueryParams.QP_START_DATE;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
public class SelectQueryParamsTest {
+ @Test
+ public void testLatestEventTimestampQuery() {
+ var params = new ProvidedRestQueryParams(
+ "abc",
+ Map.of(
+ QP_START_DATE, "0",
+ QP_END_DATE, "100",
+ QP_LIMIT, "1",
+ QP_ORDER, "DESC",
+ QP_MISSING_VALUE_BEHAVIOUR, "empty"
+ )
+ );
+
+ SelectQueryParams qp =
ProvidedRestQueryParamConverter.getSelectQueryParams(params);
+
+ String query =
qp.toQuery(DataLakeInfluxQueryBuilder.create("abc")).getCommand();
+
+ assertEquals("SELECT * FROM \"abc\" WHERE (time < 100000000 AND time > 0)
ORDER BY time DESC LIMIT 1;", query);
+ }
+
@Test
public void testWildcardTimeBoundQuery() {
var params = ProvidedQueryParameterBuilder.create("abc")
diff --git
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
index e38489b109..61cf4ed204 100644
---
a/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
+++
b/streampipes-rest/src/main/java/org/apache/streampipes/rest/impl/datalake/DataLakeResource.java
@@ -274,11 +274,7 @@ public class DataLakeResource extends
AbstractDataLakeResource {
);
}
- Map<String, Long> latestEvents = distinctMeasurementNames.stream()
- .collect(Collectors.toMap(
- measurementName -> measurementName,
- this::getLatestEvent
- ));
+ Map<String, Long> latestEvents =
this.dataExplorerQueryManagement.getLatestTimestamps(distinctMeasurementNames);
return ok(latestEvents);
}
@@ -436,25 +432,6 @@ public class DataLakeResource extends
AbstractDataLakeResource {
return new ProvidedRestQueryParams(measurementId, queryParamMap);
}
- private Long getLatestEvent(String measurementName) {
- Map<String, String> queryParams = Map.of(
- QP_START_DATE, "0",
- QP_END_DATE, String.valueOf(System.currentTimeMillis()),
- QP_LIMIT, "1",
- QP_ORDER, "DESC",
- QP_MISSING_VALUE_BEHAVIOUR, "empty"
- );
-
- try {
- return this.dataExplorerQueryManagement
- .getData(new ProvidedRestQueryParams(measurementName, queryParams),
true)
- .getLastTimestamp();
- } catch (RuntimeException e) {
- LOG.warn("Could not get latest event for measurement {}",
measurementName, e);
- return 0L;
- }
- }
-
// Checks if the parameter for missing value behaviour is set
private boolean isIgnoreMissingValues(String missingValueBehaviour) {
return "ignore".equals(missingValueBehaviour);
diff --git
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeLatestEventsBenchmarkTest.java
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeLatestEventsBenchmarkTest.java
new file mode 100644
index 0000000000..23bb3170da
--- /dev/null
+++
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeLatestEventsBenchmarkTest.java
@@ -0,0 +1,108 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.rest.impl.datalake;
+
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
+
+import org.junit.jupiter.api.Tag;
+import org.junit.jupiter.api.Test;
+
+import java.lang.reflect.Field;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Locale;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.CALLS_REAL_METHODS;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.when;
+
+@Tag("benchmark")
+class DataLakeLatestEventsBenchmarkTest {
+
+ private static final List<Integer> MEASUREMENT_COUNTS = List.of(10, 100,
1000);
+ private static final int WARMUP_RUNS = 2;
+ private static final int MEASUREMENT_RUNS = 5;
+
+ @Test
+ void benchmarkLatestEventsFanOut() throws Exception {
+
System.out.println("measurement_count,mean_duration_ms,query_management_calls");
+ for (int measurementCount : MEASUREMENT_COUNTS) {
+ for (int i = 0; i < WARMUP_RUNS; i++) {
+ runLatestEvents(measurementCount);
+ }
+
+ long totalDurationNanos = 0L;
+ for (int i = 0; i < MEASUREMENT_RUNS; i++) {
+ totalDurationNanos += runLatestEvents(measurementCount);
+ }
+
+ var meanDurationMillis =
TimeUnit.NANOSECONDS.toMicros(totalDurationNanos / MEASUREMENT_RUNS) / 1000.0;
+ System.out.printf(Locale.ROOT, "%d,%.3f,%d%n", measurementCount,
meanDurationMillis, 1);
+ }
+ }
+
+ private static long runLatestEvents(int measurementCount) throws Exception {
+ var queryManagement = mock(IDataExplorerQueryManagement.class);
+ var measurementNames = measurementNames(measurementCount);
+
when(queryManagement.getLatestTimestamps(measurementNames)).thenReturn(latestTimestamps(measurementNames));
+ var resource = dataLakeResource(queryManagement);
+
+ long startNanos = System.nanoTime();
+ var response = resource.getLatestEvents(measurementNames);
+ long durationNanos = System.nanoTime() - startNanos;
+
+ assertEquals(measurementCount, ((java.util.Map<?, ?>)
response.getBody()).size());
+ verify(queryManagement).getLatestTimestamps(measurementNames);
+ verifyNoMoreInteractions(queryManagement);
+ return durationNanos;
+ }
+
+ private static List<String> measurementNames(int measurementCount) {
+ return java.util.stream.IntStream.range(0, measurementCount)
+ .mapToObj(i -> "measure-" + i)
+ .toList();
+ }
+
+ private static java.util.Map<String, Long> latestTimestamps(List<String>
measurementNames) {
+ var latestTimestamps = new HashMap<String, Long>();
+ measurementNames.forEach(measurementName ->
latestTimestamps.put(measurementName, (long) measurementName.length()));
+ return latestTimestamps;
+ }
+
+ private static DataLakeResource
dataLakeResource(IDataExplorerQueryManagement queryManagement) throws Exception
{
+ var resource = mock(DataLakeResource.class, CALLS_REAL_METHODS);
+ doReturn(true).when(resource).checkPermissionByName(any(), eq("READ"));
+ setQueryManagement(resource, queryManagement);
+ return resource;
+ }
+
+ private static void setQueryManagement(DataLakeResource resource,
+ IDataExplorerQueryManagement
queryManagement) throws Exception {
+ Field field =
DataLakeResource.class.getDeclaredField("dataExplorerQueryManagement");
+ field.setAccessible(true);
+ field.set(resource, queryManagement);
+ }
+}
diff --git
a/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
new file mode 100644
index 0000000000..aea4c1e753
--- /dev/null
+++
b/streampipes-rest/src/test/java/org/apache/streampipes/rest/impl/datalake/DataLakeResourceTest.java
@@ -0,0 +1,84 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.rest.impl.datalake;
+
+import org.apache.streampipes.dataexplorer.api.IDataExplorerQueryManagement;
+
+import org.junit.jupiter.api.Test;
+import org.springframework.http.HttpStatus;
+
+import java.lang.reflect.Field;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.CALLS_REAL_METHODS;
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class DataLakeResourceTest {
+
+ @Test
+ void getLatestEventsRequestsLatestTimestampsForDistinctMeasurements() throws
Exception {
+ var queryManagement = mock(IDataExplorerQueryManagement.class);
+ when(queryManagement.getLatestTimestamps(List.of("a", "bb", "broken")))
+ .thenReturn(Map.of("a", 1L, "bb", 2L, "broken", 0L));
+ var resource = dataLakeResource(queryManagement, true);
+
+ var response = resource.getLatestEvents(List.of("a", "bb", "a", "broken"));
+
+ assertEquals(HttpStatus.OK, response.getStatusCode());
+ assertEquals(Map.of("a", 1L, "bb", 2L, "broken", 0L), response.getBody());
+ verify(queryManagement).getLatestTimestamps(List.of("a", "bb", "broken"));
+ verify(queryManagement, times(0)).getData(any(), eq(true));
+ }
+
+ @Test
+ void getLatestEventsRejectsUnauthorizedMeasurementBeforeQuerying() throws
Exception {
+ var queryManagement = mock(IDataExplorerQueryManagement.class);
+ var resource = dataLakeResource(queryManagement, false);
+
+ var response = resource.getLatestEvents(List.of("a"));
+
+ assertEquals(HttpStatus.BAD_REQUEST, response.getStatusCode());
+ assertEquals("No read permission for measurement a", response.getBody());
+ verify(queryManagement, times(0)).getData(any(), eq(true));
+ verify(queryManagement, times(0)).getLatestTimestamps(any());
+ }
+
+ private static DataLakeResource
dataLakeResource(IDataExplorerQueryManagement queryManagement,
+ boolean canRead) throws
Exception {
+ var resource = mock(DataLakeResource.class, CALLS_REAL_METHODS);
+ doReturn(canRead).when(resource).checkPermissionByName(any(), eq("READ"));
+ setQueryManagement(resource, queryManagement);
+ return resource;
+ }
+
+ private static void setQueryManagement(DataLakeResource resource,
+ IDataExplorerQueryManagement
queryManagement) throws Exception {
+ Field field =
DataLakeResource.class.getDeclaredField("dataExplorerQueryManagement");
+ field.setAccessible(true);
+ field.set(resource, queryManagement);
+ }
+}