Aftësitë e gjuhës Q dhe KDB+ në shembullin e shërbimit në kohë reale

Për atë që është baza KDB+, gjuha e programimit Q, çfarë dobish dhe dobësish kanë ato, mund të lexoni në artikullin tim të mëparshëm artikullin dhe shkurtimisht në hyrje. Në këtë artikull do të implementojmë një shërbim në Q, i cili do të përpunojë fluksin e të dhënave në hyrje dhe do të llogarisë funksione të ndryshme agreguese në minutë në modalitetin "të dhënave të drejtpërdrejta" (dmth. do të përfundojë të gjitha llogaritë para ardhjes së grupit të ardhshëm të të dhënave). Karakteristika kryesore e Q është se ajo është një gjuhë vektoriale, që lejon të punoheni jo vetëm me objekte të veçanta, por me grupe, grupe të grupeve dhe objekte të tjera komplekse. Gjuhë si Q dhe të afërmit e saj K, J, APL janë të njohura për shkurtësinë e tyre. Shpesh një program që zë disa ekrane kodigo në një gjuhë të njohur si Java, mund të shkruhet me to në disa rreshta. Këtë kam për qëllim të demonstroj në këtë artikull.

Aftësitë e gjuhës Q dhe KDB+ në shembullin e shërbimit në kohë reale

Hyrje

KDB+ është një bazë të dhënash kolonash e optimizuar për volumet shumë të mëdha të të dhënave, të organizuara në një mënyrë të caktuar (në radhë të parë sipas kohës). Përdoret kryesisht në organizatat financiare – banka, fonde investimi, kompani sigurimesh. Gjuha Q është gjuha e brendshme e KDB+, e cila mundëson punën efikase me këto të dhëna. Ideologjia e Q është shkurtësia dhe efikasiteti, ndërsa kuptueshmëria sakrifikohet. Kjo arsyetohet me faktin se gjuha vektoriale do të jetë e vështirë për t'u kuptuar në çdo rast, dhe shkurtësia e shënimeve lejon që të shihni një pjesë shumë më të madhe të programit në një ekran, që përfundimisht lehtëson kuptimin e tij.

Në këtë artikull do të realizojmë një program të plotë në Q dhe ndoshta do të dëshironit ta provoni atë. Për këtë do t'ju nevojitet krijimi i Q. Mund ta shkarkoni versionin 32-bit falas në faqen e kompanisë kx – www.kx.com. Aty gjithashtu, nëse jeni të interesuar, do të gjeni informacion ndihmës për Q, librin Q For Mortals dhe shumë artikuj të ndryshëm mbi këtë temë.

Vendosja e detyrës

Ekziston një burim që dërgon një tabelë me të dhëna çdo 25 milisekonda. Duke pasur parasysh se KDB+ përdoret kryesisht në financa, le të supozojmë se kjo është një tabelë tregtish (trades), e cila ka kolonat e mëposhtëm: kohë (time, në milisekonda), sim (simboli i kompanisë në bursë – IBM, AAPL, …), çmimi (price, çmimi me të cilin janë blerë aksionet), madhësia (size, madhësia e tregtisë). Intervali 25 milisekonda është zgjedhur në mënyrë arbitrare, ai nuk është shumë i vogël dhe as shumë i madh. Prania e tij do të thotë se të dhënat arrijnë në shërbim të friksonuara. Mund të ishte e lehtë të realizohej friksonimi në anën e shërbimit, përfshirë friksonimin dinamik, që varet nga ngarkesa aktuale, por për thjeshtësi do të ndalojmë në një interval fikse.

Shërbimi duhet të llogarisë çdo minutë për çdo sim të ardhur nga kolona sim një set funksionesh aggreguese – max price, avg price, sum size dhe informacione të tjera të dobishme. Për thjeshtësi, le të supozojmë se të gjitha funksionet mund të llogariten inkrimantale, pra për të marrë një vlerë të re është e mjaftueshme të di dy numra – vlerën e vjetër dhe vlerën e ardhshme. Për shembull, funksionet max, average, sum kanë këtë pronë, ndërsa funksioni median nuk e ka.

