~/posts/postgres-to-clickhouse-then-embeddings-pipeline

Postgres to ClickHouse, then the embeddings pipeline became the real migration

12 min read 2222 words
dataaianalyticsdatabasessystemsperformance

tl;dr

Backfilling ~4.8M documents from Postgres into ClickHouse pushed a RabbitMQ queue to 2,944,392 messages, because the embedding worker ran one text per encode() call. A microbatching broker, a single execution lane, in-flight de-dupe, bf16 and TensorRT with persistent compile caches made inference fast enough that ClickHouse inserts became the new bottleneck.

This started as a storage migration: move an old text-only dataset from Postgres to ClickHouse, then use ClickHouse for analytics and cheaper storage.

During the backfill, it turned into an inference optimization project.

A Grafana card showed a RabbitMQ queue that kept growing while acks per second stayed low, and the GPU on a dedicated box wasn’t nearly as busy as it should have been. The backlog wasn’t spiky. It only went one way.

This post is about the embedding service: what it does, why it was slow, what I changed, and why TensorRT plus batching did most of the heavy lifting.

The dataset

Rows migrated~4.80M (text-only rows ported out of Postgres)
Tokens109,743,122 (counted with o200k_base)
Text lengthavg 81.22, min 1, max 30,228 characters
GPURTX 4090 (24 GB VRAM, bare metal host, dockerized)

Why ClickHouse

It fit a text lake for practical reasons:

  • zstd compression, for predictable storage wins,
  • fast analytics for KPIs and time series,
  • a clean path to embedding search with cosine similarity (a dot product on normalized vectors),
  • tiered storage and cold offload,
  • JSON support for flexible metadata.

Migration strategy

Backfill, then switch:

  1. stop the scrapers,
  2. backfill the existing Postgres corpus into ClickHouse,
  3. rewrite the scrapers to write into ClickHouse directly,
  4. resume ingestion after the cutover.

That avoided dual-write complexity and kept correctness simple while the storage layer changed underneath.

The schema

The lake has two tables: one for documents, one for embeddings (a list of numbers per document, where similar texts get nearby vectors).

What I wanted from it:

  • document metadata and raw text that are queryable and cheap to scan,
  • embeddings that are append-friendly and partitioned by model and time,
  • LowCardinality where values repeat a lot,
  • monthly partitions, for easier operations and cold storage policies.
