← назад к разделу

ClickHouse в сервисе на Python это не замена PostgreSQL, а второе хранилище рядом с ней: заказы живут в транзакционной базе, а в колоночную уезжают события, по которым потом считают выручку по дням, воронки и отчёты. Отсюда две задачи интеграции: доставить события в ClickHouse так, чтобы он не захлебнулся от мелких вставок, и читать агрегаты так, чтобы долгий отчёт не положил обработчик запроса.

Официальный клиент — clickhouse-connect. Он говорит с сервером по HTTP (порт 8123, с TLS 8443), сжимает трафик, умеет вставлять списки строк, pandas и arrow и возвращать результат в удобной форме. Второй клиент, clickhouse-driver, работает по нативному протоколу на порту 9000; он чуть быстрее на больших выборках, но у официального клиента шире набор форматов и он развивается активнее. В статье основной инструмент clickhouse-connect.

Обязательно

Подключение: один клиент на процесс и вопрос сессий

Клиент создают один раз при старте и передают в репозитории; открывать соединение на запрос нельзя, это рукопожатие и аутентификация на каждый вызов.

import clickhouse_connect
from clickhouse_connect.driver.exceptions import DatabaseError


def new_client(settings: ClickHouseSettings):
    client = clickhouse_connect.get_client(
        host=settings.host,
        port=8443,
        secure=True,
        username="app",
        password=settings.password,
        database="analytics",
        compress=True,
        connect_timeout=3,
        send_receive_timeout=30,
        autogenerate_session_id=False,
        settings={"max_execution_time": 30},
    )
    client.ping()
    return client
  • autogenerate_session_id=False обязателен для клиента, которого делят. По умолчанию клиент заводит на сервере сессию и шлёт её идентификатор с каждым запросом, а ClickHouse исполняет в одной сессии только один запрос за раз: второй параллельный получает ошибку с кодом 373 Session is locked by a concurrent client. Сервис с несколькими одновременными обработчиками упрётся в неё в первый же день. Без сессии теряются только серверные SET и временные таблицы между запросами, которые сервису не нужны.
  • send_receive_timeout по умолчанию 300 секунд. Это пять минут ожидания ответа на один отчёт из обработчика HTTP. Для сервиса ставят десятки секунд, а тяжёлые отчёты уводят в фоновую задачу.
  • settings это настройки сервера, которые уедут с каждым запросом. max_execution_time ограничивает любой отчёт, который кто-то забыл ограничить сам, и действует на стороне сервера, а не только на клиенте.
  • Сжатие включено по умолчанию: события это тысячи похожих строк, по сети они ужимаются в разы.
  • Ошибки сервера приходят как DatabaseError с текстом, где есть код: Code: 60 это неизвестная таблица, 241 превышен лимит памяти запроса, 252 слишком много кусков данных. Код полезен для метрик и алертов: три 252 подряд это сигнал о записи, а не о сети.

Внутри клиента пул соединений urllib3, вызовы синхронные и безопасные для потоков (при выключенных сессиях). Из асинхронного кода их зовут через asyncio.to_thread; есть и get_async_client с теми же методами под await, он работает поверх aiohttp. Для сервиса, который пишет пачками раз в несколько секунд и читает десяток отчётов в минуту, пула потоков хватает.

Запись: почему не вставлять по одной строке

Каждый INSERT в таблицу MergeTree создаёт на диске отдельный кусок данных (part), а фоновые слияния потом склеивают их в большие. Если вставлять по строке на событие, кусков становятся тысячи, слияния не успевают, и сервер отвечает ошибкой 252 Too many parts: порог parts_to_throw_insert по умолчанию 3000 активных кусков в партиции. Правило ClickHouse простое: вставка это тысячи строк за раз и не чаще раза в секунду на таблицу.

по одной строке INSERT ×1000 1000 кусков на диске слияния не успевают: Too many parts пачкой INSERT 1000 строк один кусок слияние в фоне

Каждый INSERT создаёт кусок данных, и тысяча мелких вставок это тысяча кусков; ClickHouse рассчитан на пачки в тысячи строк.

Отсюда три способа доставки событий из Python-сервиса, и выбор между ними зависит от потока.

  1. Накопить в сервисе и отправить пачкой через client.insert. Подходит фоновому воркеру или потребителю Kafka, у которого есть буфер.
  2. Асинхронная вставка на стороне сервера: сервис шлёт по строке, а копит и пишет сам ClickHouse. Подходит, когда событий сотни в секунду и некуда их буферизовать.
  3. Kafka между сервисом и ClickHouse. Сервис публикует события в топик (из outbox или напрямую), а отдельный потребитель собирает пачки. Это основной вариант для прода: он переживает недоступность ClickHouse и не трогает путь оформления заказа.

