← Back to the section

Let's say you already know why ClickHouse is useful: fast analytics, aggregations over millions of rows that PostgreSQL takes tens of seconds to compute. Now the practical question: how do you connect it to a Spring Boot service, how does data get in, and how do you read it back out?

the consumer buffers events and inserts them into ClickHouse as one batch order-events e1 e2 e3 e4 e5 e6 reads consumer buffer e1 e2 e3 e4 e5 e6 insert insert duplicate table order_events part 1e1, e2, e3 part 2e4, e5, e6 e4e5e6crash before the offset commit:the consumer re-read the same batch block hash already seen — no third part

Single inserts shatter the table into tiny parts, and the engine cannot merge them fast enough. The consumer buffers events and sends one batch — that makes one part. If the same batch arrives again after a crash, a replicated table recognises it by the block hash and does not insert it twice.

Connecting: a second DataSource alongside PostgreSQL

Most services already have PostgreSQL as their primary database. ClickHouse is a second store that lives alongside it, not instead of it. Spring Boot lets you keep several DataSources at once.

The official driver is com.clickhouse:clickhouse-jdbc. It works over the ClickHouse HTTP interface (port 8123), not the native protocol. That's convenient: firewall rules are simpler, and no special client is needed.

Configuration:

@Configuration
public class ClickHouseConfig {

    @Bean(defaultCandidate = false)
    @ConfigurationProperties("app.clickhouse")
    DataSourceProperties clickHouseDataSourceProperties() {
        return new DataSourceProperties();
    }

    @Bean(defaultCandidate = false)
    DataSource clickHouseDataSource() {
        var hikari = clickHouseDataSourceProperties()
            .initializeDataSourceBuilder()
            .type(HikariDataSource.class)
            .build();
        hikari.setMaximumPoolSize(5);
        return hikari;
    }

    @Bean
    JdbcClient clickHouseJdbcClient(@Qualifier("clickHouseDataSource") DataSource dataSource) {
        return JdbcClient.create(dataSource);
    }

    @Bean
    JdbcTemplate clickHouseJdbcTemplate(@Qualifier("clickHouseDataSource") DataSource dataSource) {
        return new JdbcTemplate(dataSource);
    }
}
app:
  clickhouse:
    url: jdbc:clickhouse://clickhouse.internal:8123/analytics
    username: app_analytics
    password: ${CLICKHOUSE_PASSWORD}

A few details that matter here:

  • The pool is small — at most 5 connections. Analytical queries are heavy and run rarely; 50 parallel aggregations can bring the server down.
  • @Bean(defaultCandidate = false) (Spring Boot 3.4+) is mandatory: without it auto-configuration skips the primary DataSource. With it, PostgreSQL stays as it is — @Transactional and jOOQ keep working with it, never noticing ClickHouse.
  • ClickHouse has no transactions, so the new client needs no transactional wrapping at all.

Writing: why not to write directly from the service

The first idea is usually this: after saving an order, add a line clickHouse.insert(...). It looks simple — but it's a double write to two different systems with no guarantees. If ClickHouse is unavailable, the business operation breaks. If the app crashes between the two writes, data is lost or drifts out of sync. On top of that, single INSERTs into ClickHouse quickly lead to the TOO_MANY_PARTS error — the engine can't merge the small data chunks fast enough.

The reliable path looks different:

PostgreSQL (outbox in the same transaction) → Kafka → batch consumer → ClickHouse

The service writes only to PostgreSQL. The event goes to Kafka through the outbox. A separate consumer accumulates events and inserts them in batches. Each step is independent: if ClickHouse goes down, Kafka simply waits. When it comes back, the consumer resumes from where it left off.

Batch consumer: accumulate and insert

@Component
@RequiredArgsConstructor
public class OrderEventsClickHouseSink {

    private final JdbcTemplate clickHouseJdbcTemplate;

    @KafkaListener(topics = "order-events", batch = "true",
                   containerFactory = "batchContainerFactory")
    public void consume(List<ConsumerRecord<String, OrderEventPayload>> records,
                        Acknowledgment ack) {
        insertBatch(records.stream().map(r -> r.value()).toList());
        ack.acknowledge();
    }

