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 исполняет в одной сессии только один запрос за раз: второй параллельный получает ошибку с кодом 373Session 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 создаёт кусок данных, и тысяча мелких вставок это тысяча кусков; ClickHouse рассчитан на пачки в тысячи строк.
Отсюда три способа доставки событий из Python-сервиса, и выбор между ними зависит от потока.
- Накопить в сервисе и отправить пачкой через
client.insert. Подходит фоновому воркеру или потребителю Kafka, у которого есть буфер. - Асинхронная вставка на стороне сервера: сервис шлёт по строке, а копит и пишет сам ClickHouse. Подходит, когда событий сотни в секунду и некуда их буферизовать.
- 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для общего клиента (иначе код 373Session 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: пачка, отчётный запрос на известных данных, прогон миграций.
Что почитать дальше
- Как устроен ClickHouse — колонки, куски данных и слияния, из которых растут все правила записи.
- Моделирование и запросы — ключ сортировки, партиции,
ReplacingMergeTreeи агрегатные функции. - ClickHouse в проде — репликация, TTL, бэкапы и что мониторить.
- PostgreSQL или ClickHouse — где проходит граница между транзакционной и аналитической нагрузкой.
- Kafka в production на Python — потребитель
aiokafka, подтверждение смещений и остановка. - Паттерны распределённых систем на Python — outbox как источник событий для аналитики.