Enable database result streaming

This commit is contained in:
Anton Tananaev
2026-07-10 19:37:55 -07:00
parent db5bb1ce12
commit 83cb935e4c
4 changed files with 33 additions and 7 deletions

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2019 - 2025 Anton Tananaev (anton@traccar.org)
* Copyright 2019 - 2026 Anton Tananaev (anton@traccar.org)
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
@@ -579,6 +579,15 @@ public final class Keys {
List.of(KeyType.CONFIG),
20);
/**
* Number of rows fetched per round trip for streamed queries (position history and exports). On PostgreSQL this
* enables a server-side cursor so results are not fully buffered in memory.
*/
public static final ConfigKey<Integer> DATABASE_STREAM_FETCH_SIZE = new IntegerConfigKey(
"database.streamFetchSize",
List.of(KeyType.CONFIG),
1000);
/**
* SQL query to check connection status. Default value is 'SELECT 1'. For Oracle database you can use
* 'SELECT 1 FROM DUAL'.

View File

@@ -86,7 +86,7 @@ public class DatabaseStorage extends Storage {
for (int index = 0; index < values.size(); index++) {
builder.setValue(index, values.get(index));
}
Stream<T> stream = builder.executeQueryStreamed(clazz);
Stream<T> stream = builder.executeQueryStreamed(clazz, databaseType);
builder = null;
return stream;
} catch (SQLException e) {

View File

@@ -59,6 +59,8 @@ public final class QueryBuilder implements AutoCloseable {
private final String query;
private final boolean returnGeneratedKeys;
private boolean streamingTransaction;
private QueryBuilder(
Config config, DataSource dataSource, ObjectMapper objectMapper,
String query, boolean returnGeneratedKeys) throws SQLException {
@@ -234,11 +236,18 @@ public final class QueryBuilder implements AutoCloseable {
}
}
public <T> Stream<T> executeQueryStreamed(Class<T> clazz) throws SQLException {
public <T> Stream<T> executeQueryStreamed(Class<T> clazz, String databaseType) throws SQLException {
ResultSet resultSet = null;
try {
logQuery();
connection.setAutoCommit(false);
streamingTransaction = true;
statement.setFetchSize(switch (databaseType) {
case "MySQL", "MariaDB" -> Integer.MIN_VALUE;
default -> config.getInteger(Keys.DATABASE_STREAM_FETCH_SIZE);
});
resultSet = statement.executeQuery();
ResultSetMetaData resultMetaData = resultSet.getMetaData();
@@ -311,7 +320,15 @@ public final class QueryBuilder implements AutoCloseable {
try {
statement.close();
} finally {
connection.close();
try {
if (streamingTransaction) {
streamingTransaction = false;
connection.rollback();
connection.setAutoCommit(true);
}
} finally {
connection.close();
}
}
}

View File

@@ -96,7 +96,7 @@ public class QueryBuilderTest {
try (QueryBuilder query = QueryBuilder.create(config, dataSource, objectMapper,
"SELECT * FROM test_entity");
Stream<TestEntity> stream = query.executeQueryStreamed(TestEntity.class)) {
Stream<TestEntity> stream = query.executeQueryStreamed(TestEntity.class, "H2")) {
List<TestEntity> results = stream.toList();
assertEquals(1, results.size());
TestEntity entity = results.get(0);
@@ -132,7 +132,7 @@ public class QueryBuilderTest {
try (QueryBuilder query = QueryBuilder.create(config, dataSource, objectMapper,
"SELECT * FROM test_entity");
Stream<TestEntity> stream = query.executeQueryStreamed(TestEntity.class)) {
Stream<TestEntity> stream = query.executeQueryStreamed(TestEntity.class, "H2")) {
List<TestEntity> results = stream.toList();
assertEquals(1, results.size());
TestEntity loaded = results.get(0);
@@ -161,7 +161,7 @@ public class QueryBuilderTest {
try (QueryBuilder query = QueryBuilder.create(config, dataSource, objectMapper,
"SELECT * FROM test_entity ORDER BY count");
Stream<TestEntity> stream = query.executeQueryStreamed(TestEntity.class)) {
Stream<TestEntity> stream = query.executeQueryStreamed(TestEntity.class, "H2")) {
List<TestEntity> results = stream.toList();
assertEquals(3, results.size());
assertEquals("row0", results.get(0).getName());