ClickHouse в сервисе на Go это не замена PostgreSQL, а второе хранилище рядом с ней: заказы живут в транзакционной базе, а в колоночную уезжают события, по которым потом считают выручку по дням, воронки и отчёты. Отсюда две задачи интеграции: доставить события в ClickHouse так, чтобы он не захлебнулся от мелких вставок, и читать агрегаты так, чтобы долгий отчёт не положил обработчик запроса.
Клиент один, github.com/ClickHouse/clickhouse-go/v2. Он говорит с сервером по нативному протоколу (порт 9000) или по HTTP (8123) и даёт два интерфейса: свой driver.Conn и стандартный database/sql. Основной инструмент в статье нативный: он быстрее, умеет пакеты и типы ClickHouse без потерь.
Подключение: одно соединение на процесс
clickhouse.Open возвращает driver.Conn, внутри которого уже есть пул. Его создают один раз при старте и передают в репозитории; открывать соединение на запрос нельзя, это рукопожатие и аутентификация на каждый вызов.
import (
"github.com/ClickHouse/clickhouse-go/v2"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
)
func open() (driver.Conn, error) {
conn, err := clickhouse.Open(&clickhouse.Options{
Addr: []string{"clickhouse:9000"},
Auth: clickhouse.Auth{Database: "analytics", Username: "app", Password: "secret"},
Compression: &clickhouse.Compression{Method: clickhouse.CompressionLZ4},
DialTimeout: 3 * time.Second,
ReadTimeout: 30 * time.Second,
MaxOpenConns: 8,
MaxIdleConns: 4,
ConnMaxLifetime: time.Hour,
Settings: clickhouse.Settings{
"max_execution_time": 30,
},
})
if err != nil {
return nil, err
}
if err := conn.Ping(context.Background()); err != nil {
var e *clickhouse.Exception
if errors.As(err, &e) {
return nil, fmt.Errorf("clickhouse %d: %s", e.Code, e.Message)
}
return nil, err
}
return conn, nil
}
- Сжатие LZ4 включают всегда: события это тысячи похожих строк, по сети они ужимаются в разы.
Settingsэто настройки сессии, которые уедут с каждым запросом.max_execution_timeограничивает любой отчёт, который кто-то забыл ограничить сам.- Ошибки сервера приходят как
*clickhouse.Exceptionс кодом: 60 это неизвестная таблица, 241 превышен лимит памяти запроса, 252 слишком много кусков данных. Код полезен для метрик и алертов: три 252 подряд это сигнал о записи, а не о сети. - Пул маленький. Восьми соединений хватает сервису, который пишет пачками и читает десяток отчётов в минуту: ClickHouse исполняет каждый запрос многими потоками, и сотня параллельных запросов от одного сервиса его только замедлит.
Адресов может быть несколько: клиент перебирает их по порядку (или случайно, ConnOpenStrategy) и уходит на следующий при недоступности. Для кластера за балансировщиком достаточно одного адреса.
Запись: почему не вставлять по одной строке
Каждый INSERT в таблицу MergeTree создаёт на диске отдельный кусок данных (part), а фоновые слияния потом склеивают их в большие. Если вставлять по строке на событие, кусков становятся тысячи, слияния не успевают, и сервер отвечает ошибкой 252 Too many parts: порог parts_to_throw_insert по умолчанию 3000 активных кусков в партиции. Правило ClickHouse простое: вставка это тысячи строк за раз и не чаще раза в секунду на таблицу.
Отсюда три способа доставки событий из Go-сервиса, и выбор между ними зависит от потока.
- Накопить в сервисе и отправить пачкой через
PrepareBatch. Подходит фоновому воркеру или потребителю Kafka, у которого есть буфер. - Асинхронная вставка на стороне сервера: сервис шлёт по строке, а копит и пишет сам ClickHouse. Подходит, когда событий сотни в секунду и некуда их буферизовать.
- Kafka между сервисом и ClickHouse. Сервис публикует события в топик (из outbox или напрямую), а отдельный потребитель собирает пачки. Это основной вариант для прода: он переживает недоступность ClickHouse и не трогает путь оформления заказа.
Писать в ClickHouse из обработчика HTTP прямо в транзакции заказа нельзя ни одним из способов: аналитическое хранилище не должно уметь уронить оформление заказа.
Пакет: PrepareBatch, Append, Send
type Event struct {
At time.Time
OrderID uint64
Kind string
Amount float64
}
func insertBatch(ctx context.Context, conn driver.Conn, events []Event) error {
batch, err := conn.PrepareBatch(ctx, "INSERT INTO order_events (at, order_id, kind, amount)")
if err != nil {
return err
}
for _, e := range events {
if err := batch.Append(e.At, e.OrderID, e.Kind, e.Amount); err != nil {
return err
}
}
return batch.Send()
}
PrepareBatch открывает запрос, Append складывает строки в блок в памяти клиента, Send отправляет блок одним запросом. Порядок аргументов Append повторяет список колонок в INSERT. Типы клиент сверяет с типами колонок: целое число в UInt64 он приведёт сам, а вот строку вместо DateTime или nil в колонку без Nullable отвергнет на Append, ещё до отправки. Если строки удобнее хранить структурами, есть batch.AppendStruct(&e) с тегами ch:"order_id".
Пачка до max_insert_block_size строк (около миллиона) в одну партицию ложится атомарно: или весь блок, или ничего, поэтому упавший Send можно повторять целиком. Буфер в сервисе сбрасывается по двум условиям, по размеру и по времени, и обязательно при остановке процесса:
type Buffer struct {
conn driver.Conn
events []Event
maxSize int
maxWait time.Duration
last time.Time
}
func (b *Buffer) Add(ctx context.Context, e Event) error {
b.events = append(b.events, e)
if len(b.events) >= b.maxSize || time.Since(b.last) >= b.maxWait {
return b.Flush(ctx)
}
return nil
}
func (b *Buffer) Flush(ctx context.Context) error {
if len(b.events) == 0 {
return nil
}
if err := insertBatch(ctx, b.conn, b.events); err != nil {
return err
}
slog.Info("clickhouse batch", "rows", len(b.events))
b.events = b.events[:0]
b.last = time.Now()
return nil
}
Десять тысяч строк или пять секунд, что наступит раньше, это разумные значения для старта. Буфер ограничен: если ClickHouse лежит дольше, чем помещается в память, события либо отбрасываются с метрикой, либо их источником должна быть Kafka, где они подождут.
Повтор пачки без дублей
Повтор после ошибки это «хотя бы раз», и та же пачка может оказаться записанной дважды: сеть оборвалась уже после того, как сервер её принял. Защиты две.
- Дедупликация блоков. Реплицируемые таблицы (
ReplicatedMergeTree) запоминают хеши последних вставленных блоков и молча отбрасывают точный повтор. У обычногоMergeTreeэто окно выключено (non_replicated_deduplication_window = 0), его включают в настройках таблицы. Чтобы повтор считался «тем же блоком» независимо от содержимого, вставке дают явный ключ:
ctx = clickhouse.Context(ctx, clickhouse.WithSettings(clickhouse.Settings{
"insert_deduplication_token": "order-events/3/18220-18331",
}))
batch, err := conn.PrepareBatch(ctx, "INSERT INTO order_events")
Для потребителя Kafka естественный ключ это топик, партиция и диапазон смещений пачки.
- Движок
ReplacingMergeTreeпо ключу события: дубли схлопнутся при слиянии, а запросы пишут сFINALили черезargMax. Это защита от дублей любого происхождения, не только от повторов пачки, но она стоит ресурсов на чтении.
Асинхронная вставка
Когда событий немного и буферизовать их негде (например, они рождаются в разных обработчиках), пачки собирает сам сервер:
func insertAsync(ctx context.Context, conn driver.Conn, e Event) error {
return conn.AsyncInsert(ctx,
"INSERT INTO order_events (at, order_id, kind, amount) VALUES (?, ?, ?, ?)",
true, e.At, e.OrderID, e.Kind, e.Amount)
}
Сервер копит такие вставки в памяти и сбрасывает на диск по размеру (async_insert_max_data_size, 10 МиБ) или по таймеру (async_insert_busy_timeout_ms, 200 мс). Третий аргумент wait решает, чего ждёт вызов: true возвращается после записи буфера на диск (и получает ошибку, если она случилась), false сразу после приёма в буфер, и ошибка записи тогда теряется. Для событий, которые нельзя терять, только true. Дедупликация асинхронных вставок отдельная (async_insert_deduplicate) и по умолчанию выключена.
Потребитель Kafka: пачка между FetchMessage и Commit
Сервис публикует события в топик, а потребитель на kafka-go собирает их в пачку и подтверждает смещения только после успешного Send. Так недоступность ClickHouse превращается в отставание потребителя, а не в потерю событий.
func consumeToClickHouse(ctx context.Context, r *kafka.Reader,
flush func(context.Context, []kafka.Message) error) error {
const maxRows = 10000
const maxWait = 5 * time.Second
var batch []kafka.Message
deadline := time.Now().Add(maxWait)
for {
fetchCtx, cancel := context.WithDeadline(ctx, deadline)
m, err := r.FetchMessage(fetchCtx)
cancel()
if err != nil && ctx.Err() != nil {
return ctx.Err()
}
if err == nil {
batch = append(batch, m)
}
if len(batch) >= maxRows || (time.Now().After(deadline) && len(batch) > 0) {
if err := flush(ctx, batch); err != nil {
slog.Error("clickhouse flush failed, will retry", "rows", len(batch), "err", err)
time.Sleep(time.Second)
continue
}
if err := r.CommitMessages(ctx, batch...); err != nil {
return err
}
batch = batch[:0]
}
if time.Now().After(deadline) {
deadline = time.Now().Add(maxWait)
}
}
}
Два места, где здесь обычно ошибаются: подтверждение смещений до Send (при падении пачка пропадёт) и MaxBytes читателя меньше размера пачки (потребитель будет ходить за сообщениями чаще, чем нужно). Повтор flush с тем же insert_deduplication_token безопасен. Тонкости читателя kafka-go, от CommitInterval до остановки группы, разобраны в Kafka на Go в проде.
Есть и вариант без кода: табличный движок Kafka внутри ClickHouse читает топик сам, а материализованное представление перекладывает строки в MergeTree. Сервису это ничего не стоит, но разбор формата сообщения и ошибки чтения уезжают в настройки сервера, где их труднее тестировать и наблюдать.
Чтение: отдельный репозиторий аналитики
Чтение из ClickHouse живёт в своём репозитории и под своим пользователем с правами только на чтение. Запрос с параметрами, результат построчно или сразу в срез структур:
type DailyRevenue struct {
Day time.Time `ch:"day"`
Revenue float64 `ch:"revenue"`
Orders uint64 `ch:"orders"`
}
func revenueByDay(ctx context.Context, conn driver.Conn, from, to time.Time) ([]DailyRevenue, error) {
var out []DailyRevenue
err := conn.Select(ctx, &out, `
SELECT toDate(at) AS day, sum(amount) AS revenue, uniqExact(order_id) AS orders
FROM order_events
WHERE kind = 'paid' AND at >= @from AND at < @to
GROUP BY day ORDER BY day`,
clickhouse.Named("from", from), clickhouse.Named("to", to))
return out, err
}
Параметры бывают позиционные (?) и именованные (@from плюс clickhouse.Named); склеивать SQL строками нельзя по тем же причинам, что и в PostgreSQL. conn.Query с rows.Next() и rows.Scan или rows.ScanStruct нужен, когда строк много и их обрабатывают потоком, а не держат в памяти.
Типы при чтении строгие. count() и uniqExact возвращают UInt64, и сканировать их в int нельзя: клиент ответит converting UInt64 to *int is unsupported. try using *uint64. toDate это time.Time, Decimal это decimal.Decimal из shopspring/decimal, Nullable(String) это *string. Теги ch связывают поля структуры с именами колонок, поэтому у вычисляемых выражений в SELECT всегда стоит AS.
Дедлайн контекста здесь важнее, чем где-либо: отчёт за год может считаться минутами, и обработчик HTTP обязан отпустить запрос. Клиент при отмене контекста посылает серверу отмену, а max_execution_time страхует со стороны сервера. Тяжёлые отчёты вообще не считают в запросе пользователя: их считает фоновая задача, а обработчик отдаёт готовый результат из таблицы или кэша.
Схема и миграции
Схему ClickHouse ведут теми же миграциями, что и PostgreSQL, в отдельном каталоге: goose знает диалект clickhouse и работает через стандартный драйвер (clickhouse.OpenDB). Содержимое миграций проще, чем в транзакционной базе: CREATE TABLE ... ENGINE = MergeTree, ALTER TABLE ... ADD COLUMN, материализованные представления. Две особенности: у ClickHouse нет транзакций на DDL, поэтому миграция из нескольких выражений может остановиться посередине, и в кластере каждое выражение нужно с ON CLUSTER, иначе таблица появится на одном узле.
-- +goose Up
CREATE TABLE order_events
(
at DateTime,
order_id UInt64,
kind LowCardinality(String),
amount Float64
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(at)
ORDER BY (kind, at);
Тестирование с Testcontainers
Модуль testcontainers-go/modules/clickhouse поднимает сервер и отдаёт адрес нативного порта:
import tcclickhouse "github.com/testcontainers/testcontainers-go/modules/clickhouse"
c, err := tcclickhouse.Run(ctx, "clickhouse/clickhouse-server:24.8",
tcclickhouse.WithUsername("app"), tcclickhouse.WithPassword("secret"),
tcclickhouse.WithDatabase("analytics"))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { c.Terminate(ctx) })
host, err := c.ConnectionHost(ctx)
conn, err := clickhouse.Open(&clickhouse.Options{
Addr: []string{host},
Auth: clickhouse.Auth{Database: "analytics", Username: "app", Password: "secret"},
})
Что проверяют такие тесты: что пачка из тысячи строк ложится и читается обратно нужными типами, что отчётный запрос даёт ожидаемые суммы на подготовленных данных, что миграции накатываются на пустой сервер. Контейнер стартует секунд за десять, поэтому он один на пакет, в TestMain, а каждый тест работает со своей таблицей или чистит её TRUNCATE. Вставка в тесте синхронная, читать можно сразу: в отличие от Elasticsearch, у ClickHouse нет задержки видимости.
Глубже: database/sql через OpenDBрасширенное
Стандартный интерфейс нужен, когда рядом уже стоит sqlx, инструмент миграций или код, рассчитанный на *sql.DB. clickhouse.OpenDB принимает те же Options, а пачка в нём выглядит как транзакция: Begin, Prepare на INSERT, цикл Exec по строкам, Commit отправляет блок. Транзакции в смысле изоляции здесь нет, это только способ собрать пачку. Нативные типы через database/sql проходят хуже (массивы и вложенные структуры приходится описывать руками), и для нового кода стандартный интерфейс выбирают только ради совместимости.
Глубже: когда ClickHouse недоступенрасширенное
Путь записи проектируют так, чтобы недоступность была задержкой, а не потерей: события ждут в Kafka или в outbox, потребитель повторяет пачку с паузой, метрики показывают возраст самого старого непрочитанного события и число упавших Send. Буфер в памяти сервиса, если он есть, ограничен сверху и при переполнении отбрасывает события со счётчиком, а не растёт до OOM.
Путь чтения деградирует честно: дедлайн контекста в 2-5 секунд, ответ 503 для отчётного эндпоинта, кэш последнего успешного результата для витрин, которые смотрят часто. Проверка здоровья сервиса не включает ClickHouse в обязательные зависимости: оформление заказов не должно уходить из балансировщика из-за аналитики.
Коротко
- Один
driver.Connна процесс, пул на 4-8 соединений, LZ4,max_execution_timeв настройках сессии. - Мелкие вставки запрещены: каждый
INSERTэто кусок на диске, порог 3000 кусков в партиции даёт ошибку 252. - Пачка через
PrepareBatchиAppend, сброс по размеру и по времени, обязательныйFlushпри остановке. - Повтор пачки безопасен с
insert_deduplication_token; у обычного MergeTree окно дедупликации нужно включить. AsyncInsertкопит на сервере (10 МиБ или 200 мс),wait = trueвозвращает ошибки записи,wait = falseих теряет.- В проде между сервисом и ClickHouse стоит Kafka: потребитель подтверждает смещения только после успешного
Send. - Чтение в отдельном репозитории под read-only пользователем, типы строгие:
UInt64только вuint64, выражения сAS. - Долгие отчёты считает фоновая задача, обработчик отдаёт готовое; дедлайн контекста отменяет запрос и на сервере.
- Миграции
gooseс диалектомclickhouse, в кластере каждое выражениеON CLUSTER. - Тесты на контейнере: пачка, отчётный запрос на известных данных, прогон миграций.
Что почитать дальше
- Как устроен ClickHouse — колонки, куски данных и слияния, из которых растут все правила записи.
- Моделирование и запросы — ключ сортировки, партиции,
ReplacingMergeTreeи агрегатные функции. - ClickHouse в проде — репликация, TTL, бэкапы и что мониторить.
- PostgreSQL или ClickHouse — где проходит граница между транзакционной и аналитической нагрузкой.
- Kafka на Go в проде — читатель
kafka-go, подтверждение смещений и остановка потребителя. - Паттерны распределённых систем на Go — outbox как источник событий для аналитики.