Gjithashtu, ne do të supozojmë se fluksi i të dhënave në hyrje është i renditur sipas kohës. Kjo do të na japë mundësinë të punojmë vetëm me minutën e fundit. Në praktikë, mjafton të dimë të punojmë me minutën aktuale dhe me atë të mëparshme në rast se ndonjë përditësim vonohet. Për thjeshtësi, nuk do të shqyrtojmë këtë rast.

Funksionet agreguese

Më poshtë janë renditur funksionet e nevojshme agreguese. I kam marrë sa më shumë që të jetë e mundur për të rritur ngarkesën në shërbim:

  • high – max price – çmimi maksimal për minutë.
  • low – min price – çmimi minimal për minutë.
  • firstPrice – first price – çmimi i parë për minutë.
  • lastPrice – last price – çmimi i fundit për minutë.
  • firstSize – first size – madhësia e parë e kontratës për minutë.
  • lastSize – last size – madhësia e fundit e kontratës për minutë.
  • numTrades – count i – numri i kontratave për minutë.
  • volume – sum size – shuma e madhësive të kontratave për minutë.
  • pvolume – sum price – shuma e çmimeve për minutë, e nevojshme për avgPrice.
  • turnover – sum price*size – volumi total i kontratave për minutë.
  • avgPrice – pvolume%numTrades – çmimi mesatar për minutë.
  • avgSize – volume%numTrades – madhësia mesatare e kontratës për minutë.
  • vwap – turnover%volume – çmimi mesatar i peshuar sipas madhësisë së kontratës për minutë.
  • cumVolume – shuma e volumit – masa e akumuluar e transaksioneve gjatë gjithë kohës.

Le të diskutojmë një moment të qartë – si të inicializojmë këto kolona për herë të parë dhe për çdo minutë pasuese. Disa kolona si firstPrice duhet të inicializohen çdo herë me vlerën null, vlera e tyre nuk është e përcaktuar. Të tjerat si volume duhet të vendosen gjithmonë në 0. Ka gjithashtu kolona që kërkojnë një qasje të kombinuar – për shembull, cumVolume duhet të kopjohet nga minuta e kaluar, ndërsa për të parën të vendoset në 0. Do t’i caktojmë këto parametra duke përdorur llojin e të dhënave fjalor (ngjashëm si një regjistrim):

