Cassandra
Модуль предоставляет реализацию репозитория для базы данных Cassandra
поверх Java-драйвера Cassandra (org.apache.cassandra:java-driver-core версии 4.19.3, подключается модулем транзитивно).
Cassandra — это распределенная колоночная база данных, где запросы пишутся на CQL, а модель данных обычно проектируется под конкретные сценарии чтения.
В Kora модуль Cassandra предоставляет декларативные репозитории поверх CqlSession: приложение пишет CQL-запросы в @Query, а Kora на этапе компиляции генерирует код подготовки запроса, связывания параметров и отображения результата.
Общие правила для отображений, @Repository, @Query, макросов, пакетных запросов и аннотаций @Table, @Column, @Id, @Embedded описаны в разделе общих правил работы с базами данных.
Этот документ охватывает специфичные для Cassandra части: подключение драйвера, конфигурацию CqlSession, профили выполнения, UDT, отображатели и поддерживаемые сигнатуры методов.
Пошаговый разбор перед справочным описанием смотрите в разделе База данных Cassandra.
Подключение¶
Зависимость build.gradle:
Модуль:
Зависимость build.gradle.kts:
Модуль:
CassandraDatabaseModule наследует CassandraMapperModule, поэтому стандартные отображатели строк, столбцов и параметров подключаются вместе с ним.
Модуль регистрирует компонент CassandraSession, который владеет драйверной CqlSession: она открывается при старте графа приложения и закрывается при остановке.
CassandraSession является Wrapped<CqlSession>, поэтому саму CqlSession также можно внедрить в любой компонент напрямую.
Конфигурация¶
Конфигурация читается из секции 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)!
driverMetrics = true //(117)!
slo = [ 1, 10, 50, 100, 200, 500, 1000 ] //(118)!
tags = { "key1" = "value1", "key2" = "value2" } //(119)!
}
tracing {
enabled = true //(120)!
attributes = { "key1" = "value1", "key2" = "value2" } //(121)!
}
}
}
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)!
driverMetrics: true #(117)!
slo: [ 1, 10, 50, 100, 200, 500, 1000 ] #(118)!
tags: { key1: "value1", key2: "value2" } #(119)!
tracing:
enabled: true #(120)!
attributes: { key1: "value1", key2: "value2" } #(121)!
- Имя пользователя для аутентификации в
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(по умолчанию:io.koraframework.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(по умолчанию:io.koraframework.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Наименьшая ожидаемая задержка для гистограммы метрики
session.throttlingDelay(по умолчанию:1ms). - Наибольшая ожидаемая задержка для гистограммы метрики
session.throttlingDelay(по умолчанию:90s). - Количество значащих цифр для гистограммы метрики
session.throttlingDelay(по умолчанию не указано, необязательно). - Интервал обновления снимка гистограммы метрики
session.throttlingDelay(по умолчанию не указано, необязательно). - Границы
SLOдля метрикиsession.throttlingDelay(по умолчанию:io.koraframework.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 (по умолчанию:
false). - Регистрирует собственные метрики драйвера в
MeterRegistryприложения; при выключении весь блокadvanced.metricsне имеет эффекта (по умолчанию:true). - Границы
SLOметрик Kora (по умолчанию:io.koraframework.telemetry.common.TelemetryConfig.MetricsConfig#DEFAULT_SLO). - Дополнительные теги метрик Kora (по умолчанию:
{}). - Включает трассировку запросов Kora (по умолчанию:
true). - Дополнительные атрибуты трассировки Kora (по умолчанию:
{}).
Конфигурация в коде¶
Не каждая опция драйвера доступна через секцию cassandra. Недостающее можно задать в коде компонентами Configurer:
Configurer<CqlSessionBuilder>— меняет сам построитель сессии, например идентификатор клиента, собственный реестр кодеков или слушатель состояния узлов;Configurer<ProgrammaticDriverConfigLoaderBuilder>— записывает низкоуровневые опции драйвера в загрузчик конфигурации, который Kora собирает из секцииcassandra.
Оба компонента необязательны и оба применяются последними, после всего прочитанного из конфигурации, поэтому заданное здесь значение побеждает:
@Component
public final class MyCqlSessionConfigurer implements Configurer<CqlSessionBuilder> {
@Override
public CqlSessionBuilder configure(CqlSessionBuilder builder) {
return builder.withClientId(UUID.randomUUID());
}
}
@Component
public final class MyDriverOptionsConfigurer implements Configurer<ProgrammaticDriverConfigLoaderBuilder> {
@Override
public ProgrammaticDriverConfigLoaderBuilder configure(ProgrammaticDriverConfigLoaderBuilder builder) {
return builder.withString(DefaultDriverOption.RETRY_POLICY_CLASS, "DefaultRetryPolicy");
}
}
@Component
class MyCqlSessionConfigurer : Configurer<CqlSessionBuilder> {
override fun configure(builder: CqlSessionBuilder): CqlSessionBuilder {
return builder.withClientId(UUID.randomUUID())
}
}
@Component
class MyDriverOptionsConfigurer : Configurer<ProgrammaticDriverConfigLoaderBuilder> {
override fun configure(builder: ProgrammaticDriverConfigLoaderBuilder): ProgrammaticDriverConfigLoaderBuilder {
return builder.withString(DefaultDriverOption.RETRY_POLICY_CLASS, "DefaultRetryPolicy")
}
}
Configurer — это io.koraframework.common.Configurer, контракт с единственным методом T configure(T t).
Несколько кластеров¶
CassandraDatabaseModule собирает свою сессию из секции cassandra, объявляя CassandraDatabaseFactoryModule("cassandra").
Второй кластер добавляется объявлением еще одного фабричного модуля со своим путем конфигурации и своим тегом,
а репозитории выбирают его через @Repository(executorTag = …):
@KoraApp
public interface Application extends CassandraDatabaseModule {
final class Analytics { }
@Tag(Analytics.class)
@FactoryModule
default CassandraDatabaseFactoryModule analyticsCassandra() {
return new CassandraDatabaseFactoryModule("cassandraAnalytics"); //(1)!
}
}
@Repository(executorTag = Application.Analytics.class)
public interface AnalyticsRepository extends CassandraRepository { }
- Секция файла конфигурации, описывающая этот кластер; она имеет ту же структуру, что и секция
cassandra.
@KoraApp
interface Application : CassandraDatabaseModule {
class Analytics
@Tag(Analytics::class)
@FactoryModule
fun analyticsCassandra(): CassandraDatabaseFactoryModule {
return CassandraDatabaseFactoryModule("cassandraAnalytics") //(1)!
}
}
@Repository(executorTag = Application.Analytics::class)
interface AnalyticsRepository : CassandraRepository
- Секция файла конфигурации, описывающая этот кластер; она имеет ту же структуру, что и секция
cassandra.
Тег фабричного модуля распространяется на все, что он создает, поэтому компоненты Configurer для такого кластера должны иметь тот же тег.
Репозиториям, работающим с основным подключением, тег не нужен.
Общие правила описаны в разделе общих правил работы с базами данных.
Использование¶
Чтобы создать репозиторий, объявите интерфейс с @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("SELECT id, value1, value2, value3 FROM entities") //(2)!
fun findAll(): List<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) переписываются в позиционные ?, а сгенерированный код связывает каждое значение
через драйвер, например _stmt.setString(0, id), где индекс соответствует порядку параметра в запросе.
Это обеспечивает безопасность (защита от CQL-инъекций) и производительность (использование подготовленных запросов драйвером).
В отличие от реляционных баз данных, в Cassandra нет транзакций.
Когда нужно, чтобы несколько операторов применились атомарно, используйте @Batch-метод: Kora строит UNLOGGED CQL BATCH
и связывает по одному оператору на каждый элемент аргумента @Batch List<T>.
Его семантика и макросы описаны в разделе общих правил работы с базами данных.
Профиль¶
Можно переопределить общие настройки частными настройками из профиля. Предположим, есть такая конфигурация профиля someProfile:
Чтобы применить настройки из профиля someProfile, достаточно сделать следующее:
Настройки, указанные в профиле, будут применяться к каждому запросу, в частности в данном случае будет установлен таймаут 10s.
Профиль применяется только к методу, помеченному @CassandraProfile; остальные методы репозитория продолжают использовать базовую конфигурацию.
Под капотом сгенерированный код вызывает setExecutionProfileName("someProfile") у оператора, поэтому имя профиля должно существовать в секции cassandra.profiles.
Отображение¶
Можно переопределить отображение различных частей отображения и параметров запроса — Kora предоставляет для этого специальные интерфейсы.
Из коробки CassandraMapperModule предоставляет отображатели для String, Byte, Short, Integer, Long, Float, Double, Boolean,
BigDecimal, BigInteger, UUID, ByteBuffer, byte[], LocalTime, LocalDate, LocalDateTime, ZonedDateTime, Instant и CqlDuration.
Если тип не входит в этот набор или ему нужно особое представление в CQL, добавьте собственный отображатель через @Mapping.
Отображатель с зависимостями в конструкторе должен быть объявлен как @Component, чтобы контейнер мог его собрать.
Отображатель без зависимостей аннотировать @Component нельзя — Kora создает его сама, а лишнее объявление компонента завершает сборку ошибкой Multiple components match.
Отображение¶
Используйте аннотацию @EntityCassandra для оптимального отображения.
Аннотация позволяет обработчику аннотаций сгенерировать все необходимые отображатели за один раунд аннотационной обработки:
CassandraRowMapper<T>, CassandraResultSetMapper<T> для одиночного результата и CassandraResultSetMapper<List<T>> для списка.
Без этой аннотации отображатели генерируются по требованию, что может потребовать множества раундов обработки и значительно увеличить время компиляции.
Это рекомендуемый способ отображения каждого типа, возвращаемого из репозитория или связываемого в нем.
Ожидается, что все вложенные отображения и типы UDT также используют эту аннотацию.
@EntityCassandra применима только к record-ам и классам в стиле Java bean (data class в Kotlin).
Результат¶
Если нужно вручную преобразовать весь результат синхронного запроса, используйте 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();
}
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(автоматически проходит по страницам результата).
CassandraMapperModule дополнительно выводит CassandraResultSetMapper<Optional<T>> из любого CassandraRowMapper<T> в графе,
поэтому сигнатурам с Optional<T> отдельный отображатель не нужен.
Строка¶
Если нужно вручную преобразовать одну строку результата, используйте CassandraRowMapper<T>.
Этот отображатель применяется к каждой строке и подходит для возвращаемых значений вида T, Optional<T> и List<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();
}
Столбец¶
Если нужно вручную преобразовать значение столбца, используйте 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();
}
class ColumnMapper : CassandraRowColumnMapper<UUID> {
override fun apply(row: GettableByName, index: Int): UUID {
return UUID.fromString(row.getString(index)!!)
}
}
@Table("entities")
data class Entity(
@field:Id @Mapping(ColumnMapper::class) val id: UUID,
val name: String
)
@Repository
interface EntityRepository : CassandraRepository {
@Query("SELECT id, name FROM entities")
fun findAll(): List<Entity>
}
Через CassandraRowColumnMapper читается и тип @UDT, поэтому этим же интерфейсом можно взять отображение пользовательского типа под собственный контроль.
Параметр¶
Если нужно вручную преобразовать значение параметра запроса, используйте CassandraParameterColumnMapper<T>.
Он получает SettableByName<?>, индекс параметра и значение из метода репозитория.
Значение может быть null, поэтому реализация сама отвечает за то, чтобы не записывать ничего в подстановку, когда записывать нечего.
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);
}
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>
}
Асинхронный¶
Методы репозитория на Java могут возвращать CompletionStage<T> или CompletableFuture<T>.
Для таких методов Kora выполняет оператор через CqlSession#executeAsync и отображает результат с помощью CassandraAsyncResultSetMapper<T>,
который получает AsyncResultSet и возвращает CompletionStage<T>.
Для отображения с @EntityCassandra дополнительный отображатель не нужен: Kora выводит асинхронный отображатель из сгенерированного CassandraRowMapper<T>,
используя CassandraAsyncResultSetMapper.list(...) для результата List<T> и CassandraAsyncResultSetMapper.one(...) в остальных случаях.
Метод list проходит по страницам результата, запрашивая следующую, пока набор не закончится, поэтому List<T> соберет все страницы до завершения future;
one отображает только первую строку.
@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")
CompletableFuture<Entity> findById(String id);
@Query("SELECT %{return#selects} FROM %{return#table}")
CompletionStage<List<Entity>> findAll();
@Query("INSERT INTO %{entity#inserts}")
CompletionStage<Void> insert(Entity entity);
}
Собственный асинхронный отображатель подключается так же, как синхронный, через @Mapping:
final class AsyncResultMapper implements CassandraAsyncResultSetMapper<List<UUID>> {
private static final CassandraAsyncResultSetMapper<List<UUID>> DELEGATE =
CassandraAsyncResultSetMapper.list(row -> row.getUuid("id"));
@Override
public CompletionStage<List<UUID>> apply(AsyncResultSet rows) {
return DELEGATE.apply(rows);
}
}
@Repository
public interface EntityRepository extends CassandraRepository {
@Mapping(AsyncResultMapper.class)
@Query("SELECT id FROM entities")
CompletionStage<List<UUID>> getIds();
}
У репозиториев на Kotlin асинхронных сигнатур нет: методы синхронные и используют CassandraResultSetMapper / CassandraRowMapper,
см. Результат и Сигнатуры.
Ручной запрос¶
Если запрос сложно выразить одной статической @Query, вы можете объявить обычный метод с реализацией и построить CQL вручную.
Репозиторий предоставляет executor(), который возвращает CassandraExecutor:
currentSession()— активная драйвернаяCqlSession;telemetry()—DatabaseTelemetry, используемая для отчетности;query(CassandraQuery, CassandraResultSetMapper<T>),queryOne(...),queryOptional(...),queryList(...)— выполняют построенный запрос и отображают результат;query(CassandraQuery, Function<BoundStatement, T>)— выполняет построенный запрос и отдает связанный оператор вам;query(QueryContext, Function<PreparedStatement, T>)— самый низкий уровень: подготавливает сырую строкуCQL, все остальное делаете вы.
Каждый из них оборачивает вызов в телеметрию Kora, поэтому ручной запрос логируется, измеряется и трассируется точно так же, как сгенерированный.
CQL собирается построителем CassandraQuery, который держит текст запроса и значения параметров раздельно:
CassandraQuery.named()—CQLс подстановками:name, значения задаются черезbind(name, value),bindAll(map),bindIn(name, values)для конструкцийIN (:name)и условные формыcqlIf,bindIf,bindInIf;CassandraQuery.template()иCassandraQuery.template(cql, args...)—CQL, в котором уже используются позиционные подстановки?;opts(...)— параметры конкретного оператора:consistencyLevel,serialConsistencyLevel,pageSize,timeout,idempotent,tracing.
build() проверяет запрос и завершается IllegalArgumentException, если подстановка не связана, если связанный параметр нигде не используется в CQL
или если коллекция bindIn пуста.
Именованные параметры преобразуются в позиционные ?, поэтому драйвер по-прежнему получает подготовленный запрос, а значения никогда не подставляются в текст запроса.
@Repository
public interface EntityRepository extends CassandraRepository {
default List<Entity> findByFilter(@Nullable String value2) {
var query = CassandraQuery.named()
.cql("SELECT id, value1, value2, value3 FROM entities")
.cqlIf(" WHERE value2 = :value2 ALLOW FILTERING", value2 != null)
.bindIf("value2", value2, value2 != null)
.build();
return executor().queryList(query, row -> new Entity(
row.getString("id"),
row.getInt("value1"),
row.getString("value2"),
row.getString("value3")));
}
}
@Repository
interface EntityRepository : CassandraRepository {
fun findByFilter(value2: String?): List<Entity> {
val query = CassandraQuery.named()
.cql("SELECT id, value1, value2, value3 FROM entities")
.cqlIf(" WHERE value2 = :value2 ALLOW FILTERING", value2 != null)
.bindIf("value2", value2, value2 != null)
.build()
return executor().queryList(query) { row ->
Entity(
row.getString("id")!!,
row.getInt("value1"),
row.getString("value2")!!,
row.getString("value3")
)
}
}
}
Когда нужен полный контроль над подготовкой оператора, используйте перегрузку с QueryContext.
QueryContext содержит идентификатор запроса, исполняемый CQL и имя операции:
идентификатор попадает в телеметрию как атрибут db.query.text, а имя операции становится именем спана,
поэтому используйте стабильное значение вида Repository.method — именно так делают сгенерированные репозитории.
Конструктор с двумя аргументами подставляет операцию db_query.
@Repository
public interface EntityRepository extends CassandraRepository {
default List<Entity> findAllValues() {
var executor = executor();
var queryContext = new QueryContext(
"SELECT id, value1, value2, value3 FROM entities",
"SELECT id, value1, value2, value3 FROM entities",
"EntityRepository.findAllValues");
return executor.query(queryContext, statement -> executor.currentSession()
.execute(statement.bind())
.map(row -> new Entity(
row.getString("id"),
row.getInt("value1"),
row.getString("value2"),
row.getString("value3")))
.all());
}
}
@Repository
interface EntityRepository : CassandraRepository {
fun findAllValues(): List<Entity> {
val executor = executor()
val queryContext = QueryContext(
"SELECT id, value1, value2, value3 FROM entities",
"SELECT id, value1, value2, value3 FROM entities",
"EntityRepository.findAllValues"
)
return executor.query(queryContext) { statement ->
executor.currentSession()
.execute(statement.bind())
.map { row ->
Entity(
row.getString("id")!!,
row.getInt("value1"),
row.getString("value2")!!,
row.getString("value3")
)
}
.all()
}
}
}
Поскольку в Cassandra нет транзакций, query просто выполняется на текущей сессии с телеметрией; здесь нет фиксации или отката, которыми нужно управлять.
UDT¶
Поддерживаются типы UDT через аннотацию @UDT.
UDT описывает пользовательский тип Cassandra и может использоваться как поле обычного отображения.
Для каждого типа @UDT Kora генерирует чтение и запись, поэтому тип работает и в результатах, и как параметр запроса,
включая вложенные типы @UDT и List<T> из типа @UDT.
Тип @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 используется не через охватывающее отображение, а сам является результатом запроса, добавьте @EntityCassandra рядом с @UDT.
Одна @UDT генерирует чтение и запись столбца; @EntityCassandra дополнительно генерирует отображатели строки и результата,
а именно они нужны методу репозитория, возвращающему такой тип.
Макросы¶
Для упрощения написания CQL-запросов используйте макросы — они разворачиваются в CQL-конструкции на этапе компиляции.
Примеры использования показаны выше в секции Использование (методы findById и insert).
Подробная документация: Общие правила работы с базами данных — Макросы
Сигнатуры¶
Доступные из коробки сигнатуры методов репозитория:
T означает тип возвращаемого значения, либо List<T>, либо Void.
T myMethod()@Nullable T myMethod()Optional<T> myMethod()CompletionStage<T> myMethod()CompletionStageCompletableFuture<T> myMethod()CompletableFuture
Обертки CompletionStage<T> и CompletableFuture<T> также могут оборачивать List<T> или Void,
например CompletionStage<List<Entity>> или CompletableFuture<Void>.
Параметры метода могут включать обычные значения, DTO, @Batch List<T> для пакетного выполнения и CqlSession, когда методу нужен доступ к текущей сессии драйвера.
T означает тип возвращаемого значения, либо T?, либо List<T>, либо Unit.
myMethod(): TmyMethod(): T?myMethod(): List<T>myMethod()для запроса без результата
Методы репозитория на Kotlin синхронные. Метод suspend, помеченный @Query, отклоняется процессором символов
с сообщением Suspend methods are not supported by the repository generator — выполняйте блокирующий вызов, а конкурентность стройте над репозиторием.
Параметры метода могут включать обычные значения, DTO, @Batch List<T> для пакетного выполнения и CqlSession, когда методу нужен доступ к текущей сессии драйвера.
Реактивные возвращаемые типы не поддерживаются: сигнатур Mono / Flux нет, реактивного отображателя результата в модуле тоже нет.
Телеметрия¶
Логирование, метрики и трассировка настраиваются через блок telemetry в конфигурации и описаны в разделе Справочник метрик.
По умолчанию логирование и метрики запросов выключены (telemetry.logging.enabled = false, telemetry.metrics.enabled = false), а трассировка включена (telemetry.tracing.enabled = true).
Собственные метрики драйвера управляются отдельно ключом telemetry.metrics.driverMetrics и настраиваются блоком advanced.metrics.
Чтобы переопределить телеметрию полностью, можно предоставить собственные SPI-фабрики, подробнее в Общей документации по Базам данных.