Articolo tecnico

Inflate ZIP concorrente in Delphi: read gate di HotXLS

HotXLS, la libreria di componenti Excel nativa per Delphi e C++Builder, decomprime diversi fogli di lavoro XLSX contemporaneamente da un unico pacchetto ZIP aperto. Il meccanismo è TZipReadGate, una piccola classe in lxZipArchive.pas che mantiene lo stream del pacchetto più una sezione critica ed espone esattamente un metodo. Serializza la coppia seek-and-read. Tutto ciò che sta sopra quella coppia gira concorrentemente

Il problema che ha forzato questo design è uno che ogni sviluppatore Delphi che abbia aperto un workbook grande ha incontrato. Un xlsx da 80 MB sono 80 MB di XML deflated, e le parti dei fogli di lavoro al suo interno si espandono di circa cinque-dieci volte. Se il tuo percorso di apertura estrae ogni foglio in uno stream di memoria prima di analizzarlo, paghi per i byte decompressi in cima al workbook che stai costruendo, e il picco arriva prima che sia stata creata anche una sola cella. Questo articolo tratta la concorrenza a livello di pacchetto che rimuove quel passaggio di staging. Il tetto dell'allocatore che sta sopra è trattato in l'articolo sul parsing XLSX parallelo e il memory manager, e l'API leggi-una-volta-mai-materializzare è trattata in la panoramica del reader diretto in streaming

Perché il vecchio percorso di apertura metteva in staging ogni foglio in RAM

La vecchia apertura parallela in HotXLS era una pipeline a tre fasi, e la fase intermedia era l'unica che girava sui worker. La fase A percorreva l'elenco dei fogli serialmente, creava ogni foglio, leggeva la sua parte di relazione, e copiava l'intero XML del foglio decompresso in un TMemoryStream privato. La fase B distribuiva ParseWorksheetXml sul pool. La fase C tornava all'archivio sul thread chiamante per le piccole parti satellite: commenti, commenti threaded, disegni, grafici, tabelle. Quella forma fu scelta per un motivo dichiarato. Il commento in testa a lxParallelParse.pas diceva, letteralmente, che l'archivio zip e il suo stato di inflate non sono thread-safe, e le note interne andavano oltre: non preoccuparti di bloccare l'archivio, perché una volta che lo stato di inflate è serializzato per voce, il lock non compra nulla. La fase A esisteva per mantenere ogni tocco all'archivio su un solo thread. Il costo era che un workbook con otto fogli attivi teneva otto buffer XML di foglio completamente decompressi in memoria simultaneamente, e quei buffer sono gli oggetti transitori più grandi nell'intero percorso di apertura

Possono due thread decomprimere da un unico stream ZIP?

Sì, e il vecchio giudizio era sbagliato in un modo specifico e localizzabile: collassava due pezzi diversi di stato in una sola frase. Lo stato di inflate è genuinamente non condivisibile. Uno z_stream zlib porta la finestra scorrevole, le tabelle di Huffman e la posizione di bit per un membro compresso, e due thread che spingono byte attraverso lo stesso producono spazzatura. La sorgente di byte sottostante è una questione interamente diversa, e la risposta lì è che uno stream di file ha esattamente un pezzo di stato condiviso mutabile che vale la pena proteggere, il suo cursore di posizione

Il contenitore ZIP rende legale la separazione. Ogni membro in un archivio ZIP viene compresso indipendentemente: il proprio header di file locale, il proprio bit stream deflate al proprio DataOffset, il proprio CRC32 e dimensioni nella central directory. Non c'è alcun dizionario condiviso che attraversi i membri come avviene in un blocco 7z solido, quindi la voce N può essere decompressa senza toccare la voce M. Dai a ogni worker il proprio z_stream sul proprio intervallo di byte e l'unica cosa su cui collidono è il seek. Quella collisione è ciò che TZipReadGate rimuove, e l'intera classe è abbastanza breve da leggere in una sola schermata

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;

Cosa protegge TZipReadGate e cosa deliberatamente non protegge

TZipReadGate.ReadAt protegge un'operazione indivisibile, posizionare lo stream condiviso e leggere da esso, e nient'altro. TZipArchive.OpenArchive costruisce il gate su FInputStream una volta che la central directory è stata analizzata con successo, e TZipArchive.Close lo libera. Gli archivi aperti per la scrittura non ne ottengono mai uno. Ogni lettura che un worker esegue sul pacchetto passa quindi attraverso un'unica sezione critica mantenuta per la durata di una lettura bufferizzata

Tutto il resto resta fuori dal lock perché è già privato o già immutabile. TZipSubStream mantiene il proprio FPosition, così ogni worker traccia la propria posizione nella propria voce. Lo TZLibStream che TZipEntry.GetStream costruisce su quel sub-stream è per voce, creato con windowBits di -15 per deflate grezzo, e mai condiviso. La central directory è completamente analizzata prima che qualsiasi worker inizi, compreso ogni header locale, così GetEntryByName è una ricerca hash di sola lettura nel momento in cui inizia la concorrenza. L'instradamento stesso è tre righe in TZipSubStream.Read, e il ramo senza gate è ciò che mantiene ogni chiamante single-threaded esistente sul vecchio percorso di codice

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;