    private void insertBatch(List<OrderEventPayload> events) {
        var sql = """
            INSERT INTO order_events
                (event_id, event_time, event_type, region, customer_id, order_id, amount)
            VALUES (?, ?, ?, ?, ?, ?, ?)
            """;
        clickHouseJdbcTemplate.batchUpdate(sql, events.stream().map(this::toRow).toList());
    }
}

The batch insert goes through JdbcTemplate: JdbcClient has no batch mode, and the Spring documentation points to JdbcTemplate for it.

What you tune on the Kafka side: max.poll.records in the thousands (batch size), fetch.max.wait.ms — hundreds of milliseconds (so you don't hit ClickHouse too often). Commit the offset only after a successful insert. If something goes wrong, the consumer re-reads the same batch.

Async insert as a simpler option

If the data flow is small, you can enable async insert on the ClickHouse side: the server itself accumulates small inserts in a buffer and flushes them as a batch. It's turned on via async_insert=1 in the connection settings or in the query itself — the application stops having to worry about accumulation.

The downside hangs on a single setting. By default (wait_for_async_insert=1) the server answers "accepted" only once the buffer has really been written to disk — acknowledged inserts are not lost, but each one takes a little longer. Turn the wait off for speed and the answer comes back at once, and whatever is still in memory when the server dies is gone. For metrics and counters that's acceptable, for money-related events it isn't — there you still need the Kafka pipeline.

Kafka engine inside ClickHouse

A third option drops the Java consumer entirely: ClickHouse reads from Kafka itself via ENGINE = Kafka plus a materialized view, and data lands in the table directly.

The price: all the logic (mapping, retries, alerts) moves into ClickHouse DDL. The Java team doesn't see it and doesn't control it. This option works well when a separate team owns the analytics layer. If the service is responsible for its own data, an explicit Java consumer is clearer.

Idempotency: what to do about repeated inserts

Kafka delivers events at least once. After a failure the consumer re-reads part of the batch, and the same rows get inserted again. There are two defences.

Block-level deduplication. Replicated tables remember the hashes of the last inserted batches. If the same batch arrives again, ClickHouse simply ignores it. It works exactly for "crashed between the INSERT and the offset commit, then rebuilt the same batch from the same events".

ReplacingMergeTree by event id. The ReplacingMergeTree ORDER BY event_id engine collapses rows with the same key during background merges, keeping the latest version. Duplicates will land in the table but disappear on merge. Before the merge, uniq(event_id) will count correctly, count() won't.

Both defences fit into a small program: the table is a list of rows, the engine remembers block hashes.

live example

import java.util.ArrayList;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;

public class BatchSinkDemo {
    record Event(String id, String status) {
        @Override
        public String toString() {
            return id + ":" + status;
        }
    }

    static final List<Event> table = new ArrayList<>();
    static final Set<Integer> insertedBlocks = new HashSet<>();

    public static void main(String[] args) {
        List<Event> block = List.of(new Event("e1", "NEW"), new Event("e2", "PAID"));
        insert(block);
        insert(block);
        insert(List.of(new Event("e2", "SHIPPED")));

        System.out.println("rows in table:  " + table.size() + " " + table);
        System.out.println("uniq(event_id):  " + merged().size());
        System.out.println("after merge:     " + merged());
    }

    static void insert(List<Event> block) {
        if (!insertedBlocks.add(block.hashCode())) {
            System.out.println("batch " + block + " already inserted — skipped");
            return;
        }
        table.addAll(block);
        System.out.println("inserted part " + block);
    }

    static Map<String, String> merged() {
        Map<String, String> latest = new LinkedHashMap<>();
        for (Event event : table) {
            latest.put(event.id(), event.status());
        }
        return latest;
    }
}
Run

Running examples is part of paid access. There the same code runs inside the article: editor, run and check next to the paragraph. Three free days →

The repeated batch was rejected by its hash; the duplicate from a different batch survived until the merge.

Financial reports usually need both; for product analytics ReplacingMergeTree is often enough.

Reading: a separate repository for analytics

Analytical queries are separated from the rest of the code into a dedicated repository with explicit DTOs. This is the familiar read-side pattern: one layer writes, another reads.

@Repository
@RequiredArgsConstructor
public class RevenueViewRepository {

    private final JdbcClient clickHouseJdbcClient;

    public List<RevenueByRegionRow> revenueByRegion(LocalDate from, LocalDate to) {
        return clickHouseJdbcClient.sql("""
                SELECT region, sumMerge(revenue) AS revenue, uniqMerge(orders) AS orders
                FROM revenue_by_region_daily
                WHERE day BETWEEN ? AND ?
                GROUP BY region
                ORDER BY revenue DESC
                """)
            .params(from, to)
            .query((rs, i) -> new RevenueByRegionRow(
                rs.getString("region"),
                rs.getBigDecimal("revenue"),
                rs.getLong("orders")))
            .list();
    }
}

public record RevenueByRegionRow(String region, BigDecimal revenue, long orders) {}

A few hygiene rules:

  • The app_analytics user for the reading service is read-only. A different pipeline user writes to ClickHouse.
  • Set max_execution_time and max_memory_usage on the user profile: a random heavy query must not take down the cluster.
  • ClickHouse lags behind PostgreSQL by seconds to minutes. That's normal, but the API contract must be honest: "the report is current as of the last sync", not "real-time data".

For analytical queries with sumIf, argMax or FINAL, code generation doesn't help — you write the SQL by hand via JdbcClient and map it into explicit DTOs.

CDC as an alternative to outbox

The outbox pipeline works well for domain events ("order paid", "goods shipped"). But sometimes you need not an event but a mirror of a table — for example, a copy of the product catalog or a snapshot of order state.

Change Data Capture (CDC) via Debezium fits here: the tool reads the PostgreSQL change log (WAL) and publishes changes to Kafka. From there it's the same pipeline: Kafka → consumer → ClickHouse. The service itself doesn't change at all — CDC captures changes at the database level.

In ClickHouse such data is usually stored in ReplacingMergeTree by primary key with a version from the LSN or an updated_at field. When a row is updated in PostgreSQL, a new version arrives in ClickHouse, and the merge keeps the current one.

Testing with Testcontainers

Testcontainers starts a real ClickHouse server in a Docker container right inside the test:

@Testcontainers
class RevenueViewRepositoryTest {

    @Container
    static ClickHouseContainer clickHouse =
        new ClickHouseContainer("clickhouse/clickhouse-server:24.8");

    @Test
    void aggregatesRevenueByRegion() {
        insertTestEvents();
        var rows = repository.revenueByRegion(
            LocalDate.of(2026, 5, 1), LocalDate.of(2026, 5, 31));
        assertThat(rows).extracting(RevenueByRegionRow::region)
            .containsExactly("msk", "spb");
    }
}

The schema comes from the same DDL scripts used in production.

An important nuance for ReplacingMergeTree tables: the background merge is not deterministic in a test. If the test checks the "latest version" after an update, you need to either read through FINAL (as the production code does) or explicitly run OPTIMIZE TABLE ... FINAL before the check.

In short

  • The com.clickhouse:clickhouse-jdbc driver talks over HTTP (port 8123); ClickHouse is a second DataSource, and the primary PostgreSQL and @Transactional stay untouched.
  • Don't write from a handler: a double write is unreliable and single INSERTs shatter the table. The reliable path is PostgreSQL → outbox → Kafka → batch consumer → ClickHouse; for mirroring tables, CDC via Debezium replaces the outbox.
  • For small flows, async insert on the server side; for money-related data, Kafka only.
  • Idempotency: block deduplication catches a repeated batch, ReplacingMergeTree by event_id catches a duplicate from another batch.
  • Reading — a separate repository with explicit DTOs; write the SQL by hand, no code generation needed.
  • Tests — Testcontainers + real DDL; for ReplacingMergeTree, force the merge before the check.