
Cu ceva timp în urmă, ne-am confruntat cu problema curățării tuplurilor în spații . Curățarea trebuia să fie inițiată nu atunci când tarantool rămânea fără memorie, ci din timp și cu o anumită periodicitate. Pentru această sarcină, tarantool dispune de un modul, scris în Lua, numit . După o utilizare scurtă a acestui modul, am realizat că nu ni se potrivește: la curățarea constantă a unor volume mari de date, Lua se bloca în GC. Prin urmare, ne-am gândit la dezvoltarea propriului nostru modul capped expirationd, sperând că un cod scris într-o limbaj de programare nativ va rezolva cel mai bine problemele noastre.
Un exemplu bun pentru noi a fost modulul tarantool numit . Abordarea utilizată în acesta se bazează pe ideea că în spațiu se înființează un câmp separat, în care se indică timpul de viață al tuplului, altfel spus, ttl. Modulul scanează în fundal spațiul, compară ttl cu timpul curent și ia decizia dacă să ștergă tuplul sau nu. Codul modulului memcached este simplu și elegant, dar prea generic. În primul rând, nu ia în considerare tipul indexului conform căruia se face parcurgerea și ștergerea. În al doilea rând, la fiecare trecere sunt scanate toate tuplurile, ale căror număr poate fi destul de mare. Și dacă în modulul expirationd prima problemă a fost rezolvată (indexul tree a fost alocat într-o clasă separată), a doua nu a fost acordată deloc atenție. Aceste trei puncte au determinat alegerea de a scrie propriul cod.
Descriere
În documentația pentru tarantool există un foarte bun despre cum să scrii propriile proceduri stocate în limbajul C. În primul rând, sugerez să te familiarizezi cu acesta pentru a înțelege inserțiile cu comenzi și cod care vor apărea mai jos. De asemenea, merită să acorzi atenție la la obiectele disponibile în scrierea propriului modul capped, și anume la , , și .
Să începem de la distanță și să vedem cum arată modulul capped expirationd din exterior:
fiber = require('fiber')
net_box = require('net.box')
box.cfg{listen = 3300}
box.schema.func.create('libcapped-expirationd.start', {language = 'C'})
box.schema.user.grant('guest', 'execute', 'function', 'libcapped-expirationd.start')
box.schema.func.create('libcapped-expirationd.kill', {language = 'C'})
box.schema.user.grant('guest', 'execute', 'function', 'libcapped-expirationd.kill')
box.schema.space.create('tester')
box.space.tester:create_index('primary', {unique = true, parts = {1, 'unsigned'}})
capped_connection = net_box:new(3300)Pentru ușurința lansării, vom rula tarantool în directorul în care se află biblioteca noastră libcapped-expirationd.so. Din bibliotecă se exportă două funcții: start și kill. La început, trebuie să facem aceste funcții accesibile din Lua folosind box.schema.func.create și box.schema.user.grant. Apoi, vom crea un spațiu, ale cărui tupluri vor conține doar trei câmpuri: primul - un identificator unic, al doilea - un e-mail, iar al treilea - timpul de viață al tuplului. Deasupra primului câmp vom construi un index de tip tree și îl vom numi primary. Ulterior, obținem un obiect de conexiune la biblioteca noastră nativă.
După lucrările pregătitoare, lansăm funcția start:
capped_connection:call('libcapped-expirationd.start', {'non-indexed', box.space.tester.id, box.space.tester.index.primary, box.space.tester.index.primary, 3, 1024, 3600})Acest exemplu va funcționa la scanare exact la fel ca modulul expirationd, care este scris în Lua. Primul argument al funcției start primește un nume unic al sarcinii. Al doilea - identificatorul spațiului. Al treilea - indexul unic pe care se va realiza ștergerea tuplurilor. Al patrulea - indexul pe care se va face traversarea tuplurilor. Al cincilea - numărul câmpului tuplului cu timpul de viață (numerotarea începe de la 1, nu de la 0!). Al șaselea și al șaptelea - setările de scanare. 1024 - este numărul maxim de tupluri care sunt vizualizate în cadrul unei tranzacții. 3600 - timpul total de scanare în secunde.
Atenție, faptul că pentru traversare și ștergere în exemplu se folosește același index. Dacă acesta este un index tree, traversarea se realizează de la cheia mai mică la cheia mai mare. Dacă este un alt tip, de exemplu, un index hash, traversarea are loc, de regulă, într-o ordine aleatorie. La o scanare se vizualizează toate tuplurile din spațiu.
Să facem o inserare în spațiu a câtorva tupluri cu un timp de viață de 60 de secunde:
box.space.tester:insert{0, 'user0@tarantool.io', math.floor(fiber.time()) + 60}
box.space.tester:insert{1, 'user1@tarantool.io', math.floor(fiber.time()) + 60}
box.space.tester:insert{2, 'user2@tarantool.io', math.floor(fiber.time()) + 60}Să verificăm dacă inserarea a fost realizată cu succes:
tarantool> box.space.tester.index.primary:select()
---
- - [0, 'user0@tarantool.io', 1576418976]
- [1, 'user1@tarantool.io', 1576418976]
- [2, 'user2@tarantool.io', 1576418976]
...Repetăm selectul după 60+ de secunde (numărând de la începutul inserării primului tuplu) și vom observa că modulul capped expirationd a funcționat deja:
tarantool> box.space.tester.index.primary:select()
---
- []
...Oprim sarcina:
capped_connection:call('libcapped-expirationd.kill', {'non-indexed'})Să examinăm al doilea exemplu, când se folosește un index separat pentru parcurgere:
fiber = require('fiber')
net_box = require('net.box')
box.cfg{listen = 3300}
box.schema.func.create('libcapped-expirationd.start', {language = 'C'})
box.schema.user.grant('guest', 'execute', 'function', 'libcapped-expirationd.start')
box.schema.func.create('libcapped-expirationd.kill', {language = 'C'})
box.schema.user.grant('guest', 'execute', 'function', 'libcapped-expirationd.kill')
box.schema.space.create('tester')
box.space.tester:create_index('primary', {unique = true, parts = {1, 'unsigned'}})
box.space.tester:create_index('exp', {unique = false, parts = {3, 'unsigned'}})
capped_connection = net_box:new(3300)Aici totul este la fel ca în primul exemplu, cu o mică excepție. Peste al treilea câmp construim un index tree și îl numim exp. Acest index nu trebuie să fie unic, spre deosebire de indexul numit primary. Parcurgerea se va face prin indexul exp, iar ștergerea prin primary. Ne amintim că anterior atât parcurgerea, cât și ștergerea se făceau doar utilizând indexul primary.
După munca pregătitoare, apelăm funcția start cu argumente noi:
capped_connection:call('libcapped-expirationd.start', {'indexed', box.space.tester.id, box.space.tester.index.primary, box.space.tester.index.exp, 3, 1024, 3600})Din nou, vom insera în spațiu câteva tuple cu o durată de viață de 60 de secunde:
box.space.tester:insert{0, 'user0@tarantool.io', math.floor(fiber.time()) + 60}
box.space.tester:insert{1, 'user1@tarantool.io', math.floor(fiber.time()) + 60}
box.space.tester:insert{2, 'user2@tarantool.io', math.floor(fiber.time()) + 60}După 30 de secunde, prin analogie, vom adăuga câteva tuple suplimentare:
box.space.tester:insert{3, 'user3@tarantool.io', math.floor(fiber.time()) + 60}
box.space.tester:insert{4, 'user4@tarantool.io', math.floor(fiber.time()) + 60}
box.space.tester:insert{5, 'user5@tarantool.io', math.floor(fiber.time()) + 60}Să verificăm dacă inserarea a fost realizată cu succes:
tarantool> box.space.tester.index.primary:select()
---
- - [0, 'user0@tarantool.io', 1576421257]
- [1, 'user1@tarantool.io', 1576421257]
- [2, 'user2@tarantool.io', 1576421257]
- [3, 'user3@tarantool.io', 1576421287]
- [4, 'user4@tarantool.io', 1576421287]
- [5, 'user5@tarantool.io', 1576421287]
...Repetăm selectul după 60+ de secunde (numărând de la începutul inserării primului tuplu) și vom observa că modulul capped expirationd a funcționat deja:
tarantool> box.space.tester.index.primary:select()
---
- - [3, 'user3@tarantool.io', 1576421287]
- [4, 'user4@tarantool.io', 1576421287]
- [5, 'user5@tarantool.io', 1576421287]
...În spațiu au rămas tuple care mai au de trăit aproximativ 30 de secunde. Mai mult decât atât, scanarea s-a oprit la trecerea de la tuple cu identificatorul 2 și timpul de viață 1576421257 la tuple cu identificatorul 3 și timpul de viață 1576421287. Tupele cu timpul de viață 1576421287 și mai mare nu au fost vizualizate datorită ordonării cheilor indexului exp. Aceasta este economisirea pe care am dorit să o obținem de la început.
Oprim sarcina:
capped_connection:call('libcapped-expirationd.kill', {'indexed'})Implementarea
Cel mai bine despre toate particularitățile proiectului vorbesc întotdeauna sursele sale ! În cadrul publicației ne vom opri doar asupra celor mai importante aspecte, mai precis, asupra algoritmilor de parcurgere a spațiului.
Argumentele pe care le transmitem metodei start sunt păstrate într-o structură numită expirationd_task:
struct expirationd_task
{
char name[256];
uint32_t space_id;
uint32_t rm_index_id;
uint32_t it_index_id;
uint32_t it_index_type;
uint32_t field_no;
uint32_t scan_size;
uint32_t scan_time;
};Atributul name reprezintă numele sarcinii. Atributul space_id este identificatorul spațiului. Atributul rm_index_id este identificatorul indexului unic, pe care se va realiza ștergerea tuplurilor. Atributul it_index_id este identificatorul indexului, pe care se va realiza parcurgerea tuplurilor. Atributul it_index_type este tipul indexului, pe care se va realiza parcurgerea tuplurilor. Atributul field_no reprezintă numărul câmpului tuplului cu timpul de viață. Atributul scan_size este numărul maxim de tupluri, care este vizualizat în cadrul unei tranzacții. Atributul scan_time reprezintă timpul total de scanare în secunde.
Nu vom lua în considerare analiza argumentelor. Este o muncă minuțioasă, dar simplă, unde te va ajuta biblioteca . Dificultăți pot apărea doar cu indici care sunt transmiși din Lua sub forma unei structuri de date complexe cu tipul mp_map, și nu cu ajutorul tipurilor simple mp_bool, mp_double, mp_int, mp_uint și mp_array. Dar nu este necesar să analizezi întregul index. Este suficient doar să verifici unicitatea acestuia, să calculezi tipul și să extragi identificatorul.
Să enumerăm prototipurile tuturor funcțiilor utilizate pentru analiză:
bool expirationd_parse_name(struct expirationd_task *task, const char **pos);
bool expirationd_parse_space_id(struct expirationd_task *task, const char **pos);
bool expirationd_parse_rm_index_id(struct expirationd_task *task, const char **pos);
bool expirationd_parse_rm_index_unique(struct expirationd_task *task, const char **pos);
bool expirationd_parse_rm_index(struct expirationd_task *task, const char **pos);
bool expirationd_parse_it_index_id(struct expirationd_task *task, const char **pos);
bool expirationd_parse_it_index_type(struct expirationd_task *task, const char **pos);
bool expirationd_parse_it_index(struct expirationd_task *task, const char **pos);
bool expirationd_parse_field_no(struct expirationd_task *task, const char **pos);
bool expirationd_parse_scan_size(struct expirationd_task *task, const char **pos);
bool expirationd_parse_scan_time(struct expirationd_task *task, const char **pos);Și acum să trecem la cel mai important lucru — logica parcurgerii spațiului și ștergerii tuplurilor. Fiecare bloc de tupluri, de dimensiune de maximum scan_size, este parcurs și modificat într-o singură tranzacție. În cazul unui succes, această tranzacție este confirmată, iar în caz de eroare — anulată. Ultimul argument al funcției expirationd_iterate este un pointer la iteratorul de la care începe sau continuă scanarea. Acest iterator este incrementat intern până când apare o eroare, spațiul se termină sau există ocazia de a opri procesul anticipat. Funcția expirationd_expired verifică timpul de viață al tuplului, expirationd_delete — șterge tuplul, expirationd_breakable — verifică dacă trebuie să continuăm.
Codul funcției expirationd_iterate:
static bool
expirationd_iterate(struct expirationd_task *task, box_iterator_t **iterp)
{
box_iterator_t *iter = *iterp;
box_txn_begin();
for (uint32_t i = 0; i scan_size; ++i) {
box_tuple_t *tuple = NULL;
if (box_iterator_next(iter, &tuple) < 0) {
box_iterator_free(iter);
*iterp = NULL;
box_txn_rollback();
return false;
}
if (!tuple) {
box_iterator_free(iter);
*iterp = NULL;
box_txn_commit();
return true;
}
if (expirationd_expired(task, tuple))
expirationd_delete(task, tuple);
else if (expirationd_breakable(task))
break;
}
box_txn_commit();
return true;
}Codul funcției expirationd_expired:
static bool
expirationd_expired(struct expirationd_task *task, box_tuple_t *tuple)
{
const char *buf = box_tuple_field(tuple, task->field_no - 1);
if (!buf || mp_typeof(*buf) != MP_UINT)
return false;
uint64_t val = mp_decode_uint(&buf);
if (val > fiber_time64() / 1000000)
return false;
return true;
}Codul funcției expirationd_delete:
static void
expirationd_delete(struct expirationd_task *task, box_tuple_t *tuple)
{
uint32_t len;
const char *str = box_tuple_extract_key(tuple, task->space_id, task->rm_index_id, &len);
box_delete(task->space_id, task->rm_index_id, str, str + len, NULL);
}Codul funcției expirationd_breakable:
static bool
expirationd_breakable(struct expirationd_task *task)
{
return task->it_index_id != task->rm_index_id && task->it_index_type == ITER_GT;
}Aplicație
Puteți consulta codul sursă la !
Sursa: habr.com