// list ! list – создать словарь, 0n – float null, 0N – long null, `sym – тип символ, `sym1`sym2 – список символов
initWith:`sym`time`high`low`firstPrice`lastPrice`firstSize`lastSize`numTrades`volume`pvolume`turnover`avgPrice`avgSize`vwap`cumVolume!(`;00:00;0n;0n;0n;0n;0N;0N;0;0;0.0;0.0;0n;0n;0n;0);
aggCols:reverse key[initWith] except `sym`time; // список всех вычисляемых колонок, reverse объяснен ниже

Kam shtuar sym dhe time në fjalor për lehtësi, tani initWith është një rresht i gatshëm nga tabela përfundimtare e agreguar, ku mbetet të caktohen simbolet e sakta dhe koha. Mund të përdoret për të shtuar rreshta të rinj në tabelë.

Këto aggCols do na nevojiten gjatë krijimit të funksionit agregues. Lista duhet të invertohet për shkak të veçorive të rendit të llogaritjeve të shprehjeve në Q (nga djathtas në majtas). Qëllimi është të sigurohet llogaritja në drejtimin nga high në cumVolume, pasi disa kolona varen nga ato të mëparshmet.

Kolonat që duhet të kopjohen në minutën e re nga e kaluara, kolona sym është shtuar për lehtësi:

rollColumns:`sym`cumVolume;

Tani do t'i ndajmë kolonat në grupe në përputhje me mënyrën se si duhet të përditësohen. Mund të përmendim tri tipe:

  1. Akumulatorët (volume, turnover, ..) – duhet ta shtojmë vlerën hyrëse në atë të mëparshmen.
  2. Me një pikë të veçantë (high, low, ..) – vlera e parë në minutë merret nga të dhënat hyrëse, ndërsa të tjerat llogariten me funksionin.
  3. Të tjerët. Gjithmonë llogariten me funksionin.

Do të definojmë variablat për këto klasa:

accumulatorCols:`numTrades`volume`pvolume`turnover;
specialCols:`high`low`firstPrice`firstSize;

Renditja e llogaritjeve

Ne do ta përditësojmë tabelën agreguese në dy hapa. Për efikasitet, së pari do ta ngushtojmë tabelën hyrëse, që të mbetet vetëm një rresht për secilin simbol dhe minutë. Ajo që të gjitha funksionet tona janë inkrementale dhe asociative na garanton se rezultati nga ky hap shtesë nuk do të ndyshojë. Tabelën mund ta ngushtojmë me anë të selektit:

select high:max price, low:min price … by sym,time.minute from table

Ky ky method ka ka minus – seti i kolonave të llogaritur është caktuar paraprakisht. Me fat, në Q selekti është realizuar gjithashtu si funksion, ku mund të vendosen argumente të krijuara dinamikisht:

?[table;whereClause;byClause;selectClause]

Nuk do të përshkruaj në detaje formatin e argumenteve; në rastin tonë, vetëm shprehjet by dhe select do të jenë jo triviale dhe ato duhet të jenë fjalorë të tipit columns!expressions. Kështu, funksioni kompresues mund të caktohet kështu:

selExpression:`high`low`firstPrice`lastPrice`firstSize`lastSize`numTrades`volume`pvolume`turnover!parse each ("max price";"min price";"first price";"last price";"first size";"last size";"count i";"sum size";"sum price";"sum price*size"); \/\/ each është funksioni map në Q për një listë
preprocess:?[;();`sym`time!`sym`time.minute;selExpression];

Për qartësi kam përdorur funksionin parse, i cili e kthen një varg me shprehjen Q në një vlerë që mund të kalojë në funksionin eval dhe që kërkohet në selektin funksional. Gjithashtu, vlen të theksohet se preprocess është caktuar si projekcion (dmth. funksion me argumente të përcaktuara pjesërisht) i funksionit selekt, një argument (tabela) mungon. Nëse e aplikojmë preprocess në tabelë, do të marrim një tabelë të kompresuar.

Faza e dytë – është përditësimi i tabelës së agreguar. Le të shkruajmë së pari algoritmin në pseudokod:

për secilin sim në inputTable
  idx: indeksi i rreshtit në agg table për sim+currentTime;
  aggTable[idx;`high]: aggTable[idx;`high] | inputTable[sym;`high];
  aggTable[idx;`volume]: aggTable[idx;`volume] + inputTable[sym;`volume];
  …

Në Q, në vend të cikleve është e zakonshme të përdoren funksionet map/reduce. Por duke qenë se Q është një gjuhë vektoriale dhe të gjitha operacionet mund t’i aplikojmë qetësisht mbi të gjitha simbolët njëherazi, në një qasje të parë, mund të kalojmë edhe pa cikël, duke kryer operacione me të gjithë simbolët njëherazi:

idx:calcIdx inputTable;
row:aggTable idx;
aggTable[idx;`high]: row[`high] | inputTable`high;
aggTable[idx;`volume]: row[`volume] + inputTable`volume;
…

Por ne mund të shkojmë edhe më tej, në Q ka një operator unik dhe jashtëzakonisht të fuqishëm – operatori i përgjithësimit të caktimit. Ai lejon ndryshimin e grupit të vlerave në një strukturë të ndërlikuar të të dhënave duke përdorur një listë indekset, funksionesh dhe argumentesh. Në rastin tonë, ai duket kështu:

idx:calcIdx inputTable;
rows:aggTable idx;
// .[target;(idx0;idx1;..);function;argument] ~ target[idx 0;idx 1;…]: function[target[idx 0;idx 1;…];argument], në rastin tonë funksioni – është caktimi
.[aggTable;(idx;aggCols);:;flip (row[`high] | inputTable`high;row[`volume] + inputTable`volume;…);]

Fatkeqësisht, për të ndarë në një tabelë ne kemi nevojë për një listë rreshtash, jo kolona, dhe duhet të transponojmë matricën (lista e kolonave në listën e rreshtave) duke përdorur funksionin flip. Për një tabelë të madhe, kjo është e kushtueshme, prandaj, në vend të kësaj, do të përdorim caktimin e përgjithësuar për çdo kolonë veçmas, duke përdorur funksionin map (i cili duket si një apostrof):

[aggTable;;:;]'[(idx;)each aggCols; (row[`high] | inputTable`high;row[`volume] + inputTable`volume;…)];

Ne përsëri po përdorim projekcionin e funksionit. Gjithashtu vini re se në Q krijimi i një liste është gjithashtu një funksion dhe ne mund ta thërrasim atë duke përdorur funksionin each(map), për të marrë një listë listash.

Për të mos pasur një grup kolonash të llogaritura të rregullta, do të krijojmë shprehjen lart dinamikisht. Së pari do të përcaktojmë funksionet për llogaritjen e çdo kolone, duke përdorur variablat row dhe inp për të iu referuar të dhënave të agreguara dhe hyrëse:

aggExpression:`high`low`firstPrice`lastPrice`firstSize`lastSize`avgPrice`avgSize`vwap`cumVolume!
 ("row[`high]|inp`high";"row[`low]&inp`low";"row`firstPrice";"inp`lastPrice";"row`firstSize";"inp`lastSize";"pvolume%numTrades";"volume%numTrades";"turnover%volume";"row[`cumVolume]+inp`volume");

Disa kolona janë të veçanta, vlera e parë e tyre nuk duhet të llogaritet nga funksioni. Ne mund ta përcaktojmë se ajo është e para sipas kolonës row[`numTrades] – nëse është 0, atëherë vlera është e para. Në Q ka një funksion seleksioni — ?[Boolean list;list1;list2] – i cili zgjedh vlerën nga lista 1 ose 2 në varësi të kushteve në argumentin e parë:

