Cassandra
Модуль предоставляет реализацию репозитория для базы данных Cassandra с использованием драйвера DataStax.
Cassandra — это распределенная колоночная база данных, где запросы пишутся на CQL, а модель данных обычно проектируется под конкретные сценарии чтения.
В Kora модуль Cassandra предоставляет декларативные репозитории поверх CqlSession: приложение пишет CQL-запросы в @Query, а Kora на этапе компиляции генерирует код подготовки запроса, связывания параметров и отображения результата.
Общие правила для отображений, @Repository, @Query, макросов, пакетных запросов и аннотаций @Table, @Column, @Id, @Embedded описаны в разделе общих правил работы с базами данных.
Этот документ охватывает специфичные для Cassandra части: подключение драйвера, конфигурацию CqlSession, профили выполнения, UDT, отображатели и поддерживаемые сигнатуры методов.
Пошаговый разбор перед справочным описанием смотрите в разделе База данных Cassandra.
Подключение¶
Зависимость build.gradle:
Модуль:
Зависимость build.gradle.kts:
Модуль:
Конфигурация¶
Конфигурация читается из секции cassandra и описывается интерфейсом CassandraConfig.
Как минимум необходимо указать basic.contactPoints. Остальные параметры необязательны или передаются драйверу только при явной настройке.
Пример простой конфигурации:
cassandra {
basic {
contactPoints = "127.0.0.1:9042, 127.0.0.2:9042" //(1)!
dc = "datacenter1" //(2)!
sessionKeyspace = "test-db" //(3)!
request {
timeout = "5s" //(4)!
}
}
auth {
login = "username" //(5)!
password = "password" //(6)!
}
}
- Адреса узлов
Cassandraдля подключения к базе данных (обязательно, без значения по умолчанию) - Имя датацентра
Cassandra(по умолчанию не указано, необязательно) - Имя
keyspaceдля подключения (по умолчанию не указано, необязательно) - Таймаут выполнения запроса для подключения (по умолчанию не указано, необязательно)
- Имя пользователя для подключения (по умолчанию не указано, необязательно)
- Пароль для подключения (по умолчанию не указано, необязательно)
cassandra:
basic:
contactPoints: "127.0.0.1:9042, 127.0.0.2:9042" #(1)!
dc: "datacenter1" #(2)!
sessionKeyspace: "test-db" #(3)!
request:
timeout: "5s" #(4)!
auth:
login: "username" #(5)!
password: "password" #(6)!
- Адреса узлов
Cassandraдля подключения к базе данных (обязательно, без значения по умолчанию) - Имя датацентра
Cassandra(по умолчанию не указано, необязательно) - Имя
keyspaceдля подключения (по умолчанию не указано, необязательно) - Таймаут выполнения запроса для подключения (по умолчанию не указано, необязательно)
- Имя пользователя для подключения (по умолчанию не указано, необязательно)
- Пароль для подключения (по умолчанию не указано, необязательно)
Пример полной конфигурации
Полная конфигурация с примерами значений. Описания параметров являются общими для примеров HOCON и YAML.
cassandra {
auth {
login = "username" //(1)!
password = "password" //(2)!
}
basic {
contactPoints = [ "127.0.0.1:9042", "127.0.0.2:9042" ] //(3)!
sessionName = "some-session-name" //(4)!
dc = "datacenter1" //(5)!
sessionKeyspace = "test-db" //(6)!
loadBalancingPolicy.slowReplicaAvoidance = true //(7)!
cloud.secureConnectBundle = "/location/of/secure/connect/bundle" //(8)!
request {
timeout = "5s" //(9)!
consistency = "LOCAL_ONE" //(10)!
pageSize = 5000 //(11)!
serialConsistency = "LOCAL_SERIAL" //(12)!
defaultIdempotence = false //(13)!
}
}
advanced {
sessionLeak.threshold = 4 //(14)!
connection {
connectTimeout = "10s" //(15)!
initQueryTimeout = "10s" //(16)!
setKeyspaceTimeout = "10s" //(17)!
maxRequestsPerConnection = 1024 //(18)!
maxOrphanRequests = 256 //(19)!
warnOnInitError = true //(20)!
pool {
localSize = 10 //(21)!
remoteSize = 10 //(22)!
}
}
reconnectOnInit = false //(23)!
reconnectionPolicy {
baseDelay = "1s" //(24)!
maxDelay = "60s" //(25)!
}
loadBalancingPolicy.dcFailover {
maxNodesPerRemoveDc = 1 //(26)!
allowForLocalConsistencyLevels = false //(27)!
}
sslEngineFactory {
cipherSuites = [ "TLS_RSA_WITH_AES_128_CBC_SHA", "TLS_RSA_WITH_AES_256_CBC_SHA" ] //(28)!
hostnameValidation = true //(29)!
keystorePath = "/path/to/client.keystore" //(30)!
keystorePassword = "password" //(31)!
truststorePath = "/path/to/client.truststore" //(32)!
truststorePassword = "password" //(33)!
}
timestampGenerator {
forceJavaClock = false //(34)!
driftWarning.threshold = "1s" //(35)!
driftWarning.interval = "10s" //(36)!
}
protocol {
version = "V4" //(37)!
compression = "lz4" //(38)!
maxFrameLength = 268435456 //(39)!
}
request {
warnIfSetKeyspace = true //(40)!
trace {
attempts = 5 //(41)!
interval = "1ms" //(42)!
consistency = "ONE" //(43)!
}
logWarnings = true //(44)!
}
metrics {
idGenerator {
name = "TaggingMetricIdGenerator" //(45)!
prefix = "my-app" //(46)!
}
node {
enabled = [ "bytes-sent", "bytes-received", "open-connections" ] //(47)!
cqlMessages {
lowestLatency = "1ms" //(48)!
highestLatency = "90s" //(49)!
significantDigits = 1 //(50)!
refreshInterval = "10s" //(51)!
slo = [ 1, 10, 50, 100, 200, 500, 1000 ] //(52)!
}
}
session {
enabled = [ "connected-nodes", "cql-requests", "cql-client-timeouts" ] //(53)!
cqlRequests {
lowestLatency = "1ms" //(54)!
highestLatency = "90s" //(55)!
significantDigits = 1 //(56)!
refreshInterval = "10s" //(57)!
slo = [ 1, 10, 50, 100, 200, 500, 1000 ] //(58)!
}
throttlingDelay {
lowestLatency = "1ms" //(59)!
highestLatency = "90s" //(60)!
significantDigits = 1 //(61)!
refreshInterval = "10s" //(62)!
slo = [ 1, 10, 50, 100, 200, 500, 1000 ] //(63)!
}
}
publishPercentileHistogram = false //(64)!
}
socket {
tcpNoDelay = true //(65)!
keepAlive = false //(66)!
reuseAddress = true //(67)!
lingerInterval = 0 //(68)!
receiveBufferSize = 65535 //(69)!
sendBufferSize = 65535 //(70)!
}
heartbeat {
interval = "30s" //(71)!
timeout = "2m" //(72)!
}
metadata {
schema {
enabled = true //(73)!
requestTimeout = "20s" //(74)!
requestPageSize = 20 //(75)!
refreshedKeyspaces = [ "ks1", "ks2" ] //(76)!
debouncer.window = "1s" //(77)!
debouncer.maxEvents = 20 //(78)!
}
topologyEventDebouncer.window = "1s" //(79)!
topologyEventDebouncer.maxEvents = 20 //(80)!
tokenMapEnabled = true //(81)!
}
controlConnection {
timeout = "10s" //(82)!
schemaAgreement {
interval = "200ms" //(83)!
timeout = "10s" //(84)!
warnOnFailure = true //(85)!
}
}
preparedStatements {
prepareOnAllNodes = true //(86)!
reprepareOnUp {
enabled = true //(87)!
checkSystemTable = false //(88)!
maxStatements = 0 //(89)!
maxParallelism = 100 //(90)!
timeout = "20s" //(91)!
}
preparedCache.weakValues = false //(92)!
}
netty {
ioGroup.size = 0 //(93)!
ioGroup.shutdown {
quietPeriod = 2 //(94)!
timeout = 15 //(95)!
unit = "SECONDS" //(96)!
}
adminGroup.size = 2 //(97)!
adminGroup.shutdown {
quietPeriod = 2 //(98)!
timeout = 15 //(99)!
unit = "SECONDS" //(100)!
}
timer.tickDuration = "100ms" //(101)!
timer.ticksPerWheel = 2048 //(102)!
daemon = false //(103)!
}
coalescer.rescheduleInterval = "10ms" //(104)!
resolveContactPoints = false //(105)!
throttler {
throttlerClass = "ConcurrencyLimitingRequestThrottler" //(106)!
maxConcurrentRequests = 1024 //(107)!
maxRequestsPerSecond = 10000 //(108)!
maxQueueSize = 10000 //(109)!
drainInterval = "1ms" //(110)!
}
}
profiles {
someProfile {
basic.request.timeout = "10s" //(111)!
basic.request.consistency = "LOCAL_QUORUM" //(112)!
advanced.request.trace.attempts = 3 //(113)!
advanced.request.trace.consistency = "ONE" //(114)!
}
}
telemetry {
logging.enabled = false //(115)!
metrics {
enabled = true //(116)!
slo = [ 1, 10, 50, 100, 200, 500, 1000 ] //(117)!
tags = { "key1" = "value1", "key2" = "value2" } //(118)!
}
tracing {
enabled = true //(119)!
attributes = { "key1" = "value1", "key2" = "value2" } //(120)!
}
}
}
- Имя пользователя для аутентификации в
Cassandra(по умолчанию не указано, необязательно). - Пароль для аутентификации в
Cassandra(по умолчанию не указано, необязательно). - Адреса узлов
Cassandraв форматеhost:port(обязательно, без значения по умолчанию). - Имя сессии драйвера, используемое в логах, метриках и диагностике (по умолчанию не указано, необязательно).
- Локальный датацентр для политики балансировки нагрузки (по умолчанию не указано, необязательно).
keyspace, который будет установлен для сессии после подключения (по умолчанию не указано, необязательно).- Включает избегание медленных реплик в стандартной политике балансировки нагрузки (по умолчанию не указано, необязательно).
- Путь или
URLкSecure Connect Bundleдля подключения кDataStax Astra/ облачной Cassandra (по умолчанию не указано, необязательно). - Обычный таймаут запроса (по умолчанию не указано, необязательно).
- Уровень согласованности обычного запроса, например
ONE,LOCAL_ONE,LOCAL_QUORUM,QUORUM,ALL(по умолчанию не указано, необязательно). - Размер страницы результата, то есть максимальное количество строк, запрашиваемых за один сетевой обмен (по умолчанию не указано, необязательно).
- Уровень последовательной согласованности для облегченных транзакций
LWT:SERIALилиLOCAL_SERIAL(по умолчанию не указано, необязательно). - Значение идемпотентности запроса по умолчанию; влияет на то, можно ли безопасно применять повторные попытки и спекулятивное выполнение (по умолчанию не указано, необязательно).
- Порог предупреждения об утечке сессии драйвера (по умолчанию не указано, необязательно).
- Таймаут открытия сетевого соединения с узлом (по умолчанию не указано, необязательно).
- Таймаут запросов, которые драйвер выполняет при инициализации соединения (по умолчанию не указано, необязательно).
- Таймаут установки
keyspaceна соединении (по умолчанию не указано, необязательно). - Максимальное количество одновременных запросов на одно соединение (по умолчанию не указано, необязательно).
- Максимальное количество запросов, ответ на которые уже не ожидается, но которые все еще могут завершиться внутри драйвера (по умолчанию не указано, необязательно).
- Логирует предупреждение при неудачной инициализации соединения для отдельного узла (по умолчанию не указано, необязательно).
- Размер пула соединений для узлов локального датацентра (по умолчанию не указано, необязательно).
- Размер пула соединений для удаленных узлов (по умолчанию не указано, необязательно).
- Разрешает повторную попытку инициализации, когда во время запуска все
contactPointsне отвечают (по умолчанию не указано, необязательно). - Начальная задержка политики переподключения (по умолчанию не указано, необязательно).
- Максимальная задержка политики переподключения (по умолчанию не указано, необязательно).
- Максимальное количество узлов удаленного датацентра, которые могут использоваться для отказоустойчивости (по умолчанию не указано, необязательно).
- Разрешает переключение на удаленный датацентр для локальных уровней согласованности (по умолчанию не указано, необязательно).
- Разрешенные наборы шифров для
SSL/TLS(по умолчанию не указано, необязательно). - Проверяет, что имя хоста узла соответствует сертификату
SSL/TLS(по умолчанию не указано, необязательно). - Путь к клиентскому keystore (по умолчанию не указано, необязательно).
- Пароль клиентского keystore (по умолчанию не указано, необязательно).
- Путь к truststore (по умолчанию не указано, необязательно).
- Пароль truststore (по умолчанию не указано, необязательно).
- Принудительно использует системные часы Java для генерации временных меток запросов (по умолчанию не указано, необязательно).
- Порог предупреждения о смещении временной метки в будущее (по умолчанию не указано, необязательно).
- Минимальный интервал между предупреждениями о смещении временной метки (по умолчанию не указано, необязательно).
- Версия бинарного протокола Cassandra, например
V4(по умолчанию не указано, необязательно). - Алгоритм сжатия протокола, например
lz4илиsnappy(по умолчанию не указано, необязательно). - Максимальный размер кадра протокола в байтах (по умолчанию не указано, необязательно).
- Логирует предупреждение, когда запрос явно меняет
keyspace(по умолчанию не указано, необязательно). - Количество попыток получить информацию трассировки запроса из Cassandra (по умолчанию не указано, необязательно).
- Интервал между попытками получить информацию трассировки запроса (по умолчанию не указано, необязательно).
- Уровень согласованности для запросов к таблицам трассировки (по умолчанию не указано, необязательно).
- Логирует предупреждения, возвращаемые Cassandra вместе с ответом на запрос (по умолчанию не указано, необязательно).
- Имя генератора идентификаторов метрик драйвера (по умолчанию:
TaggingMetricIdGenerator). - Префикс имен метрик драйвера (по умолчанию не указано, необязательно).
- Включенные метрики уровня узла (по умолчанию:
open-connections,in-flight,bytes-received,bytes-sent,write-timeouts,read-timeouts,aborted-requests). - Наименьшая ожидаемая задержка для гистограммы метрики
node.cqlMessages(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
node.cqlMessages(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
node.cqlMessages(по умолчанию не указано, необязательно). - Интервал обновления снимка для гистограммы метрики
node.cqlMessages(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиnode.cqlMessages(по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Включенные метрики уровня сессии (по умолчанию:
connected-nodes,cql-requests,cql-client-timeouts,cql-prepared-cache-size,throttling.delay,throttling.queue-size). - Наименьшая ожидаемая задержка для гистограммы метрики
session.cqlRequests(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
session.cqlRequests(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
session.cqlRequests(по умолчанию не указано, необязательно). - Интервал обновления снимка для гистограммы метрики
session.cqlRequests(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиsession.cqlRequests(по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Наименьшая ожидаемая задержка для гистограммы метрики
session.throttlingDelay(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
session.throttlingDelay(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
session.throttlingDelay(по умолчанию не указано, необязательно). - Интервал обновления снимка для гистограммы метрики
session.throttlingDelay(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиsession.throttlingDelay(по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Публикует процентильные гистограммы для метрик драйвера (по умолчанию:
false). - Включает
TCP_NODELAY, что отключает алгоритм Нейгла (по умолчанию не указано, необязательно). - Включает
SO_KEEPALIVEдля TCP-сокетов (по умолчанию не указано, необязательно). - Включает
SO_REUSEADDRдля TCP-сокетов (по умолчанию не указано, необязательно). - Значение
SO_LINGERдля TCP-сокетов (по умолчанию не указано, необязательно). - Размер буфера приема TCP-сокета в байтах (по умолчанию не указано, необязательно).
- Размер буфера отправки TCP-сокета в байтах (по умолчанию не указано, необязательно).
- Интервал отправки
heartbeatпо простаивающему соединению (по умолчанию не указано, необязательно). - Таймаут ожидания ответа на
heartbeat(по умолчанию не указано, необязательно). - Включает загрузку и обновление метаданных схемы (по умолчанию не указано, необязательно).
- Таймаут запросов метаданных схемы (по умолчанию не указано, необязательно).
- Размер страницы запросов метаданных схемы (по умолчанию не указано, необязательно).
- Список имен
keyspace, метаданные схемы которых обновляются драйвером (по умолчанию не указано, необязательно). - Окно для объединения событий обновления схемы перед обработкой (по умолчанию не указано, необязательно).
- Максимальное количество событий обновления схемы, которое может накопиться в окне (по умолчанию не указано, необязательно).
- Окно для объединения событий изменения топологии кластера (по умолчанию не указано, необязательно).
- Максимальное количество событий изменения топологии, которое может накопиться в окне (по умолчанию не указано, необязательно).
- Включает карту токенов для маршрутизации запросов к владельцам данных (по умолчанию не указано, необязательно).
- Таймаут служебного
control connection(по умолчанию не указано, необязательно). - Интервал проверки
schema agreementмежду узлами (по умолчанию не указано, необязательно). - Максимальное время ожидания
schema agreement(по умолчанию не указано, необязательно). - Логирует предупреждение, если
schema agreementне достигнута вовремя (по умолчанию не указано, необязательно). - Подготавливает запрос на всех узлах после того, как он успешно подготовлен на одном узле (по умолчанию не указано, необязательно).
- Повторно подготавливает запросы на узле, который снова стал доступен (по умолчанию не указано, необязательно).
- Проверяет системную таблицу
system.prepared_statementsперед повторной подготовкой запроса (по умолчанию не указано, необязательно). - Максимальное количество запросов для повторной подготовки;
0означает отсутствие ограничения на стороне драйвера (по умолчанию не указано, необязательно). - Максимальное количество параллельных запросов повторной подготовки (по умолчанию не указано, необязательно).
- Таймаут повторной подготовки запросов на одном узле (по умолчанию не указано, необязательно).
- Хранит значения кэша подготовленных запросов через слабые ссылки (по умолчанию не указано, необязательно).
- Количество потоков
Nettyдля сетевого ввода-вывода;0позволяет драйверу выбрать автоматически (по умолчанию не указано, необязательно). - Период затишья для плавной остановки
ioGroup(по умолчанию не указано, необязательно). - Максимальное время ожидания остановки
ioGroup(по умолчанию не указано, необязательно). - Единица измерения для параметров остановки
ioGroup(по умолчанию не указано, необязательно). - Количество потоков
Nettyдля административных задач драйвера (по умолчанию не указано, необязательно). - Период затишья для плавной остановки
adminGroup(по умолчанию не указано, необязательно). - Максимальное время ожидания остановки
adminGroup(по умолчанию не указано, необязательно). - Единица измерения для параметров остановки
adminGroup(по умолчанию не указано, необязательно). - Длительность одного тика таймера
Nettyдля отложенных задач драйвера (по умолчанию не указано, необязательно). - Количество тиков в колесе таймера
Netty(по умолчанию не указано, необязательно). - Делает потоки
Nettyдемон-потоками (по умолчанию не указано, необязательно). - Интервал перепланирования для объединения сообщений перед отправкой (по умолчанию не указано, необязательно).
- Разрешает драйверу разрешать
contactPointsчерез DNS во время запуска (по умолчанию не указано, необязательно). - Класс ограничителя запросов драйвера (по умолчанию не указано, необязательно).
- Максимальное количество одновременных запросов для ограничителя (по умолчанию не указано, необязательно).
- Максимальное количество запросов в секунду для ограничителя (по умолчанию не указано, необязательно).
- Максимальный размер очереди запросов ограничителя (по умолчанию не указано, необязательно).
- Интервал, с которым ограничитель освобождает запросы из очереди (по умолчанию не указано, необязательно).
- Переопределение
basic.request.timeoutдля профиляsomeProfile(по умолчанию не указано, необязательно). - Переопределение
basic.request.consistencyдля профиляsomeProfile(по умолчанию не указано, необязательно). - Переопределение
advanced.request.trace.attemptsдля профиляsomeProfile(по умолчанию не указано, необязательно). - Переопределение
advanced.request.trace.consistencyдля профиляsomeProfile(по умолчанию не указано, необязательно). - Включает логирование запросов Kora (по умолчанию:
false). - Включает метрики запросов Kora (по умолчанию:
true). - Границы
SLOметрик Kora (по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Дополнительные теги метрик Kora (по умолчанию:
{}). - Включает трассировку запросов Kora (по умолчанию:
true). - Дополнительные атрибуты трассировки Kora (по умолчанию:
{}).
cassandra:
auth:
login: "username" #(1)!
password: "password" #(2)!
basic:
contactPoints: [ "127.0.0.1:9042", "127.0.0.2:9042" ] #(3)!
sessionName: "some-session-name" #(4)!
dc: "datacenter1" #(5)!
sessionKeyspace: "test-db" #(6)!
loadBalancingPolicy:
slowReplicaAvoidance: true #(7)!
cloud:
secureConnectBundle: "/location/of/secure/connect/bundle" #(8)!
request:
timeout: "5s" #(9)!
consistency: "LOCAL_ONE" #(10)!
pageSize: 5000 #(11)!
serialConsistency: "LOCAL_SERIAL" #(12)!
defaultIdempotence: false #(13)!
advanced:
sessionLeak:
threshold: 4 #(14)!
connection:
connectTimeout: "10s" #(15)!
initQueryTimeout: "10s" #(16)!
setKeyspaceTimeout: "10s" #(17)!
maxRequestsPerConnection: 1024 #(18)!
maxOrphanRequests: 256 #(19)!
warnOnInitError: true #(20)!
pool:
localSize: 10 #(21)!
remoteSize: 10 #(22)!
reconnectOnInit: false #(23)!
reconnectionPolicy:
baseDelay: "1s" #(24)!
maxDelay: "60s" #(25)!
loadBalancingPolicy:
dcFailover:
maxNodesPerRemoveDc: 1 #(26)!
allowForLocalConsistencyLevels: false #(27)!
sslEngineFactory:
cipherSuites: [ "TLS_RSA_WITH_AES_128_CBC_SHA", "TLS_RSA_WITH_AES_256_CBC_SHA" ] #(28)!
hostnameValidation: true #(29)!
keystorePath: "/path/to/client.keystore" #(30)!
keystorePassword: "password" #(31)!
truststorePath: "/path/to/client.truststore" #(32)!
truststorePassword: "password" #(33)!
timestampGenerator:
forceJavaClock: false #(34)!
driftWarning:
threshold: "1s" #(35)!
interval: "10s" #(36)!
protocol:
version: "V4" #(37)!
compression: "lz4" #(38)!
maxFrameLength: 268435456 #(39)!
request:
warnIfSetKeyspace: true #(40)!
trace:
attempts: 5 #(41)!
interval: "1ms" #(42)!
consistency: "ONE" #(43)!
logWarnings: true #(44)!
metrics:
idGenerator:
name: "TaggingMetricIdGenerator" #(45)!
prefix: "my-app" #(46)!
node:
enabled: [ "bytes-sent", "bytes-received", "open-connections" ] #(47)!
cqlMessages:
lowestLatency: "1ms" #(48)!
highestLatency: "90s" #(49)!
significantDigits: 1 #(50)!
refreshInterval: "10s" #(51)!
slo: [ 1, 10, 50, 100, 200, 500, 1000 ] #(52)!
session:
enabled: [ "connected-nodes", "cql-requests", "cql-client-timeouts" ] #(53)!
cqlRequests:
lowestLatency: "1ms" #(54)!
highestLatency: "90s" #(55)!
significantDigits: 1 #(56)!
refreshInterval: "10s" #(57)!
slo: [ 1, 10, 50, 100, 200, 500, 1000 ] #(58)!
throttlingDelay:
lowestLatency: "1ms" #(59)!
highestLatency: "90s" #(60)!
significantDigits: 1 #(61)!
refreshInterval: "10s" #(62)!
slo: [ 1, 10, 50, 100, 200, 500, 1000 ] #(63)!
publishPercentileHistogram: false #(64)!
socket:
tcpNoDelay: true #(65)!
keepAlive: false #(66)!
reuseAddress: true #(67)!
lingerInterval: 0 #(68)!
receiveBufferSize: 65535 #(69)!
sendBufferSize: 65535 #(70)!
heartbeat:
interval: "30s" #(71)!
timeout: "2m" #(72)!
metadata:
schema:
enabled: true #(73)!
requestTimeout: "20s" #(74)!
requestPageSize: 20 #(75)!
refreshedKeyspaces: [ "ks1", "ks2" ] #(76)!
debouncer:
window: "1s" #(77)!
maxEvents: 20 #(78)!
topologyEventDebouncer:
window: "1s" #(79)!
maxEvents: 20 #(80)!
tokenMapEnabled: true #(81)!
controlConnection:
timeout: "10s" #(82)!
schemaAgreement:
interval: "200ms" #(83)!
timeout: "10s" #(84)!
warnOnFailure: true #(85)!
preparedStatements:
prepareOnAllNodes: true #(86)!
reprepareOnUp:
enabled: true #(87)!
checkSystemTable: false #(88)!
maxStatements: 0 #(89)!
maxParallelism: 100 #(90)!
timeout: "20s" #(91)!
preparedCache:
weakValues: false #(92)!
netty:
ioGroup:
size: 0 #(93)!
shutdown:
quietPeriod: 2 #(94)!
timeout: 15 #(95)!
unit: "SECONDS" #(96)!
adminGroup:
size: 2 #(97)!
shutdown:
quietPeriod: 2 #(98)!
timeout: 15 #(99)!
unit: "SECONDS" #(100)!
timer:
tickDuration: "100ms" #(101)!
ticksPerWheel: 2048 #(102)!
daemon: false #(103)!
coalescer:
rescheduleInterval: "10ms" #(104)!
resolveContactPoints: false #(105)!
throttler:
throttlerClass: "ConcurrencyLimitingRequestThrottler" #(106)!
maxConcurrentRequests: 1024 #(107)!
maxRequestsPerSecond: 10000 #(108)!
maxQueueSize: 10000 #(109)!
drainInterval: "1ms" #(110)!
profiles:
someProfile:
basic:
request:
timeout: "10s" #(111)!
consistency: "LOCAL_QUORUM" #(112)!
advanced:
request:
trace:
attempts: 3 #(113)!
consistency: "ONE" #(114)!
telemetry:
logging:
enabled: false #(115)!
metrics:
enabled: true #(116)!
slo: [ 1, 10, 50, 100, 200, 500, 1000 ] #(117)!
tags: { key1: "value1", key2: "value2" } #(118)!
tracing:
enabled: true #(119)!
attributes: { key1: "value1", key2: "value2" } #(120)!
- Имя пользователя для аутентификации в
Cassandra(по умолчанию не указано, необязательно). - Пароль для аутентификации в
Cassandra(по умолчанию не указано, необязательно). - Адреса узлов
Cassandraв форматеhost:port(обязательно, без значения по умолчанию). - Имя сессии драйвера, используемое в логах, метриках и диагностике (по умолчанию не указано, необязательно).
- Локальный датацентр для политики балансировки нагрузки (по умолчанию не указано, необязательно).
keyspace, который будет установлен для сессии после подключения (по умолчанию не указано, необязательно).- Включает избегание медленных реплик в стандартной политике балансировки нагрузки (по умолчанию не указано, необязательно).
- Путь или
URLкSecure Connect Bundleдля подключения кDataStax Astra/ облачной Cassandra (по умолчанию не указано, необязательно). - Обычный таймаут запроса (по умолчанию не указано, необязательно).
- Уровень согласованности обычного запроса, например
ONE,LOCAL_ONE,LOCAL_QUORUM,QUORUM,ALL(по умолчанию не указано, необязательно). - Размер страницы результата, то есть максимальное количество строк, запрашиваемых за один сетевой обмен (по умолчанию не указано, необязательно).
- Уровень последовательной согласованности для облегченных транзакций
LWT:SERIALилиLOCAL_SERIAL(по умолчанию не указано, необязательно). - Значение идемпотентности запроса по умолчанию; влияет на то, можно ли безопасно применять повторные попытки и спекулятивное выполнение (по умолчанию не указано, необязательно).
- Порог предупреждения об утечке сессии драйвера (по умолчанию не указано, необязательно).
- Таймаут открытия сетевого соединения с узлом (по умолчанию не указано, необязательно).
- Таймаут запросов, которые драйвер выполняет при инициализации соединения (по умолчанию не указано, необязательно).
- Таймаут установки
keyspaceна соединении (по умолчанию не указано, необязательно). - Максимальное количество одновременных запросов на одно соединение (по умолчанию не указано, необязательно).
- Максимальное количество запросов, ответ на которые уже не ожидается, но которые все еще могут завершиться внутри драйвера (по умолчанию не указано, необязательно).
- Логирует предупреждение при неудачной инициализации соединения для отдельного узла (по умолчанию не указано, необязательно).
- Размер пула соединений для узлов локального датацентра (по умолчанию не указано, необязательно).
- Размер пула соединений для удаленных узлов (по умолчанию не указано, необязательно).
- Разрешает повторную попытку инициализации, когда во время запуска все
contactPointsне отвечают (по умолчанию не указано, необязательно). - Начальная задержка политики переподключения (по умолчанию не указано, необязательно).
- Максимальная задержка политики переподключения (по умолчанию не указано, необязательно).
- Максимальное количество узлов удаленного датацентра, которые могут использоваться для отказоустойчивости (по умолчанию не указано, необязательно).
- Разрешает переключение на удаленный датацентр для локальных уровней согласованности (по умолчанию не указано, необязательно).
- Разрешенные наборы шифров для
SSL/TLS(по умолчанию не указано, необязательно). - Проверяет, что имя хоста узла соответствует сертификату
SSL/TLS(по умолчанию не указано, необязательно). - Путь к клиентскому keystore (по умолчанию не указано, необязательно).
- Пароль клиентского keystore (по умолчанию не указано, необязательно).
- Путь к truststore (по умолчанию не указано, необязательно).
- Пароль truststore (по умолчанию не указано, необязательно).
- Принудительно использует системные часы Java для генерации временных меток запросов (по умолчанию не указано, необязательно).
- Порог предупреждения о смещении временной метки в будущее (по умолчанию не указано, необязательно).
- Минимальный интервал между предупреждениями о смещении временной метки (по умолчанию не указано, необязательно).
- Версия бинарного протокола Cassandra, например
V4(по умолчанию не указано, необязательно). - Алгоритм сжатия протокола, например
lz4илиsnappy(по умолчанию не указано, необязательно). - Максимальный размер кадра протокола в байтах (по умолчанию не указано, необязательно).
- Логирует предупреждение, когда запрос явно меняет
keyspace(по умолчанию не указано, необязательно). - Количество попыток получить информацию трассировки запроса из Cassandra (по умолчанию не указано, необязательно).
- Интервал между попытками получить информацию трассировки запроса (по умолчанию не указано, необязательно).
- Уровень согласованности для запросов к таблицам трассировки (по умолчанию не указано, необязательно).
- Логирует предупреждения, возвращаемые Cassandra вместе с ответом на запрос (по умолчанию не указано, необязательно).
- Имя генератора идентификаторов метрик драйвера (по умолчанию:
TaggingMetricIdGenerator). - Префикс имен метрик драйвера (по умолчанию не указано, необязательно).
- Включенные метрики уровня узла (по умолчанию:
open-connections,in-flight,bytes-received,bytes-sent,write-timeouts,read-timeouts,aborted-requests). - Наименьшая ожидаемая задержка для гистограммы метрики
node.cqlMessages(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
node.cqlMessages(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
node.cqlMessages(по умолчанию не указано, необязательно). - Интервал обновления снимка для гистограммы метрики
node.cqlMessages(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиnode.cqlMessages(по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Включенные метрики уровня сессии (по умолчанию:
connected-nodes,cql-requests,cql-client-timeouts,cql-prepared-cache-size,throttling.delay,throttling.queue-size). - Наименьшая ожидаемая задержка для гистограммы метрики
session.cqlRequests(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
session.cqlRequests(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
session.cqlRequests(по умолчанию не указано, необязательно). - Интервал обновления снимка для гистограммы метрики
session.cqlRequests(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиsession.cqlRequests(по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Наименьшая ожидаемая задержка для гистограммы метрики
session.throttlingDelay(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
session.throttlingDelay(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
session.throttlingDelay(по умолчанию не указано, необязательно). - Интервал обновления снимка для гистограммы метрики
session.throttlingDelay(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиsession.throttlingDelay(по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Публикует процентильные гистограммы для метрик драйвера (по умолчанию:
false). - Включает
TCP_NODELAY, что отключает алгоритм Нейгла (по умолчанию не указано, необязательно). - Включает
SO_KEEPALIVEдля TCP-сокетов (по умолчанию не указано, необязательно). - Включает
SO_REUSEADDRдля TCP-сокетов (по умолчанию не указано, необязательно). - Значение
SO_LINGERдля TCP-сокетов (по умолчанию не указано, необязательно). - Размер буфера приема TCP-сокета в байтах (по умолчанию не указано, необязательно).
- Размер буфера отправки TCP-сокета в байтах (по умолчанию не указано, необязательно).
- Интервал отправки
heartbeatпо простаивающему соединению (по умолчанию не указано, необязательно). - Таймаут ожидания ответа на
heartbeat(по умолчанию не указано, необязательно). - Включает загрузку и обновление метаданных схемы (по умолчанию не указано, необязательно).
- Таймаут запросов метаданных схемы (по умолчанию не указано, необязательно).
- Размер страницы запросов метаданных схемы (по умолчанию не указано, необязательно).
- Список имен
keyspace, метаданные схемы которых обновляются драйвером (по умолчанию не указано, необязательно). - Окно для объединения событий обновления схемы перед обработкой (по умолчанию не указано, необязательно).
- Максимальное количество событий обновления схемы, которое может накопиться в окне (по умолчанию не указано, необязательно).
- Окно для объединения событий изменения топологии кластера (по умолчанию не указано, необязательно).
- Максимальное количество событий изменения топологии, которое может накопиться в окне (по умолчанию не указано, необязательно).
- Включает карту токенов для маршрутизации запросов к владельцам данных (по умолчанию не указано, необязательно).
- Таймаут служебного
control connection(по умолчанию не указано, необязательно). - Интервал проверки
schema agreementмежду узлами (по умолчанию не указано, необязательно). - Максимальное время ожидания
schema agreement(по умолчанию не указано, необязательно). - Логирует предупреждение, если
schema agreementне достигнута вовремя (по умолчанию не указано, необязательно). - Подготавливает запрос на всех узлах после того, как он успешно подготовлен на одном узле (по умолчанию не указано, необязательно).
- Повторно подготавливает запросы на узле, который снова стал доступен (по умолчанию не указано, необязательно).
- Проверяет системную таблицу
system.prepared_statementsперед повторной подготовкой запроса (по умолчанию не указано, необязательно). - Максимальное количество запросов для повторной подготовки;
0означает отсутствие ограничения на стороне драйвера (по умолчанию не указано, необязательно). - Максимальное количество параллельных запросов повторной подготовки (по умолчанию не указано, необязательно).
- Таймаут повторной подготовки запросов на одном узле (по умолчанию не указано, необязательно).
- Хранит значения кэша подготовленных запросов через слабые ссылки (по умолчанию не указано, необязательно).
- Количество потоков
Nettyдля сетевого ввода-вывода;0позволяет драйверу выбрать автоматически (по умолчанию не указано, необязательно). - Период затишья для плавной остановки
ioGroup(по умолчанию не указано, необязательно). - Максимальное время ожидания остановки
ioGroup(по умолчанию не указано, необязательно). - Единица измерения для параметров остановки
ioGroup(по умолчанию не указано, необязательно). - Количество потоков
Nettyдля административных задач драйвера (по умолчанию не указано, необязательно). - Период затишья для плавной остановки
adminGroup(по умолчанию не указано, необязательно). - Максимальное время ожидания остановки
adminGroup(по умолчанию не указано, необязательно). - Единица измерения для параметров остановки
adminGroup(по умолчанию не указано, необязательно). - Длительность одного тика таймера
Nettyдля отложенных задач драйвера (по умолчанию не указано, необязательно). - Количество тиков в колесе таймера
Netty(по умолчанию не указано, необязательно). - Делает потоки
Nettyдемон-потоками (по умолчанию не указано, необязательно). - Интервал перепланирования для объединения сообщений перед отправкой (по умолчанию не указано, необязательно).
- Разрешает драйверу разрешать
contactPointsчерез DNS во время запуска (по умолчанию не указано, необязательно). - Класс ограничителя запросов драйвера (по умолчанию не указано, необязательно).
- Максимальное количество одновременных запросов для ограничителя (по умолчанию не указано, необязательно).
- Максимальное количество запросов в секунду для ограничителя (по умолчанию не указано, необязательно).
- Максимальный размер очереди запросов ограничителя (по умолчанию не указано, необязательно).
- Интервал, с которым ограничитель освобождает запросы из очереди (по умолчанию не указано, необязательно).
- Переопределение
basic.request.timeoutдля профиляsomeProfile(по умолчанию не указано, необязательно). - Переопределение
basic.request.consistencyдля профиляsomeProfile(по умолчанию не указано, необязательно). - Переопределение
advanced.request.trace.attemptsдля профиляsomeProfile(по умолчанию не указано, необязательно). - Переопределение
advanced.request.trace.consistencyдля профиляsomeProfile(по умолчанию не указано, необязательно). - Включает логирование запросов Kora (по умолчанию:
false). - Включает метрики запросов Kora (по умолчанию:
true). - Границы
SLOметрик Kora (по умолчанию:ru.tinkoff.kora.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Дополнительные теги метрик Kora (по умолчанию:
{}). - Включает трассировку запросов Kora (по умолчанию:
true). - Дополнительные атрибуты трассировки Kora (по умолчанию:
{}).
Конфигурация в коде¶
Драйвер можно настроить вручную в коде, зарегистрировав компонент CassandraConfigurer.
Метод configure получает CqlSessionBuilder и ProgrammaticDriverConfigLoaderBuilder,
поэтому вы можете настроить построитель сессии и переопределить низкоуровневые параметры драйвера, которые не доступны через секцию конфигурации cassandra:
Использование¶
Чтобы создать репозиторий, объявите интерфейс с @Repository и унаследуйте CassandraRepository.
Такой репозиторий получает доступ к CqlSession через сгенерированный код и использует @Query для выполнения CQL-запросов.
Параметры запроса связываются по имени: :id, :entity.field, :filter.value.
Отображения описываются с помощью общих аннотаций баз данных и помечаются @EntityCassandra,
чтобы Kora сгенерировала отображатель на этапе компиляции (см. Отображение):
@Repository
public interface EntityRepository extends CassandraRepository {
@EntityCassandra
@Table("entities")
record Entity(@Id String id,
@Column("value1") int field1,
String value2,
@Nullable String value3) {}
@Query("SELECT %{return#selects} FROM %{return#table} WHERE id = :id") //(1)!
@Nullable
Entity findById(String id);
@Query("SELECT id, value1, value2, value3 FROM entities") //(2)!
List<Entity> findAll();
@Query("INSERT INTO %{entity#inserts}") //(3)!
void insert(Entity entity);
}
- Использует макрос
%{return#selects}и%{return#table}. Разворачивается в запрос: Метод использует макросы дляSELECT. Подробнее: Общие правила работы с базами данных — Макросы - Поля перечислены вручную без использования макросов — это допустимо, но требует поддержки при изменении отображения.
- Использует макрос
%{entity#inserts}. Разворачивается в запрос:Метод использует макросы дляINSERT INTO entities(id, value1, value2, value3) VALUES(:entity.id, :entity.value1, :entity.value2, :entity.value3)INSERT. Подробнее: Общие правила работы с базами данных — Макросы
@Repository
interface EntityRepository : CassandraRepository {
@EntityCassandra
@Table("entities")
data class Entity(
@field:Id val id: String,
@field:Column("value1") val field1: Int,
val value2: String,
val value3: String?
)
@Query("SELECT %{return#selects} FROM %{return#table} WHERE id = :id") //(1)!
fun findById(id: String): Entity?
@Query("INSERT INTO %{entity#inserts}") //(3)!
fun insert(entity: Entity)
}
- Использует макрос
%{return#selects}и%{return#table}. Разворачивается в запрос: Метод использует макросы дляSELECT. Подробнее: Общие правила работы с базами данных — Макросы - Поля перечислены вручную без использования макросов — это допустимо, но требует поддержки при изменении отображения.
- Использует макрос
%{entity#inserts}. Разворачивается в запрос:Метод использует макросы дляINSERT INTO entities(id, value1, value2, value3) VALUES(:entity.id, :entity.value1, :entity.value2, :entity.value3)INSERT. Подробнее: Общие правила работы с базами данных — Макросы
CQL остается под контролем разработчика: вы сами пишете текст запроса, тогда как Kora берет на себя только связывание параметров,
выполнение запроса и отображение результата.
Общие правила для отображений, @Table, @Column, @Id, @Embedded, @Batch и макросов описаны в разделе
Общие правила работы с базами данных.
Связывание параметров: Kora выполняет типизированное внедрение аргументов в CQL-запрос на этапе компиляции.
Параметры запроса (например, :id, :entity.field1) заменяются в сгенерированном коде на соответствующие вызовы драйвера Cassandra.
Например, для параметра String id будет сгенерировано что-то вроде statement.setString(1, id), где индекс соответствует порядку параметра в запросе.
Это обеспечивает безопасность (защита от CQL-инъекций) и производительность (использование подготовленных запросов драйвером).
В отличие от реляционных баз данных, в Cassandra нет транзакций.
Когда нужно, чтобы несколько операторов применились атомарно, используйте @Batch-метод (CQL BATCH), как показано выше;
его семантика и макросы описаны в разделе общих правил работы с базами данных.
Профиль¶
Можно переопределить общие настройки частными настройками из профиля. Предположим, есть такая конфигурация профиля someProfile:
Чтобы применить настройки из профиля someProfile, достаточно сделать следующее:
Настройки, указанные в профиле, будут применяться к каждому запросу, в частности в данном случае будет установлен таймаут 10s.
Профиль применяется только к методу, помеченному @CassandraProfile; остальные методы репозитория продолжают использовать базовую конфигурацию.
Отображение¶
Можно переопределить отображение различных частей отображения и параметров запроса — Kora предоставляет для этого специальные интерфейсы.
Из коробки CassandraModule предоставляет отображатели для распространенных типов: String, числовых типов, Boolean, BigDecimal, BigInteger, UUID, ByteBuffer, LocalDate, LocalTime, LocalDateTime, ZonedDateTime и Instant.
Если тип не входит в этот набор или ему нужно особое представление в CQL, добавьте собственный отображатель через @Mapping.
Отображение¶
Используйте аннотацию @EntityCassandra для оптимального отображения.
Аннотация позволяет обработчику аннотаций сгенерировать все необходимые отображатели за один раунд аннотационной обработки.
Без этой аннотации отображатели генерируются по требованию, что может потребовать множества раундов обработки и значительно увеличить время компиляции.
Это рекомендуемый способ отображения каждого типа, возвращаемого из репозитория или связываемого в нем.
Ожидается, что все вложенные отображения и типы UDT также используют эту аннотацию.
Результат¶
Если нужно вручную преобразовать весь результат синхронного запроса, используйте CassandraResultSetMapper<T>.
Он получает ResultSet и возвращает значение метода репозитория: одиночный объект, список, Optional<T> или другой поддерживаемый тип.
final class ResultMapper implements CassandraResultSetMapper<List<UUID>> {
@Override
public List<UUID> apply(ResultSet rows) {
var result = new ArrayList<UUID>();
for (var row : rows) {
result.add(row.getUuid("id"));
}
return result;
}
}
@Repository
public interface EntityRepository extends CassandraRepository {
@Mapping(ResultMapper.class)
@Query("SELECT id FROM entities")
List<UUID> getIds();
}
В Kotlin отображатели нужно писать только для типов T?, поэтому в интерфейсах тип указывается как @Nullable.
class ResultMapper : CassandraResultSetMapper<List<UUID>> {
override fun apply(rows: ResultSet): List<UUID> {
return rows.map { it.getUuid("id") }
}
}
@Repository
interface EntityRepository : CassandraRepository {
@Mapping(ResultMapper::class)
@Query("SELECT id FROM entities")
fun getIds(): List<UUID>
}
Каждый интерфейс отображателя результата также предоставляет статические фабричные методы, которые строят полный отображатель результата из CassandraRowMapper<T>,
так что один отображатель строки можно переиспользовать в разных сигнатурах:
CassandraResultSetMapper—singleResultSetMapper,optionalResultSetMapper,listResultSetMapper;CassandraAsyncResultSetMapper—one,list(автоматически проходит по страницам результата);CassandraReactiveResultSetMapper—flux,mono,monoVoid,monoList.
Строка¶
Если нужно вручную преобразовать одну строку результата, используйте CassandraRowMapper<T>.
Этот отображатель применяется к каждой строке и подходит для возвращаемых значений вида T, Optional<T>, List<T>, Flux<T> и Flow<T>.
final class RowMapper implements CassandraRowMapper<UUID> {
@Override
public UUID apply(Row row) {
return UUID.fromString(row.getString(0));
}
}
@Repository
public interface EntityRepository extends CassandraRepository {
@Mapping(RowMapper.class)
@Query("SELECT id FROM entities")
List<UUID> findAll();
}
В Kotlin отображатели нужно писать только для типов T?, поэтому в интерфейсах тип указывается как @Nullable.
Столбец¶
Если нужно вручную преобразовать значение столбца, предлагается использовать CassandraRowColumnMapper:
public final class ColumnMapper implements CassandraRowColumnMapper<UUID> {
@Override
public UUID apply(GettableByName row, int index) {
return UUID.fromString(row.getString(index));
}
}
@Table("entities")
public record Entity(@Mapping(ColumnMapper.class) @Id UUID id, String name) { }
@Repository
public interface EntityRepository extends CassandraRepository {
@Query("SELECT id, name FROM entities")
List<Entity> findAll();
}
В Kotlin отображатели нужно писать только для типов T?, поэтому в интерфейсах тип указывается как @Nullable.
class ColumnMapper : CassandraRowColumnMapper<UUID> {
override fun apply(row: GettableByName, index: Int): UUID {
return UUID.fromString(row.getString(index))
}
}
@Table("entities")
data class Entity(
@Id @Mapping(ColumnMapper::class) val id: UUID,
val name: String
)
@Repository
interface EntityRepository : CassandraRepository {
@Query("SELECT id, name FROM entities")
fun findAll(): List<Entity>
}
Параметр¶
Если нужно вручную преобразовать значение параметра запроса, используйте CassandraParameterColumnMapper<T>.
Он получает SettableByName<?>, индекс параметра и значение из метода репозитория.
public final class ParameterMapper implements CassandraParameterColumnMapper<UUID> {
@Override
public void apply(SettableByName<?> stmt, int index, @Nullable UUID value) {
if (value != null) {
stmt.setString(index, value.toString());
}
}
}
@Repository
public interface EntityRepository extends CassandraRepository {
@Query("SELECT id, name FROM entities WHERE id = :id")
List<Entity> findById(@Mapping(ParameterMapper.class) UUID id);
}
В Kotlin отображатели нужно писать только для типов T?, поэтому в интерфейсах тип указывается как @Nullable.
class ParameterMapper : CassandraParameterColumnMapper<UUID?> {
override fun apply(stmt: SettableByName<*>, index: Int, value: UUID?) {
if (value != null) {
stmt.setString(index, value.toString())
}
}
}
@Repository
interface EntityRepository : CassandraRepository {
@Query("SELECT id, name FROM entities WHERE id = :id")
fun findById(@Mapping(ParameterMapper::class) id: UUID): List<Entity>
}
Асинхронный¶
Для CompletionStage<T> и CompletableFuture<T> используйте CassandraAsyncResultSetMapper<T>, который получает AsyncResultSet и возвращает CompletionStage<T>.
Его метод list автоматически запрашивает последующие страницы результата, поэтому результат List<T> собирает все страницы перед завершением.
Для реактивных типов Mono<T> / Flux<T> используйте CassandraReactiveResultSetMapper<T, P>, который получает ReactiveResultSet и возвращает нужный Publisher.
final class ReactiveResultMapper implements CassandraReactiveResultSetMapper<UUID, Flux<UUID>> {
@Override
public Flux<UUID> apply(ReactiveResultSet rows) {
return Flux.from(rows).map(r -> UUID.fromString(r.getString(0)));
}
}
@Repository
public interface EntityRepository extends CassandraRepository {
@Mapping(ReactiveResultMapper.class)
@Query("SELECT id FROM entities")
Flux<UUID> getIds();
}
class ReactiveResultMapper : CassandraReactiveResultSetMapper<UUID, Flux<UUID>> {
override fun apply(rows: ReactiveResultSet): Flux<UUID> {
return Flux.from(rows).map { r -> UUID.fromString(r.getString(0)) }
}
}
@Repository
interface EntityRepository : CassandraRepository {
@Mapping(ReactiveResultMapper::class)
@Query("SELECT id FROM entities")
fun getIds(): Flux<UUID>
}
Ручной запрос¶
Если запрос сложно выразить одной статической @Query, вы можете объявить обычный метод с реализацией и построить CQL вручную.
Репозиторий предоставляет getCassandraConnectionFactory(), а CassandraConnectionFactory#query выполняет такой запрос:
он подготавливает оператор через текущую CqlSession, оборачивает выполнение в телеметрию Kora и возвращает значение, полученное из колбэка.
Метод доступа currentSession() возвращает активную CqlSession, а telemetry() возвращает DataBaseTelemetry, используемую для отчетности.
QueryContext содержит идентификатор запроса и итоговый CQL.
Идентификатор передается в телеметрию, поэтому используйте стабильное имя, например Repository.method.
Связывайте значения через BoundStatement, полученный из подготовленного оператора; никогда не конкатенируйте значения напрямую в строку запроса.
@Repository
public interface EntityRepository extends CassandraRepository {
default List<Entity> findByFilter(@Nullable String value2) {
var sql = new StringBuilder("SELECT id, value1, value2, value3 FROM entities");
if (value2 != null) {
sql.append(" WHERE value2 = ? ALLOW FILTERING");
}
var connectionFactory = getCassandraConnectionFactory();
var queryContext = new QueryContext("EntityRepository.findByFilter", sql.toString());
return connectionFactory.query(queryContext, statement -> {
var boundStatement = (value2 != null)
? statement.bind(value2)
: statement.bind();
var resultSet = connectionFactory.currentSession().execute(boundStatement);
var result = new ArrayList<Entity>();
for (var row : resultSet) {
result.add(new Entity(
row.getString("id"),
row.getInt("value1"),
row.getString("value2"),
row.getString("value3")));
}
return result;
});
}
}
@Repository
interface EntityRepository : CassandraRepository {
fun findByFilter(value2: String?): List<Entity> {
val sql = StringBuilder("SELECT id, value1, value2, value3 FROM entities")
if (value2 != null) {
sql.append(" WHERE value2 = ? ALLOW FILTERING")
}
val connectionFactory = cassandraConnectionFactory
val queryContext = QueryContext("EntityRepository.findByFilter", sql.toString())
return connectionFactory.query(queryContext) { statement ->
val boundStatement = if (value2 != null) statement.bind(value2) else statement.bind()
val resultSet = connectionFactory.currentSession().execute(boundStatement)
resultSet.map { row ->
Entity(
row.getString("id"),
row.getInt("value1"),
row.getString("value2"),
row.getString("value3")
)
}
}
}
}
Поскольку в Cassandra нет транзакций, query просто выполняется на текущей сессии с телеметрией; здесь нет фиксации или отката, которыми нужно управлять.
UDT¶
Поддерживаются типы UDT через аннотацию @UDT.
UDT описывает пользовательский тип Cassandra и может использоваться как поле обычного отображения.
Тип @UDT отображается как любое другое отображение, поэтому охватывающее отображение помечается @EntityCassandra.
Для следующей схемы, где username — пользовательский тип, хранящийся в столбце FROZEN:
CREATE TYPE IF NOT EXISTS username(first text, last text);
CREATE TABLE IF NOT EXISTS entities_udt
(
id VARCHAR,
name FROZEN<username>,
PRIMARY KEY (id)
);
отображение и репозиторий выглядят так:
@Repository
public interface EntityRepository extends CassandraRepository {
@EntityCassandra
record Entity(String id, Name name) {
@UDT
record Name(String first, String last) {}
}
@Query("SELECT * FROM entities_udt WHERE id = :id")
@Nullable
Entity findById(String id);
@Query("""
INSERT INTO entities_udt(id, name)
VALUES (:entity.id, :entity.name)
""")
void insert(Entity entity);
}
@Repository
interface EntityRepository : CassandraRepository {
@EntityCassandra
data class Entity(val id: String, val name: Name) {
@UDT
data class Name(val first: String, val last: String)
}
@Query("SELECT * FROM entities_udt WHERE id = :id")
fun findById(id: String): Entity?
@Query("""
INSERT INTO entities_udt(id, name)
VALUES (:entity.id, :entity.name)
""")
fun insert(entity: Entity)
}
Если тип UDT используется не через охватывающее отображение, а как самостоятельный тип Cassandra, генерацию отображателя можно включить явно с помощью @EntityCassandra.
Это полезно, когда отображатель нужен как отдельный компонент графа.
Макросы¶
Для упрощения написания CQL-запросов используйте макросы — они разворачиваются в CQL-конструкции на этапе компиляции.
Примеры использования показаны выше в секции Использование (методы findById и insert).
Подробная документация: Общие правила работы с базами данных — Макросы
Сигнатуры¶
Доступные из коробки сигнатуры методов репозитория:
T означает тип возвращаемого значения, либо List<T>, либо Void.
T myMethod()@Nullable T myMethod()Optional<T> myMethod()CompletionStage<T> myMethod()CompletionStageCompletableFuture<T> myMethod()CompletableFutureMono<T> myMethod()Project Reactor (требует зависимость)Flux<T> myMethod()Project Reactor (требует зависимость)
Обертки CompletionStage<T>, CompletableFuture<T> и Mono<T> также могут оборачивать List<T>,
например CompletionStage<List<Entity>> или Mono<List<Entity>>.
Параметры метода могут включать обычные значения, DTO, @Batch List<T> для пакетного выполнения и CqlSession, когда методу нужен доступ к текущей сессии драйвера.
T означает тип возвращаемого значения, либо T?, либо List<T>, либо Unit.
myMethod(): Tsuspend myMethod(): TKotlin Coroutine (требует зависимость какimplementation)myMethod(): Flow<T>Kotlin Coroutine (требует зависимость какimplementation)
Параметры метода могут включать обычные значения, DTO, @Batch List<T> для пакетного выполнения и CqlSession, когда методу нужен доступ к текущей сессии драйвера.
Телеметрия¶
Логирование, метрики и трассировка настраиваются через блок telemetry в конфигурации и описаны в разделе Справочник метрик.
Чтобы переопределить телеметрию полностью, можно предоставить собственные SPI-фабрики, подробнее в Общей документации по Базам данных.