Писать в ClickHouse из обработчика HTTP прямо в транзакции заказа нельзя ни одним из способов: аналитическое хранилище не должно уметь уронить оформление заказа.

Пачка: insert со списком строк

from datetime import datetime

from pydantic import BaseModel

COLUMNS = ["at", "order_id", "kind", "amount"]


class Event(BaseModel):
    at: datetime
    order_id: int
    kind: str
    amount: float


def insert_batch(client, events: list[Event], token: str | None = None) -> int:
    rows = [[e.at, e.order_id, e.kind, e.amount] for e in events]
    settings = {"insert_deduplication_token": token} if token else None
    summary = client.insert("order_events", rows, column_names=COLUMNS, settings=settings)
    return summary.written_rows

insert принимает список строк, каждая — список значений в порядке column_names, собирает из них один блок и отправляет одним запросом. Типы клиент сверяет с типами колонок: int в UInt64 он приведёт сам, а строку вместо DateTime или None в колонку без Nullable отвергнет до отправки. С datetime без часового пояса осторожнее: по умолчанию клиент считает его местным временем процесса (naive_datetime_insert = local), и контейнер с другим TZ сдвинет все события; передают осведомлённые datetime в UTC. Если данные уже лежат в pandas, есть insert_df, и это самый быстрый путь для больших объёмов.

Пачка до max_insert_block_size строк (около миллиона) в одну партицию ложится атомарно: или весь блок, или ничего, поэтому упавший insert можно повторять целиком. Буфер в сервисе сбрасывается по двум условиям, по размеру и по времени, и обязательно при остановке процесса:

import asyncio
import time


class EventBuffer:
    def __init__(self, client, max_size: int = 10_000, max_wait: float = 5.0) -> None:
        self.client = client
        self.events: list[Event] = []
        self.max_size = max_size
        self.max_wait = max_wait
        self.last_flush = time.monotonic()
        self.lock = asyncio.Lock()

    async def add(self, event: Event) -> None:
        async with self.lock:
            self.events.append(event)
            if len(self.events) >= self.max_size or time.monotonic() - self.last_flush >= self.max_wait:
                await self._flush()

    async def flush(self) -> None:
        async with self.lock:
            await self._flush()

    async def _flush(self) -> None:
        if not self.events:
            return
        batch, self.events = self.events, []
        await asyncio.to_thread(insert_batch, self.client, batch)
        log.info("clickhouse batch", rows=len(batch))
        self.last_flush = time.monotonic()

Десять тысяч строк или пять секунд, что наступит раньше, это разумные значения для старта; flush зовут из lifespan при остановке. Буфер ограничен: если ClickHouse лежит дольше, чем помещается в память, события либо отбрасываются с метрикой, либо их источником должна быть Kafka, где они подождут.

Повтор пачки без дублей

Повтор после ошибки это «хотя бы раз», и та же пачка может оказаться записанной дважды: сеть оборвалась уже после того, как сервер её принял. Защиты две.

  • Дедупликация блоков. Реплицируемые таблицы (ReplicatedMergeTree) запоминают хеши последних вставленных блоков и молча отбрасывают точный повтор. У обычного MergeTree это окно выключено (non_replicated_deduplication_window = 0), его включают в настройках таблицы. Чтобы повтор считался «тем же блоком» независимо от содержимого, вставке дают явный ключ через insert_deduplication_token, как в insert_batch выше. Для потребителя Kafka естественный ключ это топик, партиция и диапазон смещений пачки: order-events/3/18220-18331.
  • Движок ReplacingMergeTree по ключу события: дубли схлопнутся при слиянии, а запросы пишут с FINAL или через argMax. Это защита от дублей любого происхождения, не только от повторов пачки, но она стоит ресурсов на чтении.

Асинхронная вставка

Когда событий немного и буферизовать их негде (например, они рождаются в разных обработчиках), пачки собирает сам сервер:

def insert_async(client, event: Event) -> None:
    client.insert(
        "order_events",
        [[event.at, event.order_id, event.kind, event.amount]],
        column_names=COLUMNS,
        settings={"async_insert": 1, "wait_for_async_insert": 1},
    )