// high -> ?[isFirst;inp`high;row[`high]|inp`high]
// @ - тоже обобщенное присваивание для случая когда индекс неглубокий
@[`aggExpression;specialCols;{[x;y]"?[isFirst;inp`",y,";",x,"]"};string specialCols];

Këtu kam thirrur një caktim të përgjithshëm me funksionin tim (shprehja në kllapat e treguesve). Ajo merr vlerën aktuale (argumenti i parë) dhe një argument shtesë, që e transfert në parametrin e 4-të.

Veçmas, do të shtojmë kolonat akumuluese, pasi për to funksioni është i njëjtë:

// volume -> row[`volume]+inp`volume
aggExpression[accumulatorCols]:{"row[`",x,"]+inp`",x } each string accumulatorCols;

Kjo është një caktim i zakonshëm për standardet e Q, vetëm se unë caktoj menjëherë një listë vlerash. Në fund, le të krijojmë funksionin kryesor:

// ":",/:aggExprs ~ map[{":",x};aggExpr] => ":row[`high]|inp`high" присвоим вычисленное значение переменной, потому что некоторые колонки зависят от уже вычисленных значений
// string[cols],'exprs ~ map[,;string[cols];exprs] => "high:row[`high]|inp`high" завершим создание присваивания. ,’ расшифровывается как map[concat]
// ";" sv exprs – String from Vector (sv), соединяет список строк вставляя “;” посредине
updateAgg:value "{[aggTable;idx;inp] row:aggTable idx; isFirst_0=row`numTrades; .[aggTable;;:;]'[(idx;)each aggCols;(",(";"sv string[aggCols],'":",/:aggExpression aggCols),")]}";

Me këtë shprehje krijoj dinamikisht një funksion nga një varg, i cili përmban shprehjen që përmenda më parë. Rezultati do të duket kështu:

{[aggTable;idx;inp] rows:aggTable idx; isFirst_0=row`numTrades; .[aggTable;;:;]'[(idx;)each aggCols ;(cumVolume:row[`cumVolume]+inp`cumVolume;… ; high:?[isFirst;inp`high;row[`high]|inp`high])]}

Rendi i llogaritjes së kolonave është i invertuar, sepse në Q rendi i llogaritjes është nga e djathta në të majtë.

Tani, ne kemi dy funksione kryesore që na nevojiten për llogaritjet, mbetet të shtojmë pak infrastrukturë dhe shërbimi është gati.

Hapat përfundimtarë

Ne kemi funksione preprocess dhe updateAgg, të cilat bëjnë të gjithë punën. Por është e nevojshme gjithashtu të sigurojmë kalimin e duhur përmes minutash dhe të llogarisim indekset për agregimin. Së pari, do të përcaktojmë funksionin init:

