Выпуск Picopyn 3.0.0

Команда 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