Команда Picodata выпустила Picopyn 3.0.0 — Python-драйвер для Picodata с пулом соединений, автоматическим обнаружением узлов кластера и маршрутизацией запросов, учитывающей шардирование данных. Версия совместима с Picodata 26.1.X, как и предыдущие. Пакет можно установить из PyPI или исходного кода. Подробный changelog можно изучить в документации драйвера, в текущей же статье рассмотрим основные изменения с примерами.
Последние релизы (2.0.0 и
2.1.0) фокусировались на
работе с топологией кластера и маршрутизации запросов в асинхронном
драйвере. Главное в 3.0.0 — синхронный драйвер получил ровно те же
возможности: сбор топологии, приведение состава пула к ней, кеш
метаданных запросов и автоматическую маршрутизацию в execute. Кроме
того, появился явный API маршрутизации по ключу шардирования, а
маршрутизация запросов профилирована и оптимизирована.
Несовместимые изменения
Удалены Client и async Pool.connect()
Client был тонкой обёрткой над Pool, которая ничего не добавляла, и
устарел в 2.1.0. В
3.0.0 оба класса, синхронный и асинхронный, удалены. Вместе с ними
удалён устаревший синоним Pool.connect() у асинхронного пула — нужно
использовать Pool.open().
# picopyn < 3.0.0
from picopyn.asynchronous import Client as AsyncClient
client = AsyncClient(dsn, pool_size=10)
await client.connect()
# picopyn >= 3.0.0
from picopyn.asynchronous import Pool as AsyncPool
pool = AsyncPool(dsn, max_size=10, enable_discovery=True)
await pool.open()
При обновлении на новую версию драйвера важно учесть два отличия:
pool_sizeклиента — этоmax_sizeпула;- у клиента режим discovery был включён всегда, а пул по умолчанию
создаётся с
enable_discovery=False, поэтому прежнее поведение нужно запросить явно.
Для синхронного режима обновление выглядит так же: SyncPool(dsn, max_size=10, enable_discovery=True) и pool.open().
Синхронный API
Отслеживание топологии
Синхронный пул собирает топологию кластера в фоновом потоке: тиры, репликасеты, инстансы и карту бакетов. Сбор топологии кластера включён по умолчанию — аналогично поведению асинхронного пула. В настройках пула можно явно включить/выключить сбор топологии, а также изменить интервал обновления в секундах.
from picopyn.synchronous import Pool
pool = Pool(
dsn="postgresql://admin:pass@host1:5432,host2:5432",
enable_discovery=True,
max_size=10,
topology_update_interval=60,
)
pool.open()
print(pool.topology)
Приведение состава пула к топологии и перебалансировка
Синхронный пул теперь тоже подстраивает свой состав соединений под
кластер Picodata: соединения с выбывшими узлами вытесняются, пул
добирает новые соединения до своего максимального размера max_size, а
в конце цикла выполняется проход ребалансировки — соединение
переносится с самого нагруженного узла-кандидата на наименее
нагруженный. Как и в асинхронном пуле, число переносов за цикл
ограничено настраиваемым параметром rebalance_pool_divisor (по
умолчанию 10, то есть не более max(1, max_size // 10) переносов за
одно обновление топологии).
Кеш метаданных запроса и маршрутизация в execute
Синхронный пул держит собственное техническое соединение, на котором
получает от Picodata метаданные запроса (тир и параметры ключа
шардирования) и кеширует их. Размер кеша задаётся
query_metadata_cache_size.
На основе этого кеша и карты бакетов Pool.execute сам выбирает
соединение с мастером нужного репликасета. Условия те же, что и в
асинхронном драйвере:
- запрос параметризован;
- запрос выполняет операцию
INSERT; - вставляется одна строка данных;
- включены и сбор топологии, и сервис метаданных;
- запрос выполняется не в первый раз.
from picopyn.synchronous import Pool
pool = Pool(
dsn="postgresql://admin:pass@host1:5432,host2:5432",
enable_discovery=True,
max_size=10,
)
pool.open()
query = 'INSERT INTO "warehouse" VALUES (%s, %s)'
# первый вызов уходит на соединение, выбранное балансировкой,
# метаданные запрашиваются в фоне и не задерживают выполнение
pool.execute(query, (1, "first"))
# повторные вызовы того же запроса уходят сразу на мастер нужного репликасета
pool.execute(query, (2, "second"))
pool.close()
При устаревании метаданных запроса в кеше, Picodata отдаёт специальную ошибку, которую перехватывает драйвер. После этого он вытесняет запись из кеша, чтобы следующий вызов пересчитал метаданные.
Корректное завершение работы пула
Раньше синхронный close() закрывал только свободные соединения и
возвращался сразу, оставляя выданные пользователю соединения. Дождаться
их было нельзя, а если владелец соединения так и не возвращал его в пул
— закрыть их было нельзя вовсе.
Теперь close() дожидается, в том числе, и выданных соединений в
пределах timeout (по умолчанию час), после чего закрывает оставшиеся
принудительно; тот же таймаут ограничивает остановку трекера топологии.
Конкурентные вызовы close() присоединяются к уже идущему завершению.
pool.close(timeout=10)
Новый метод terminate() закрывает всё немедленно, не дожидаясь ничего,
и завершает уже запущенный close(). Трекер топологии и сервис
метаданных тоже получили terminate(), поэтому завершение работы не
зависает на обновлении топологии или на подготовке запроса: их прерывает
новый Connection.cancel().
Маршрутизация по ключу шардирования
Раньше отправить на нужный узел можно было либо параметризованный
однострочный INSERT (автоматически), либо проходить цикл маршрутизации,
вызывая все промежуточные функции вручную. В 3.0.0 появился явный API:
фабрика ключей шардирования для таблицы.
Pool.create_sharding_key_factory(table_name) один раз читает схему
распределения таблицы из системной таблицы
_pico_table,
после чего создание ключей выполняется локально, без обращений к базе.
Полученный
ShardingKey
содержит тир и bucket_id, по которым пул берёт соединение с текущим
мастером репликасета.
В примерах ниже используется таблица warehouse, распределённая по полю
item:
CREATE TABLE warehouse (
id INTEGER NOT NULL,
item TEXT NOT NULL,
type TEXT NOT NULL,
PRIMARY KEY (id))
USING memtx DISTRIBUTED BY (item);
Async:
factory = await pool.create_sharding_key_factory("warehouse")
key = factory.create_key({"item": "bricks"})
async with pool.acquire_by_sharding_key(key) as conn:
await conn.execute(
'SELECT * FROM "warehouse" WHERE item = $1',
"bricks",
)
Sync:
factory = pool.create_sharding_key_factory("warehouse")
key = factory.create_key({"item": "bricks"})
with pool.connection_by_sharding_key(key) as conn:
cursor = conn.cursor()
cursor.execute(
'SELECT * FROM "warehouse" WHERE item = %s',
("bricks",),
)
Особенности:
- в
create_keyможно передать всю строку целиком — поля, не входящие в ключ шардирования, игнорируются; - маршрутизация здесь явная: если топология не позволяет разрешить ключ или подходящее соединение не освободилось, метод выдаёт исключение, а не переходит молча на произвольное соединение;
- фабрика — это снимок схемы: после пересоздания таблицы или изменения схемы распределения её также нужно пересоздать;
- на текущий момент инвалидация схемы таблицы (пересоздание фабрики) является ответственностью пользователя.
Новые настройки пула
routing_acquire_timeout
Ранее, механизм shard-aware routing запрашивал соединение с мастером репликасета, в котором хранится бакет, и, если все существующие соединения с этим узлом были заняты/отсутствовали, мгновенно переходил на обычное соединение (то есть, то, которое пул отдавал без условий, по своей стратегии балансировки). Из-за этого под нагрузкой доля запросов, которые удалось направить на нужные узлы, могла значительно снижаться.
Новый параметр routing_acquire_timeout (float) задаёт, сколько секунд
механизм маршрутизации готов ждать освобождения соединения с таким
мастером, прежде чем взять обычное соединение из пула. По
умолчанию 0 — прежнее поведение, одна попытка без ожидания.
pool = Pool(
dsn="postgresql://admin:pass@host1:5432,host2:5432",
enable_discovery=True,
max_size=10,
routing_acquire_timeout=0.015,
)
Таймаут расходуется только тогда, когда соединения с нужным мастером в пуле есть, но все заняты. Если соединения с этим узлом нет вовсе (то есть ни среди свободных, ни среди занятых), механизм маршрутизации сразу выбирает обычное соединение.
Значение стоит выбирать исходя из нагрузки на кластер: если пул перегружен
и запросы и регулярно ждут свободного соединения, значение 0 избавит от
дополнительного ожидания.
Отключение сервиса метаданных
Пул всегда создавал сервис метаданных запросов, то есть открывал
техническое соединение к Picodata и держал кеш, даже когда shard-aware
routing
не был нужен. Теперь query_metadata_cache_size=None полностью
выключает сервис: ни кеша, ни технического соединения, а
get_query_metadata бросает RuntimeError. Значение меньше 1
отклоняется сразу при создании пула.
pool = Pool(
dsn="postgresql://admin:pass@host1:5432,host2:5432",
enable_discovery=True,
query_metadata_cache_size=None,
)
Планы на ближайшие релизы
- Сбор метрик драйвера
- Добавление возможности менять параметры пула (например, интервал обновления топологии или размер кеша метаданных) “на ходу”, а не только в момент запуска пула
- Поддержка Picodata 26.2.X