init:{
  tradeAgg:: 0#enlist[initWith]; // krijojmë një tabelë të tipizuar bosh, enlist e kthen fjalorin në tabelë, dhe 0# do të thotë të marrim 0 elemente prej saj
  currTime::00:00; // do të fillojmë nga 0, :: do të thotë se po i japim vlerë një variabli global
  currSyms::`u#`symbol$(); // `u# - e kthen listën në pemë, për të përshpejtuar kërkimin e elementeve
  offset::0; // indeksi në tradeAgg, ku fillon minuta aktuale 
  rollCache:: `sym xkey update `u#sym from rollColumns#tradeAgg; // cache për vlerat e fundit të kolonave roll, tabela me çelësin sym
 }

Po ashtu do të përcaktojmë funksionin roll, i cili do të ndryshojë minutën aktuale:

roll:{[tm]
  if[currTime>tm; :init[]]; // nëse kalojmë mesnatën, thjesht do të thërrasim init
  rollCache,::offset _ rollColumns#tradeAgg; // përditësojmë cache – merrni kolonat roll nga aggTable, i prisni, i shtoni në rollCache
  offset::count tradeAgg;
  currSyms::`u#`$();
 }

Na nevojitet një funksion për të shtuar simbole të reja:

addSyms:{[syms]
  currSyms,::syms; // do shtojmë në listën e njohur
  // do shtojmë në tabelën sym, time dhe rollColumns duke përdorur caktimin e përgjithshëm.
  // Funksioni ^ zëvendëson vlerat e paracaktuara për kolonat roll, nëse simboli nuk është në cache. vlera e rrotës kthen listën e kolonave në tabelë.
  `tradeAgg upsert @[count[syms]#enlist initWith;`sym`time,cols rc;:;(syms;currTime), (initWith cols rc)^value flip rc:rollCache ([] sym: syms)];
 }

Dhe, përfundimisht, funksioni upd (emri tradicional i këtij funksioni për shërbimet Q), i cili thirret nga klienti për të shtuar të dhëna:

upd:{[tblName;data] // tblName nuk na nevojitet, por zakonisht shërbimi përpunon disa tabela 
  tm:exec distinct time from data:() xkey preprocess data; // preprocess & llogarit kohën
  updMinute[data] each tm; // do shtojmë të dhënat për çdo minutë
};
updMinute:{[data;tm]
  if[tm<>currTime; roll tm; currTime::tm]; // ndërroni minutën, nëse është e nevojshme
  data:select from data where time=tm; // filtrimi
  if[count msyms:syms where not (syms:data`sym)in currSyms; addSyms msyms]; // simbole të reja
  updateAgg[`tradeAgg;offset+currSyms?syms;data]; // do ta përditësojmë tabelën e agreguar. Funksioni ? kërkon indeksin e elementeve të listës së djathtë në listën e majtë.
 };

Ja, dhe kjo është gjithçka. Kjo është kodin e plotë të shërbimit tonë, siç u premtua, vetëm disa rreshta:

initWith:`sym`time`high`low`firstPrice`lastPrice`firstSize`lastSize`numTrades`volume`pvolume`turnover`avgPrice`avgSize`vwap`cumVolume!(`;00:00;0n;0n;0n;0n;0N;0N;0;0;0.0;0.0;0n;0n;0n;0);
aggCols:reverse key[initWith] except `sym`time;
rollColumns:`sym`cumVolume;

accumulatorCols:`numTrades`volume`pvolume`turnover;
specialCols:`high`low`firstPrice`firstSize;

selExpression:`high`low`firstPrice`lastPrice`firstSize`lastSize`numTrades`volume`pvolume`turnover!parse each ("max price";"min price";"first price";"last price";"first size";"last size";"count i";"sum size";"sum price";"sum price*size");
preprocess:?[;();`sym`time!`sym`time.minute;selExpression];