Сервер копит такие вставки в памяти и сбрасывает на диск по размеру (async_insert_max_data_size, 10 МиБ) или по таймеру (async_insert_busy_timeout_ms, 200 мс). wait_for_async_insert решает, чего ждёт вызов: 1 возвращается после записи буфера на диск (и получает ошибку, если она случилась), 0 сразу после приёма в буфер, и ошибка записи тогда теряется. Для событий, которые нельзя терять, только 1. Дедупликация асинхронных вставок отдельная (async_insert_deduplicate) и по умолчанию выключена.

Потребитель Kafka: пачка между getmany и commit

Сервис публикует события в топик, а потребитель на aiokafka собирает их в пачку и подтверждает смещения только после успешной вставки. Так недоступность ClickHouse превращается в отставание потребителя, а не в потерю событий.

from aiokafka import AIOKafkaConsumer


async def consume_to_clickhouse(consumer: AIOKafkaConsumer, client, stop: asyncio.Event) -> None:
    while not stop.is_set():
        batches = await consumer.getmany(timeout_ms=5000, max_records=10_000)
        records = [r for partition_records in batches.values() for r in partition_records]
        if not records:
            continue
        events = [Event.model_validate_json(r.value) for r in records]
        token = "order-events/" + "/".join(f"{tp.partition}:{rs[0].offset}-{rs[-1].offset}" for tp, rs in batches.items())
        while True:
            try:
                await asyncio.to_thread(insert_batch, client, events, token)
                break
            except DatabaseError as e:
                log.error("clickhouse flush failed, will retry", rows=len(events), error=str(e)[:200])
                await asyncio.sleep(1)
        await consumer.commit()

Потребитель создают с enable_auto_commit=False: иначе смещения уедут по таймеру до того, как пачка легла в ClickHouse, и при падении процесса она пропадёт. getmany возвращает всё, что накопилось за пять секунд, но не больше max_records, и это удобнее, чем собирать пачку из одиночных __anext__. Повтор insert_batch с тем же токеном безопасен. Тонкости потребителя, от max_poll_interval_ms до остановки группы, разобраны в Kafka на Python в проде.

Есть и вариант без кода: табличный движок Kafka внутри ClickHouse читает топик сам, а материализованное представление перекладывает строки в MergeTree. Сервису это ничего не стоит, но разбор формата сообщения и ошибки чтения уезжают в настройки сервера, где их труднее тестировать и наблюдать.

Чтение: отдельный репозиторий аналитики

Чтение из ClickHouse живёт в своём репозитории и под своим пользователем с правами только на чтение. Запрос с параметрами, результат построчно или словарями:

from datetime import date, datetime
from decimal import Decimal


class DailyRevenue(BaseModel):
    day: date
    revenue: float
    orders: int


REVENUE_BY_DAY = """
    SELECT toDate(at) AS day, sum(amount) AS revenue, uniqExact(order_id) AS orders
    FROM order_events
    WHERE kind = 'paid' AND at >= {from:DateTime} AND at < {to:DateTime}
    GROUP BY day ORDER BY day
"""


def revenue_by_day(client, start: datetime, end: datetime) -> list[DailyRevenue]:
    result = client.query(REVENUE_BY_DAY, parameters={"from": start, "to": end})
    return [DailyRevenue.model_validate(row) for row in result.named_results()]

Параметры пишут в виде {имя:Тип} — это серверная подстановка ClickHouse, значения уезжают отдельно от текста запроса, и склеивать SQL строками не нужно по тем же причинам, что и в PostgreSQL. result.named_results() отдаёт словари по именам колонок, result.result_rows — списки, query_df — DataFrame. Когда строк много и их обрабатывают потоком, берут query_rows_stream в блоке with: он отдаёт строки по мере прихода, не собирая весь результат в памяти.

Типы при чтении. count() и uniqExact возвращают UInt64, в Python это обычный int; toDate это date, DateTime без часового пояса в колонке приходит наивным datetime в UTC, Decimal это decimal.Decimal, Nullable(String) это str | None, Array это list. У вычисляемых выражений в SELECT всегда стоит AS, иначе в named_results колонка будет называться sum(amount).

Таймаут здесь важнее, чем где-либо: отчёт за год может считаться минутами, и обработчик HTTP обязан отпустить запрос. send_receive_timeout клиента обрывает ожидание на стороне сервиса, max_execution_time страхует со стороны сервера и останавливает сам запрос. Тяжёлые отчёты вообще не считают в запросе пользователя: их считает фоновая задача, а обработчик отдаёт готовый результат из таблицы или кэша.

Схема и миграции