CREATE TABLE lake.docs
(
  source LowCardinality(String),
  source_id String,
  doc_uuid UUID,
  updated_at DateTime64(3),
  ingest_at DateTime64(3),
  author String,
  text String CODEC(ZSTD(3)),
  meta_json String CODEC(ZSTD(3))
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(ingest_at)
ORDER BY (source, source_id, updated_at)
SETTINGS index_granularity = 8192;

CREATE TABLE lake.embeddings
(
  model LowCardinality(String),
  doc_uuid UUID,
  updated_at DateTime64(3),
  embedding Array(Float32) CODEC(ZSTD(6))
)
ENGINE = MergeTree
PARTITION BY (model, toYYYYMM(updated_at))
ORDER BY (model, doc_uuid)
SETTINGS index_granularity = 8192;
  • docs is partitioned by ingest month, and ordering by (source, source_id, updated_at) makes versions and timelines easy to query.
  • embeddings is partitioned by (model, month(updated_at)) to keep multi-model history clean, and ordered by (model, doc_uuid) to match how it’s joined and read.
  • Embeddings are Array(Float32) with a stronger zstd level than text, because vectors are big and compress well.

What embedding it all would cost

The dataset holds 109,743,122 tokens (with o200k_base), and the number keeps growing with ingestion. Even at this snapshot, it’s big enough that embedding costs are easy to underestimate.

Before going all in on self-hosted embeddings, I worked out what it would cost to embed the whole dataset once with a few hosted APIs, at standard pricing and at batch pricing where it exists:

Estimated cost to embed 109,743,122 tokens once (USD)
loading chart…
Estimated total API cost to embed the full dataset once. Batch pricing is shown where the provider publishes it.

This isn’t an argument that hosted APIs are bad. It’s why a local pipeline has to be efficient: at this scale, embeddings easily become a recurring cost, a backlog problem, or both.

Where embeddings fit in

The embedding worker is an async service. It consumes RabbitMQ messages, embeds text with BAAI/bge-m3 (FlagEmbedding), then writes the results or triggers downstream work. It handles two independent queues:

  • to_embed: unrelated to the migration; this text gets indexed into ElasticSearch (embedding happens here, indexing elsewhere).
  • to_embed_text: the migration path; every text document inserted into ClickHouse (lake.docs and lake.embeddings).

A quality note: I don’t chunk long documents in this iteration. The raw text is stored in full, but embeddings use a max_length of 1024, so long texts are effectively truncated in their embedding. For backfill throughput and stability, that was an acceptable trade. If retrieval quality becomes the priority, chunking is the obvious next step.

The signal that something was wrong

The Grafana panel showed queue depth climbing while acks per second stayed low. From the 1-minute export on 2025-12-26:

TimeQueue depth
14:3513,807first above 10k
15:05100,485first above 100k
16:14512,420first above 500k
16:38808,778first above 800k
17:081,005,023first above 1M
21:002,009,989first above 2M
21:532,500,356first above 2.5M
22:012,589,117
22:502,944,392peak

The steepest 1-hour increase in the export was 619,317 messages, from 16:02 to 17:02.

RabbitMQ queue depth (to_embed_text) during backfill
loading chart…
5-minute sampling from a 1-minute Grafana export of RabbitMQ queue depth (to_embed_text), on 2025-12-26.

The exact slope matters less than what it means: the worker was consuming messages more slowly than they were being produced, all the time.

Why it was slow

The original design was correct but not GPU-friendly. Each message handler:

  1. ran inference on exactly one text,
  2. then did the storage writes,
  3. and only then acked and moved on.

There was concurrency at the message level, but no batching at the inference level. Even with several handlers, their inference calls weren’t coordinated, so the GPU rarely saw a real batch. And when I noticed the backlog, the configuration was effectively serial anyway: concurrency 1, and RabbitMQ delivering one unacked message at a time.

So the rhythm was: run a tiny encode call, block on database writes and network round trips, repeat. That can leave even a very strong GPU mostly idle.

before: one text per encode()GPUwritesafter: microbatch (≤ 10 ms window, ≤ 12 texts)requestsGPUwritesillustrative, not to scale
Before: each encode() call handles one text, and the GPU sits idle while the writes for that message finish. After: requests from several handlers pile up for at most 10 ms (or 12 texts), then run as one batch. Illustrative, not to scale.

What I changed

In short: I turned a per-message, single-item encode() loop into something that behaves like a GPU service. Inference goes through a microbatching broker and a single execution lane, with in-flight de-dupe and an optional cache in front. I tightened precision and attention settings, and TensorRT (with persistent compile caches) did the real acceleration. After that, inference stopped being the slowest part, and the bottleneck moved to ClickHouse inserts, which is a much better problem to have.

Most of this lives in services/embedding/main.py, plus the container and compose changes that make TensorRT practical.

1. One broker for all inference (microbatching)

A new InferenceBroker centralizes inference requests so they can be batched across handlers:

  • handlers call the broker’s embed,
  • the broker collects requests for up to microbatch_max_wait_ms (10 ms),
  • runs a single model.encode over the batch,
  • and hands each handler its own result.

Settings used here: enabled, microbatch_max_wait_ms = 10, microbatch_max_batch_size = 12.

Transformer inference has a fixed cost per call (kernel launches, scheduling, Python overhead). Batching spreads that cost over more texts and keeps the GPU busier.

The catch: microbatching only helps if several requests arrive inside the window. If the pipeline is IO-bound or configured with very low concurrency, batches often shrink back to 1.

model.encode()BrokerHandler BHandler Amodel.encode()BrokerHandler BHandler Await up to 10msembed(text A)embed(text B)encode([A, B, ...])vectorsvector(A)vector(B)

2. A single execution lane, on purpose

All model execution goes through one lane (one thread). Concurrent inference calls on the same model and GPU tend to cause contention and allocator pressure. One lane makes latency more predictable and batching more effective.

The trade-off: in isolation, single-lane inference was about 20% slower in raw microbenchmarks. End to end, it was still a net win once batching and TensorRT were in place.

3. In-flight de-dupe and an optional cache

The broker tracks in-flight requests keyed by model settings and normalized text. If the same text arrives while it’s already being embedded, both callers share the same future. That matters because about 5% of the content is duplicate.

There’s also an optional in-process LRU cache with the same key, capped around 256 MB in this deployment.

4. Precision and attention backend

A config key picks the inference precision (auto, fp16, bf16, fp32). On this hardware, bf16 was about 5% faster than the previous precision behaviour.

There’s also a best-effort attention backend selector (SDPA variants, plus flash attention if installed). It helped, but it wasn’t the main improvement here.

5. TensorRT, the main time saver

TensorRT made the biggest difference to steady-state inference.

In the container:

  • the base image moved to nvcr.io/nvidia/tensorrt:25.08-py3,
  • torch-tensorrt is installed at build time, so it’s always there,
  • a build fix copies .python-version before uv sync, so the interpreter choice is stable and the venv isn’t unexpectedly recreated.

At runtime, the config picks inference_backend (torch or tensorrt). With tensorrt, the service tries torch.compile with the TensorRT backend. If that fails, it logs and falls back to plain torch instead of crashing.

The TensorRT warmup (compile) took about 15 seconds, and the compiled artifacts, about 4.3 GB, are kept in mounted cache directories.

6. Persistent compile caches

So containers don’t recompile on every recreation, compose mounts persistent cache directories for the torch inductor cache, the torch extensions cache, the HuggingFace cache (model files) and the TensorRT cache.

That’s the difference between a predictable worker restart and one that spends time rebuilding the world.

7. Faster config snapshots

Config snapshots used to be deep-copied through a JSON encode and decode. That’s correct, but slow, and it can coerce types unexpectedly. It’s now copy.deepcopy, which is simpler and faster for plain Python dict trees.

Results (inference only)

Inference latency after optimization
loading chart…
Min/mean/p95/max; Qwen values scaled vs bge-m3.

No single number here is the point. The worker finally behaved like it had a GPU: inference got fast and consistent enough that the rest of the pipeline started to dominate.

The next bottleneck: ClickHouse inserts

With fast inference, the limit moved to ClickHouse insert latency. That’s not surprising:

  • the to_embed_text path does several writes per message (docs and embeddings),
  • slow inserts leave the handler blocked on IO between messages,
  • blocked handlers pick up the next message later,
  • and fewer requests land in each microbatch window, so batches get smaller.

The pipeline is now IO-paced. That’s still better than model-paced, because the embedding engine isn’t the main limit anymore.

To drain multi-million-message backlogs faster, the next work is mostly about write amplification:

  • fewer round trips (stage docs and embeddings more efficiently),
  • more concurrency, carefully, while keeping in-flight memory bounded,
  • checking RabbitMQ delivery settings so the broker gets enough parallel demand to form real batches.

This was supposed to be a ClickHouse post about compression and analytics. It still is, but the most useful part turned out to be this: at scale, an embedding pipeline isn’t just a model call. It’s a queueing system, a batching system, a compilation system and a storage system, and once you treat it like one, the fixes are straightforward.

hash: 127c
EOF