Skip to content
JackSparrow414
Go back

JDBC Streaming Reads and Correct Streaming Processing, with Examples

Table of contents

Open Table of contents

Article body

This article uses PostgreSQL. For MySQL, see the official MySQL JDBC documentation.

Enabling Streaming Reads in PostgreSQL

Official documentation

Why Does OOM Still Occur After Enabling Streaming Reads?

Most developers know that some JDBC drivers return an entire query result to the application at once, which can easily cause an out-of-memory error (OOM). The PostgreSQL JDBC documentation explains:

By default the driver collects all the results for the query at once. This can be inconvenient for large data sets so the JDBC driver provides a means of basing a ResultSet on a database cursor and only fetching a small number of rows.

Most also know that streaming reads avoid loading everything at once. It takes two lines of code; setting fetchSize is the key:

connection.setAutoCommit(false);
statement.setFetchSize(50);

What many explanations omit is how to return or process the data afterward. Without thinking it through enough, I wrote the following service code:

public class JdbcStreamService {

	public List<User> getUserList() {
		List<User> result = new ArrayList();
		String url = "jdbc:postgresql://localhost:5432/postgres?user=postgres&password=12345&?currentSchema=public";
        Connection conn = DriverManager.getConnection(url);
        // Disable auto-commit
        conn.setAutoCommit(false);
        Statement st = conn.createStatement();
        // Enable streaming reads
        st.setFetchSize(50);
        ResultSet rs = st.executeQuery("SELECT * FROM public.user");
        try {
            while (rs.next()) {
                int userId = rs.getInt(1);
                String password = rs.getString(2);
                String roles = rs.getString(3);
                String introduction = rs.getString(4);
                User user = new User();
                user.setUserId(userId);
                user.setPassword(password);
                user.setRoles(roles);
                user.setIntroduction(introduction);
                result.add(user);
            }
        } finally {
            rs.close();
            st.close();
            conn.close();
        }
        return result;
	}
}

I was delighted: streaming reads were enabled, so surely I no longer needed to worry about OOM. I put every result into a List and returned it to the Controller. Unfortunately, after starting the application and calling the endpoint, this approach still caused OOM.

Why?

Although the database results are read in batches, the User objects produced by each batch are never consumed. They remain in the List, and the JVM cannot garbage-collect them.

Thinking it over, I realized that I had not understood streams deeply enough. A data stream is like water: it has a source and a destination, and it must be consumed. Consumption can mean placing data in some storage or discarding it. That storage might be memory, a file on disk, or another endpoint reached over the network. Once I understood this, I knew how to write the code correctly. I strongly recommend this Stack Overflow explanation of streams.

Each batch fetched from the database is converted from ResultSet rows into entities. We need to consume that entity stream while iterating over the ResultSet. The entities already reside in JVM memory, so to avoid OOM, write them to a file or return them over the network to the requester—usually a browser.

Streaming Reads into a File

Suppose we need to generate a report from a database query with a very large result set. This is a common requirement.

The following code queries through JPA and sets fetchSize. Each iteration reads a User, writes it to the file, and removes it from the persistence context.

Note: you must detach the entity; otherwise too many entities in the persistence context can also cause OOM. Documentation

if you are processing a huge number of objects and need to manage memory efficiently, the evict() method can be used to remove the object and its collections from the first-level cache

public boolean streamJdbcResultToFile(String fileName) {
        EntityManager entityManager = JPAUtil.acquireEntityManager();
        jakarta.persistence.Query query = entityManager.createQuery(JPQL, User.class);
        // Convert the JPA Query to a Hibernate Query
        query.unwrap(Query.class);
        query.setHint(AvailableHints.HINT_FETCH_SIZE, FETCH_SIZE);
        try (FileWriter fileWriter = getFileWriter(fileName)) {
            // Use JPA getResultStream
            query.getResultStream().forEach(object -> {
                User user = (User) object;
                writeToFile(fileWriter, user);
                // Must detach; otherwise too many entities in the persistence context can cause OOM
                entityManager.detach(user);
            });
        }finally {
            entityManager.close();
        }
        return true;
}

Streaming Reads into a JSON Response

Here we return a large number of JSON objects to the frontend as a stream.

The JSON returned to the frontend has a structure like this:

{
  "users": [
    {
      "userId": 1,
      "password": "KRRLAZg0IzxPY232lJefu9l6Hts8HO1cTfLIF38jrqPWNjIT78nVYlNrCsOl",
      "roles": "a5cdt8YMmXdxc5jkwFKxJBvkCZMR25ljGvVr79h33R0rB1SQKoIm6AllSGVL5Xk119phqMPKvYSPvxkXkc9W0PClhAybnPNGKm9jGey6P8IuisUNP5xvDZpKuPj00kyQ9lKSU6zr5qJN1i5U0dhOqAqPUOqpuluZNdwtDuVkaFI8sqFhCdYO6bUtSMCbiuyAOzFkn05t",
      "introduction": "8m1IXJOmixm2joDScCW2LVZwJtsNdBuG3NUzlMCMjtlYnvMJ6SzEkxRATmbq4mcb7WQ1NPCWPDjKsGfAnZpewihL7Ih95IxpGzcN8vq58m9DRZjiQyIhS7DrH60chEUuEV2qaf80hxq7p3P8gmPRXTtidt9lVOT7fhj0hN0SMqJtnyXoaQ6WFOVXehGkQyZfjWAJJiyu"
    }
  ]
}

This uses a RESTEasy response stream. For a plain Servlet application, replace it with ServletOutputStream.

public StreamingOutput streamJdbcResultToResponse() {
     StreamingOutput stream = output -> {
        EntityManager entityManager = JPAUtil.acquireEntityManager();
        Query<User> query = entityManager.createQuery(JPQL, User.class).unwrap(Query.class);
        query.setFetchSize(FETCH_SIZE);
        // Both streams are closed automatically
        try(ScrollableResults<User> scrollableResults = query.scroll();
        // Use JsonGenerator
            JsonGenerator jsonGenerator = OBJECT_MAPPER.getFactory().createGenerator(output);) {
            // Start the JSON object
            jsonGenerator.writeStartObject();
            // JSON object key
            jsonGenerator.writeArrayFieldStart("users");

            while (scrollableResults.next()) {
                User user = scrollableResults.get();
                jsonGenerator.writeObject(user);
                entityManager.detach(user);
            }
            // Finish
            jsonGenerator.writeEndArray();
            jsonGenerator.writeEndObject();
        }
    };
    return stream;
}

Testing

For testing, enable or comment out the fetchSize code and configure a small JVM heap in IDEA, for example:

-Xms50m -Xmx50m

Observe when OOM occurs and compare streaming reads plus streaming processing under a small heap, whether writing a file or writing a response to the frontend.

Other Helpful Articles and Code

The examples here do not use Spring. For streaming results to the frontend in Spring, see the following project or search Stack Overflow:

Example Code

The source includes a complete PostgreSQL docker-compose file and configuration for generating test data. To try the full example, install Docker, JDK 11, and Tomcat 10 locally. The README explains how to run and test it.


Share this post:

Previous Post
Kafka (Part 6): Oracle-to-PostgreSQL CDC with Kafka Connect and Debezium, and Cache Consistency
Next Post
Building Elastic Stack from the Official Documentation: A Three-Node Elasticsearch Cluster, Kibana, Filebeat, Metricbeat, and Migration Without Downtime

Comments

Questions, corrections, and experiences are welcome. Sign in with GitHub to comment; both language versions share this discussion.

Comments are available on the live site only.