У Alembic нет диалекта ClickHouse, и схему ведут отдельно от миграций PostgreSQL: каталог migrations/clickhouse/ с нумерованными SQL-файлами и короткий раннер, который применяет новые файлы по порядку и записывает их имена в таблицу версий (client.command на каждое выражение); готовые раннеры вроде clickhouse-migrations делают то же. Содержимое миграций проще, чем в транзакционной базе: CREATE TABLE ... ENGINE = MergeTree, ALTER TABLE ... ADD COLUMN, материализованные представления. Две особенности: у ClickHouse нет транзакций на DDL, поэтому миграция из нескольких выражений может остановиться посередине, и в кластере каждое выражение нужно с ON CLUSTER, иначе таблица появится на одном узле.

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

Контейнер ClickHouse открывает и нативный порт, и HTTP; клиенту нужен второй:

import pytest
from testcontainers.clickhouse import ClickHouseContainer


@pytest.fixture(scope="session")
def ch_client():
    with ClickHouseContainer("clickhouse/clickhouse-server:24.8", username="app", password="secret", dbname="analytics") as container:
        client = clickhouse_connect.get_client(
            host=container.get_container_host_ip(),
            port=int(container.get_exposed_port(8123)),
            username="app",
            password="secret",
            database="analytics",
            autogenerate_session_id=False,
        )
        apply_migrations(client)
        yield client

Контейнер ждёт ответа Ok. на HTTP-порту, поэтому к моменту выхода из конструктора сервер принимает запросы. Что проверяют такие тесты: что пачка из тысячи строк ложится и читается обратно нужными типами, что отчётный запрос даёт ожидаемые суммы на подготовленных данных, что миграции накатываются на пустой сервер. Контейнер стартует секунд за десять, поэтому он один на сессию, а каждый тест работает со своей таблицей или чистит её TRUNCATE. Вставка в тесте синхронная, читать можно сразу: в отличие от Elasticsearch, у ClickHouse нет задержки видимости.

Дополнительно: при первом чтении можно пропустить

Глубже: pandas и Arrow вместо списковрасширенное

Когда события приходят таблицей — выгрузка из PostgreSQL, файл партнёра, результат обработки, — список списков лишний шаг. client.insert_df("order_events", df) отправляет DataFrame напрямую, сопоставляя колонки по именам, client.insert_arrow принимает таблицу Arrow без копирования в объекты Python, а query_df и query_arrow возвращают результат в том же виде. Это и быстрее, и экономнее по памяти на миллионах строк, но тянет за собой pandas или pyarrow в образ сервиса; для потока событий из Kafka списки строк остаются проще.

Глубже: когда ClickHouse недоступенрасширенное

Путь записи проектируют так, чтобы недоступность была задержкой, а не потерей: события ждут в Kafka или в outbox, потребитель повторяет пачку с паузой, метрики показывают возраст самого старого непрочитанного события и число упавших вставок. Буфер в памяти сервиса, если он есть, ограничен сверху и при переполнении отбрасывает события со счётчиком, а не растёт до OOM.

Путь чтения деградирует честно: send_receive_timeout в единицы секунд для витрин, ответ 503 для отчётного эндпоинта, кэш последнего успешного результата для витрин, которые смотрят часто. Проверка здоровья сервиса не включает ClickHouse в обязательные зависимости: оформление заказов не должно уходить из балансировщика из-за аналитики.

Коротко

  • Один клиент clickhouse-connect на процесс по HTTP, autogenerate_session_id=False для общего клиента (иначе код 373 Session is locked), send_receive_timeout в десятки секунд вместо 300, max_execution_time в настройках сессии.
  • Мелкие вставки запрещены: каждый INSERT это кусок на диске, порог 3000 кусков в партиции даёт ошибку 252.
  • Пачка через insert со списком строк и column_names, сброс по размеру и по времени под asyncio.Lock, обязательный flush при остановке.
  • Повтор пачки безопасен с insert_deduplication_token; у обычного MergeTree окно дедупликации нужно включить.
  • Асинхронная вставка копит на сервере (10 МиБ или 200 мс), wait_for_async_insert=1 возвращает ошибки записи, 0 их теряет.
  • В проде между сервисом и ClickHouse стоит Kafka: getmany, вставка, затем commit при выключенном автокоммите.
  • Чтение в отдельном репозитории под read-only пользователем, параметры {имя:Тип}, результат через named_results или query_df, выражения с AS.
  • Долгие отчёты считает фоновая задача, обработчик отдаёт готовое; таймаут клиента и max_execution_time сервера работают вместе.
  • Миграции отдельным каталогом SQL и своим раннером, без транзакций DDL, в кластере каждое выражение ON CLUSTER.
  • Тесты на контейнере через порт 8123: пачка, отчётный запрос на известных данных, прогон миграций.

Что почитать дальше