aggExpression:`high`low`firstPrice`lastPrice`firstSize`lastSize`avgPrice`avgSize`vwap`cumVolume!("row[`high]|inp`high";"row[`low]&inp`low";"row`firstPrice";"inp`lastPrice";"row`firstSize";"inp`lastSize";"pvolume%numTrades";"volume%numTrades";"turnover%volume";"row[`cumVolume]+inp`volume");
@[`aggExpression;specialCols;{"?[isFirst;inp`",y,";",x,"]"};string specialCols];
aggExpression[accumulatorCols]:{"row[`",x,"]+inp`",x } each string accumulatorCols;
updateAgg:value "{[aggTable;idx;inp] row:aggTable idx; isFirst_0=row`numTrades; .[aggTable;;:;]'[(idx;)each aggCols;(",(";"sv string[aggCols],'":",/:aggExpression aggCols),")]}"; / '

init:{
  tradeAgg::0#enlist[initWith];
  currTime::00:00;
  currSyms::`u#`symbol$();
  offset::0;
  rollCache:: `sym xkey update `u#sym from rollColumns#tradeAgg;
 };
roll:{[tm]
  if[currTime>tm; :init[]];
  rollCache,::offset _ rollColumns#tradeAgg;
  offset::count tradeAgg;
  currSyms::`u#`$();
 };
addSyms:{[syms]
  currSyms,::syms;
  `tradeAgg upsert @[count[syms]#enlist initWith;`sym`time,cols rc;:;(syms;currTime),(initWith cols rc)^value flip rc:rollCache ([] sym: syms)];
 };

upd:{[tblName;data] updMinute[data] each exec distinct time from data:() xkey preprocess data};
updMinute:{[data;tm]
  if[tm<>currTime; roll tm; currTime::tm];
  data:select from data where time=tm;
  if[count msyms:syms where not (syms:data`sym)in currSyms; addSyms msyms];
  updateAgg[`tradeAgg;offset+currSyms?syms;data];
 };

Testimi

Do t'i kontrollojmë performancën e shërbimit. Për këtë, do ta nisim atë në një proces të veçantë (vendosni kodin në skedarin service.q) dhe thërrisni funksionin init:

q service.q –p 5566

q)init[]

Në një konsol tjetër, nisni procesin e dytë Q dhe lidheni me të parin:

h:hopen `:host:5566
h:hopen 5566 \/\/ nëse të dy janë në të njëjtin host

Fillimisht do të krijojmë një listë simboresh – 10000 të tilla dhe do të shtojmë një funksion për të krijuar një tabelë të rastësishme. Në konsolën e dytë:

syms:`IBM`AAPL`GOOG,-9997?`8
rnd:{[n;t] ([] sym:n?syms; time:t+asc n#til 25; price:n?10f; size:n?10)}

Shtova në listën e simboleve tre prej tyre të vërtetë, për ta bërë më të lehtë gjetjen e tyre në tabelë. Funksioni rnd krijon një tabelë të rastësishme me n rreshta, ku koha ndryshon nga t në t+25 milisekonda.

Tani mund të provojmë të dërgojmë të dhëna në shërbim (do të shtojmë orët e para të dhjeta):

{h (`upd;`trade;rnd[10000;x])} each `time$00:00 + til 60*10

Mund të kontrolloni në shërbim nëse tabela u përditësua:

c 25 200
select from tradeAgg where sym=`AAPL
-20#select from tradeAgg where sym=`AAPL

Rezultati:

sym|time|high|low|firstPrice|lastPrice|firstSize|lastSize|numTrades|volume|pvolume|turnover|avgPrice|avgSize|vwap|cumVolume
--|--|--|--|--|--------------------------------
AAPL|09:27|9.258904|9.258904|9.258904|9.258904|8|8|1|8|9.258904|74.07123|9.258904|8|9.258904|2888
AAPL|09:28|9.068162|9.068162|9.068162|9.068162|7|7|1|7|9.068162|63.47713|9.068162|7|9.068162|2895
AAPL|09:31|4.680449|0.2011121|1.620827|0.2011121|1|5|4|14|9.569556|36.84342|2.392389|3.5|2.631673|2909
AAPL|09:33|2.812535|2.812535|2.812535|2.812535|6|6|1|6|2.812535|16.87521|2.812535|6|2.812535|2915
AAPL|09:34|5.099025|5.099025|5.099025|5.099025|4|4|1|4|5.099025|20.3961|5.099025|4|5.099025|2919

Tani do t’i bëjmë një test të ngarkesës për të zbuluar se sa të dhëna shërbimi mund të përpunojë në minutë. Ju kujtoj se kemi caktuar një interval për përditësime prej 25 milisekondash. Sipas kësaj, shërbimi duhet (në mesatare) të përfundojë të paktën në 20 milisekonda për përditësim, për të lënë kohë për përdoruesit që të kërkojnë të dhënat. Ju lutem, futni të dhënat e mëposhtme në procesin e dytë:

tm:10:00:00.000
stressTest:{[n] 1 string[tm]," "; times,::h ({st:.z.T; upd[`trade;x]; .z.T-st};rnd[n;tm]); tm+:25}
start:{[n] times::(); do[4800;stressTest[n]]; -1 " "; `min`avg`med`max!(min times;avg times;med times;max times)}

4800 – këto janë dy minuta. Mund të provojmë të fillojmë së pari me 1000 rreshta çdo 25 milisekonda:

start 1000

Në rastin tim, rezultati është rreth disa milisekondash për përditësim. Kështu që menjëherë do ta rris numrin e rreshtave në 10,000:

start 10000

Rezultati:

min| 00:00:00.004
avg| 9.191458
med| 9f
max| 00:00:00.030

Sërish asgjë e veçantë, megjithatë kjo është 24 milion rreshta në minutë, 400 mijë në sekondë. Më shumë se 25 milisekonda përditësimi u ndal vetëm 5 herë, duket se gjatë kalimit të minutës. Të rrisim në 100,000:

start 100000

Rezultati:

min| 00:00:00.013
avg| 25.11083
med| 24f
max| 00:00:00.108
q)sum times
00:02:00.532

Siç e shohim, shërbimi është ngushtë duke u përballur, por megjithatë arrin të mbajë një nivel. Një volum i tillë të dhënash (240 milion rreshta në minutë) është jashtëzakonisht i madh, në raste të tilla është zakonshme të nisen disa klonë (ose madje dhjetëra klonë) të shërbimit, secili prej të cilëve përpunon vetëm një pjesë të simboleve. Megjithatë, rezultati është mbresëlënës për një gjuhë interpretuese, e cila është kryesisht e orientuar ndaj ruajtjes së të dhënave.

Mund të lindë pyetja, pse koha rritet në mënyrë jo lineare me madhësinë e çdo përditësimi. Arsyeja është se funksioni i kompresimit - është faktikisht funksioni C, i cili punon shumë më efektivisht se updateAgg. Duke filluar nga një madhësi e caktuar përditësimi (në rreth 10.000), updateAgg arrin kufirin e tij dhe më tej koha e ekzekutimit të tij nuk varet nga madhësia e përditësimit. Pikërisht për shkak të hapit paraprak Q, shërbimi është në gjendje të procesojë kështu volumeve të mëdha të të dhënave. Kjo thekson sa e rëndësishme është, gjatë punës me të dhëna të mëdha, të zgjidhni algoritmin e duhur. Një pikë tjetër është ruajtja e duhur e të dhënave në memorie. Nëse të dhënat do të ruheshin jo kolonalisht ose nuk ishin të renditura sipas kohës, atëherë do të njihemi me një fenomen të tillë si humbja e keshit TLB - mungesa e adresës së faqes së memories në keshin e adresave të procesorit. Kërkimi i adresës merr diku rreth 30 herë më shumë kohë në rastin e dështimit dhe në rastin e të dhënave të shpërndara mund të ngadalësojë shërbimin disa herë.

Përfundimi

Në këtë artikull, tregova se baza KDB+ dhe gjuha Q janë të përshtatshme jo vetëm për ruajtjen e të dhënave të mëdha dhe aksesin e thjeshtë në to përmes selektimit, por edhe për krijimin e shërbimeve të përpunimit të të dhënave që mund të përballojnë qindra milion rreshta/gigabajt të dhënash, edhe në një proces të vetëm Q. Gjuha Q lejon implementimin shumë të shkurtër dhe efikas të algoritmeve të lidhura me përpunimin e të dhënave përmes natyrës së saj vektoriale, interpretuesit të brendshëm të dialektit SQL dhe një grup të suksesshëm funksionesh biblioteke.

Dua të theksoj se ajo që përmenda më sipër është vetëm një pjesë e mundësive të Q, ai ka edhe karakteristika të tjera unike. Për shembull, një protokoll IPC jashtëzakonisht i thjeshtë, i cili fshin kufirin ndërmjet proceseve të veçanta Q dhe lejon kombinimin e qindra këtyre proceseve në një rrjet të vetëm, që mund të ndodhet në dhjetëra servera në skajet e ndryshme të botës.

Burimi: habr.com

Bli një hosting të besueshëm për faqet me mbrojtje DDoS, VPS VDS serverë 🔥 Bli një hosting të besueshëm për faqet me mbrojtje DDoS, VPS VDS serverë | ProHoster