HotXLS, biblioteca de componente Excel nativă pentru Delphi și C++Builder, inflatează mai multe foi de calcul XLSX în același timp dintr-un singur pachet ZIP deschis. Mecanismul este TZipReadGate, o clasă mică din lxZipArchive.pas care păstrează fluxul pachetului plus o singură secțiune critică și expune exact o metodă. Serializează perechea seek-and-read. Tot ce se află deasupra acelei perechi rulează concurent
Problema care a forțat acest design este una pe care a întâlnit-o orice dezvoltator Delphi care a deschis un registru de calcul mare. Un xlsx de 80 MB este 80 MB de XML deflate, iar părțile de foaie de calcul din interior se expandează aproximativ de cinci până la zece ori. Dacă calea dumneavoastră de deschidere extrage fiecare foaie de calcul într-un flux de memorie înainte de a o analiza, plătiți pentru octeții inflatați pe deasupra registrului de calcul pe care îl construiți, iar vârful sosește înainte ca o singură celulă să fi fost creată. Acest articol tratează concurența la nivel de pachet care elimină acel pas de staging. Plafonul alocatorului care stă deasupra este tratat în articolul despre analiza paralelă XLSX și managerul de memorie, iar API-ul de citire o dată, niciodată materializat, este tratat în prezentarea cititorului direct în streaming
De ce vechea cale de deschidere puse fiecare foaie de calcul în RAM
Deschiderea paralelă originală în HotXLS era un pipeline în trei faze, iar faza de mijloc era singura care rula pe workeri. Faza A parcurgea lista de foi serial, crea fiecare foaie de calcul, îi citea partea de relații și copia întregul XML de foaie de calcul inflatat într-un TMemoryStream privat. Faza B distribuia ParseWorksheetXml pe pool. Faza C se întorcea la arhivă pe firul apelant pentru părțile satelit mici: comentarii, comentarii cu fire, desene, diagrame, tabele. Acea formă a fost aleasă dintr-un motiv declarat. Comentariul de antet de pe lxParallelParse.pas obișnuia să spună, aproape textual, că arhiva zip și starea ei de inflate nu sunt thread-safe, iar notele interne mergeau mai departe: nu vă deranjați să blocați arhiva, deoarece odată ce starea de inflate este serializată per intrare, blocajul nu aduce nimic. Faza A exista pentru a păstra fiecare atingere a arhivei pe un singur fir. Costul era că un registru de calcul cu opt foi ocupate ținea opt buffer-e XML de foaie de calcul complet inflatate în memorie simultan, iar acele buffer-e sunt cele mai mari obiecte tranzitorii din întreaga cale de deschidere
Pot două fire să inflateze dintr-un singur flux ZIP?
Da, iar judecata veche era greșită într-un mod specific, localizabil: colapsa două piese diferite de stare într-o singură propoziție. Starea de inflate chiar nu este partajabilă. Un z_stream zlib poartă fereastra glisantă, tabelele Huffman și poziția de bit pentru un membru comprimat, iar două fire care împing octeți prin același produc gunoi. Sursa de octeți subiacentă este o întrebare complet diferită, iar răspunsul acolo este că un flux de fișier are exact o singură piesă de stare partajată mutabilă care merită protejată, cursorul său de poziție
Containerul ZIP face separarea legală. Fiecare membru dintr-o arhivă ZIP este comprimat independent: propriul antet de fișier local, propriul flux de biți deflate la propriul DataOffset, propriile CRC32 și dimensiuni din directorul central. Nu există niciun dicționar partajat care se întinde peste membri așa cum are un bloc 7z solid, așa că intrarea N poate fi inflatată fără a atinge intrarea M. Dați fiecărui worker propriul z_stream peste propriul interval de octeți, iar singurul lucru pe care se ciocnesc este seek-ul. Acea coliziune este ceea ce elimină TZipReadGate, iar întreaga clasă este suficient de scurtă pentru a fi citită într-un singur ecran
type
TZipReadGate = class
private
FBaseStream: TStream;
FLock: TRTLCriticalSection;
public
constructor Create(ABaseStream: TStream);
destructor Destroy; override;
function ReadAt(AOffset: Int64; var Buffer; Count: Longint): Longint;
end;
function TZipReadGate.ReadAt(AOffset: Int64; var Buffer;
Count: Longint): Longint;
begin
if Count <= 0 then
begin
Result := 0;
Exit;
end;
EnterCriticalSection(FLock);
try
FBaseStream.Position := AOffset;
Result := FBaseStream.Read(Buffer, Count);
finally
LeaveCriticalSection(FLock);
end;
end;
Ce protejează TZipReadGate și ce nu protejează deliberat
TZipReadGate.ReadAt păzește o singură operație indivizibilă, poziționarea fluxului partajat și citirea din el, și nimic altceva. TZipArchive.OpenArchive construiește poarta peste FInputStream odată ce directorul central s-a analizat cu succes, iar TZipArchive.Close o eliberează. Arhivele deschise pentru scriere nu primesc niciodată una. Fiecare citire pe care un worker o efectuează pe pachet trece deci printr-o singură secțiune critică ținută pe durata unei citiri bufferate
Tot restul rămâne în afara blocajului deoarece este deja privat sau deja imutabil. TZipSubStream își păstrează propriul FPosition, așa că fiecare worker își urmărește propriul loc în propria intrare. TZLibStream pe care TZipEntry.GetStream îl construiește peste acel sub-flux este per intrare, creat cu windowBits de -15 pentru deflate brut, și niciodată partajat. Directorul central este complet analizat înainte ca vreun worker să înceapă, inclusiv fiecare antet local, așa că GetEntryByName este o căutare hash doar-citire până când începe concurența. Rutarea în sine este trei linii în TZipSubStream.Read, iar ramura fără poartă este ceea ce păstrează fiecare apelant existent single-threaded pe vechea cale de cod
function TZipSubStream.Read(var buffer; Count: longint): longint;
var
rest: Int64;
rc: longint;
begin
rest := FSize - FPosition;
if (Count > rest) then
Count := rest;
if FReadGate <> nil then
rc := FReadGate.ReadAt(FOffset + FPosition, buffer, Count)
else
begin
FBaseStream.Position := FOffset + FPosition;
rc := FBaseStream.Read(buffer, Count);
end;
FPosition := FPosition + rc;
Result := rc;
end;
Cât costă poarta sub contenție?
Mai puțin decât sugerează expresia "blocaj global pe arhivă", datorită granularității pe care TZLibStream se întâmplă să o folosească. Buffer-ul lui de intrare este BufferSize, definit ca $4000, așa că ReadInputBuffer extrage 16 KB de octeți comprimați per reumplere și îi predă lui zng_inflate. O achiziție de blocaj acoperă deci 16 KB de intrare deflate, ceea ce pentru XML de foaie de calcul se expandează în ceva de ordinul a 100 KB de markup pe care worker-ul apoi îl decodează și analizează fără să țină nimic. Blocajul este ținut pentru o citire poziționată față de cache-ul sistemului de operare; lucrul pe care îl protejează este măsurat în milisecunde
Limita onestă este acolo unde acel raport se inversează. Intrările stocate, nu deflate, se citesc prin poartă unu-la-unu, fără niciun lucru de inflate care să ascundă latența, așa că un pachet plin de membri stocați ar serializa mult mai puternic. Un fișier rece pe media lentă lărgește secțiunea critică, deoarece citirea din interior este acum un transfer real de disc, nu o lovitură de cache. Iar dincolo de câțiva workeri, poarta nu este oricum primul lucru pe care îl loviți: analiza foilor de calcul este intensivă în alocare, iar managerul de memorie Delphi serializează alocările peste fire cu mult înainte ca poarta de citire să devină constrângerea. De aceea TXLSXWorkbook.ParallelParseThreads este implicit un plafon automat, nu un fir per nucleu
Corpul worker-ului, și bucla de golire ușor de uitat
Cu poarta în loc, HotXLS a șters complet staging-ul din Faza A. Worker-ul acum își deschide propriul flux de intrare și îl alimentează direct parserului. Două câmpuri tranzitorii poartă intrările: FParZip păstrează arhiva pe durata fazei paralele, FParSheetPartNames păstrează numele de părți, iar ambele sunt șterse în blocul finally astfel încât niciun pointer învechit nu supraviețuiește unei deschideri eșuate. Fluxul care revine din TZipArchive.OpenFile este un TZipVerifiedStream care înfășoară un TZLibStream care înfășoară un TZipSubStream, iar eliberarea celui exterior eliberează lanțul
procedure TXLSXWorkbook.ParseSheetJob(AIndex: Integer);
var
Stream: TStream;
DrainBuffer: array [0..32767] of Byte;
PartName: WideString;
begin
PartName := WideString(FParSheetPartNames[AIndex]);
Stream := FParZip.OpenFile(PartName);
if Stream = nil then
Exit;
try
ParseWorksheetXml(Stream, FSheets.ByPos[AIndex], FParSst,
FParRels[AIndex], FParFontMap, FParFillMap, FParBorderMap,
FParNumFmtMap, FParAlignMap, FParProtMap, FParDateMap);
// Consume any trailing bytes so the ZIP entry size and CRC are verified.
while Stream.Read(DrainBuffer, SizeOf(DrainBuffer)) > 0 do
;
finally
Stream.Free;
end;
end;
Bucla de golire este detaliul pe care o portare directă a codului vechi l-ar renunța, iar renunțarea la ea dezactivează tacit verificarea integrității. TZipVerifiedStream acumulează un CRC32 curent pe măsură ce octeții trec și apelează VerifyComplete doar când poziția lui ajunge la dimensiunea necomprimată înregistrată în directorul central; de acolo vin excepțiile de nepotrivire de dimensiune și de nepotrivire CRC32, plus o citire de sondare de un octet care prinde o intrare mai lungă decât cea declarată. Un cititor XML se oprește la elementul de închidere și de obicei lasă necitit un rând nou sau câțiva octeți de spațiu alb final, așa că fără golire poziția nu ajunge niciodată la dimensiunea declarată, iar verificările nu se declanșează niciodată. Citirea restului într-un buffer temporar nu costă nimic și le restaurează. Când existau fluxurile de staging, XlsxCopyStreamAll făcea asta din întâmplare
Ce mai rulează serial, și flag-ul care dezactivează totul
Faza A supraviețuiește, minus extracția. Continuă să creeze fiecare foaie de calcul și să îi citească relațiile pe firul apelant, ceea ce lasă fiecare hartă partajată imutabilă odată ce workerii încep. Faza C continuă să parcurgă foile serial ulterior pentru comentarii, desene, diagrame și tabele, iar garda ei s-a schimbat dintr-o verificare null pe vechiul array de staging la zip.Exists față de numele părții. Intrările partajate doar-citire pe care le ating workerii, tabelul de șiruri partajat și hărțile cellXf, sunt complete înainte ca Faza B să înceapă și niciodată scrise în timpul ei
var
Wb: TXLSXWorkbook;
begin
Wb := TXLSXWorkbook.Create;
try
Wb.ParallelParse := True; // default; False forces one sheet at a time
Wb.ParallelParseThreads := 4; // 0 selects the automatic cap
Wb.Open('quarterly-consolidation.xlsx');
// ... workbook is identical either way ...
finally
Wb.Free;
end;
end;
Setarea ParallelParse la False înainte de Open distribuie aceeași procedură de job cu un număr de fire de unu, iar RunParallelJobs degenerează la o simplă buclă pe firul apelant. Merită cunoscut din două motive: este răspunsul dintr-o linie dacă vreo problemă de threading apare vreodată pe teren, și înseamnă că căile serială și paralelă partajează un singur corp de cod de analiză, în loc să diveargă. Excepțiile worker-ilor sunt capturate, indexul de job cel mai mic câștigă, iar eroarea este reridicată pe firul apelant după ce fiecare worker se alătură, așa că o foaie de calcul coruptă tot apare ca o singură excepție în locul așteptat. Reglajul general al căii de deschidere înconjurătoare este tratat în ghidul pentru performanța registrelor de calcul mari în Delphi
Poarta de citire, faza de deschidere paralelă și accesul de intrare în streaming descrise aici vin ca parte a componentei Excel HotXLS standard pentru Delphi și C++Builder, cu sursă completă; pagina de produs conține referința completă TXLSXWorkbook inclusiv proprietățile de deschidere paralelă