Quanto costa il gate sotto contesa?

Meno di quanto suggerisca la frase "lock globale sull'archivio", grazie alla granularità che TZLibStream capita di usare. Il suo buffer di input è BufferSize, definito come $4000, così ReadInputBuffer preleva 16 KB di byte compressi per ricarica e li consegna a zng_inflate. Un'acquisizione del lock copre quindi 16 KB di input deflate, che per XML di foglio si espande in qualcosa dell'ordine di 100 KB di markup che il worker poi decodifica e analizza senza tenere nulla. Il lock viene mantenuto per una lettura posizionata contro la cache del sistema operativo; il lavoro che protegge si misura in millisecondi

Il confine onesto è dove quel rapporto si inverte. Le voci memorizzate piuttosto che deflated vengono lette attraverso il gate uno a uno senza alcun lavoro di inflate a nascondere la latenza, così un pacchetto pieno di membri memorizzati si serializzerebbe molto più duramente. Un file freddo su un supporto lento allarga la sezione critica, perché la lettura al suo interno è ora un vero trasferimento su disco piuttosto che un colpo di cache. E oltre una manciata di worker il gate non è comunque il primo ostacolo: il parsing dei fogli è pesante in allocazioni, e il memory manager Delphi serializza le allocazioni tra thread ben prima che il read gate diventi il vincolo. Ecco perché TXLSXWorkbook.ParallelParseThreads ha come predefinito un tetto automatico invece di un thread per core

Il corpo del worker, e il ciclo di drenaggio facile da dimenticare

Con il gate in posizione, HotXLS ha eliminato del tutto lo staging della fase A. Il worker ora apre il proprio stream di voce e lo alimenta direttamente al parser. Due campi transitori portano gli input: FParZip mantiene l'archivio per la durata della fase parallela, FParSheetPartNames mantiene i nomi delle parti, ed entrambi vengono azzerati nel blocco finally così che nessun puntatore obsoleto sopravviva a un'apertura fallita. Lo stream che torna da TZipArchive.OpenFile è un TZipVerifiedStream che avvolge uno TZLibStream che avvolge uno TZipSubStream, e liberare quello esterno libera l'intera catena

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;

Il ciclo di drenaggio è il dettaglio che una trasposizione diretta del vecchio codice tralascerebbe, e tralasciarlo disabilita silenziosamente il controllo di integrità. TZipVerifiedStream accumula un CRC32 corrente man mano che i byte passano e chiama VerifyComplete solo quando la sua posizione raggiunge la dimensione decompressa registrata nella central directory; è da lì che vengono le eccezioni di mismatch di dimensione e di CRC32, più una lettura di sonda di un byte che intercetta una voce più lunga di quanto dichiarato. Un lettore XML si ferma all'elemento di chiusura e di solito lascia non letti un a capo o alcuni byte di spazio bianco finale, quindi senza il drenaggio la posizione non raggiunge mai la dimensione dichiarata e i controlli non scattano mai. Leggere il resto in un buffer temporaneo non costa nulla e li ripristina. Quando esistevano gli stream di staging, XlsxCopyStreamAll faceva questo per caso

Cosa gira ancora serialmente, e il flag che disattiva tutto

La fase A sopravvive, meno l'estrazione. Crea ancora ogni foglio e legge le sue relazioni sul thread chiamante, il che è ciò che lascia ogni mappa condivisa immutabile una volta che i worker iniziano. La fase C percorre ancora i fogli serialmente in seguito per commenti, disegni, grafici e tabelle, e la sua protezione è cambiata da un controllo nullo sul vecchio array di staging a zip.Exists contro il nome della parte. Gli input condivisi di sola lettura che i worker toccano, la tabella stringhe condivisa e le mappe cellXf, sono completi prima che inizi la fase B e non vengono mai scritti durante essa

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;

Impostare ParallelParse a False prima di Open distribuisce la stessa procedura di job con un conteggio thread di uno, e RunParallelJobs degenera in un semplice ciclo sul thread chiamante. Questo vale la pena saperlo per due motivi: è la risposta di una riga se mai emerge una preoccupazione di threading sul campo, e significa che i percorsi seriale e parallelo condividono un unico corpo di codice di parsing invece di divergere. Le eccezioni dei worker vengono catturate, l'indice di job più basso vince, e l'errore viene rilanciato sul thread chiamante dopo che ogni worker si è unito, così un foglio corrotto emerge comunque come una singola eccezione nel punto atteso. La messa a punto generale del percorso di apertura circostante è trattata in la guida alle prestazioni dei workbook di grandi dimensioni in Delphi

Il read gate, la fase di apertura parallela e l'accesso alle voci in streaming descritti qui sono distribuiti come parte del componente Excel HotXLS standard per Delphi e C++Builder, con codice sorgente completo; la pagina prodotto porta il riferimento completo di TXLSXWorkbook incluse le proprietà di apertura parallela