O HotXLS, a biblioteca de componentes Excel nativa para Delphi e C++Builder, infla várias folhas de trabalho XLSX ao mesmo tempo a partir de um único pacote ZIP aberto. O mecanismo é TZipReadGate, uma pequena classe em lxZipArchive.pas que detém o stream do pacote mais uma secção crítica e expõe exatamente um método. Serializa o par seek-and-read. Tudo acima desse par corre em concorrência
O problema que forçou este desenho é um que todo o programador Delphi que já abriu um livro de trabalho grande já encontrou. Um xlsx de 80 MB são 80 MB de XML deflacionado, e as partes de folha de trabalho lá dentro expandem entre cinco a dez vezes. Se o seu percurso de abertura extrai cada folha de trabalho para um stream de memória antes de a analisar, paga pelos bytes inflacionados em cima do livro de trabalho que está a construir, e o pico chega antes de uma única célula ter sido criada. Este artigo é sobre a concorrência ao nível do pacote que remove esse passo intermédio. O teto do alocador que se situa acima disso está coberto em o artigo sobre análise paralela de XLSX e o gestor de memória, e a API de ler-uma-vez, nunca-materializar está coberta em o percurso pelo leitor direto de streaming
Por que armazenava o percurso de abertura antigo cada folha de trabalho em RAM
A abertura paralela original no HotXLS era um pipeline de três fases, e a fase intermédia era a única que corria em workers. A Fase A percorria a lista de folhas em série, criava cada folha de trabalho, lia a sua parte de relação, e copiava o XML de folha de trabalho inflacionado inteiro para um TMemoryStream privado. A Fase B distribuía ParseWorksheetXml pela pool. A Fase C voltava ao arquivo na thread de chamada para as pequenas partes satélite: comentários, comentários encadeados, desenhos, gráficos, tabelas. Essa forma foi escolhida por uma razão declarada. O comentário de cabeçalho em lxParallelParse.pas costumava dizer, em tantas palavras, que o arquivo zip e o seu estado de inflação não são thread-safe, e as notas internas iam mais longe: não se dê ao trabalho de bloquear o arquivo, porque uma vez que o estado de inflação seja serializado por entrada, o bloqueio não compra nada. A Fase A existia para manter cada toque ao arquivo numa única thread. O custo era que um livro de trabalho com oito folhas ativas mantinha oito buffers de XML de folha de trabalho totalmente inflacionados em memória simultaneamente, e esses buffers são os maiores objetos transitórios em todo o percurso de abertura
Podem duas threads inflar a partir de um stream ZIP?
Sim, e o julgamento antigo estava errado de uma forma específica e localizável: colapsava dois pedaços diferentes de estado numa frase. O estado de inflação é genuinamente não partilhável. Um z_stream zlib transporta a janela deslizante, as tabelas de Huffman e a posição de bits de um membro comprimido, e duas threads a empurrar bytes através do mesmo produzem lixo. A fonte de bytes subjacente é uma questão inteiramente diferente, e a resposta aí é que um stream de ficheiro tem exatamente um pedaço de estado mutável partilhado que vale a pena proteger, o seu cursor de posição
O contentor ZIP torna a separação legal. Cada membro num arquivo ZIP é comprimido independentemente: o seu próprio cabeçalho de ficheiro local, o seu próprio stream de bits deflate no seu próprio DataOffset, o seu próprio CRC32 e tamanhos no diretório central. Não há dicionário partilhado a abranger membros da forma que um bloco 7z sólido tem, pelo que a entrada N pode ser inflada sem tocar na entrada M. Dê a cada worker o seu próprio z_stream sobre o seu próprio intervalo de bytes e a única coisa em que colidem é o seek. Essa colisão é o que TZipReadGate remove, e a classe inteira é suficientemente curta para se ler num único ecrã
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;
O que protege TZipReadGate e o que deliberadamente não protege
TZipReadGate.ReadAt protege uma operação indivisível, posicionar o stream partilhado e ler dele, e nada mais. TZipArchive.OpenArchive constrói o gate sobre FInputStream assim que o diretório central é analisado com sucesso, e TZipArchive.Close liberta-o. Os arquivos abertos para escrita nunca recebem um. Cada leitura que um worker executa no pacote passa por isso através de uma única secção crítica mantida pela duração de uma leitura em buffer
Tudo o resto fica fora do bloqueio porque já é privado ou já é imutável. TZipSubStream mantém o seu próprio FPosition, pelo que cada worker acompanha o seu próprio lugar na sua própria entrada. O TZLibStream que TZipEntry.GetStream constrói sobre esse substream é por entrada, criado com windowBits de -15 para deflate em bruto, e nunca partilhado. O diretório central está totalmente analisado antes de qualquer worker começar, incluindo cada cabeçalho local, pelo que GetEntryByName é uma pesquisa de hash apenas de leitura no momento em que a concorrência começa. O encaminhamento em si são três linhas em TZipSubStream.Read, e o ramo sem gate é o que mantém cada chamador de thread única existente no caminho de código antigo
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 custa o gate sob contenção?
Menos do que a frase "bloqueio global no arquivo" sugere, por causa da granularidade que TZLibStream calha de usar. O seu buffer de entrada é BufferSize, definido como $4000, pelo que ReadInputBuffer puxa 16 KB de bytes comprimidos por reabastecimento e entrega-os a zng_inflate. Uma aquisição de bloqueio cobre por isso 16 KB de entrada deflate, que para XML de folha de trabalho expande para algo na ordem de 100 KB de marcação que o worker depois descodifica e analisa sem manter nada. O bloqueio é mantido para uma leitura posicionada contra uma cache do sistema operativo; o trabalho que protege é medido em milissegundos
A fronteira honesta é onde esse rácio se inverte. As entradas armazenadas em vez de deflacionadas passam pelo gate um para um sem trabalho de inflação para esconder a latência, pelo que um pacote cheio de membros armazenados serializaria muito mais rigidamente. Um ficheiro frio em suporte lento alarga a secção crítica, porque a leitura lá dentro é agora uma transferência de disco real em vez de um acerto de cache. E para lá de um punhado de workers o gate não é o que se atinge primeiro de qualquer forma: a análise de folha de trabalho é intensiva em alocação, e o gestor de memória do Delphi serializa alocações entre threads bem antes de o read gate se tornar a restrição. É por isso que TXLSXWorkbook.ParallelParseThreads tem por predefinição um teto automático em vez de uma thread por núcleo
O corpo do worker, e o ciclo de drenagem fácil de esquecer
Com o gate no lugar, o HotXLS eliminou a preparação da Fase A por completo. O worker agora abre o seu próprio stream de entrada e alimenta-o diretamente ao analisador. Dois campos transitórios transportam as entradas: FParZip mantém o arquivo durante a fase paralela, FParSheetPartNames mantém os nomes das partes, e ambos são limpos no bloco finally para que nenhum ponteiro obsoleto sobreviva a uma abertura falhada. O stream que volta de TZipArchive.OpenFile é um TZipVerifiedStream a envolver um TZLibStream a envolver um TZipSubStream, e libertar o exterior liberta a cadeia
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;
O ciclo de drenagem é o detalhe que uma migração direta do código antigo deixaria cair, e deixá-lo cair desativa silenciosamente a verificação de integridade. TZipVerifiedStream acumula um CRC32 em execução à medida que os bytes passam, e só chama VerifyComplete quando a sua posição alcança o tamanho descomprimido registado no diretório central; é daí que vêm as exceções de discrepância de tamanho e de discrepância de CRC32, mais uma leitura de sonda de um byte que apanha uma entrada mais longa do que a declarada. Um leitor de XML para na etiqueta de fecho e normalmente deixa por ler uma nova linha ou alguns bytes de espaço em branco final, pelo que sem a drenagem a posição nunca alcança o tamanho declarado e as verificações nunca disparam. Ler o remanescente para um buffer de rascunho não custa nada e restaura-as. Quando os streams de preparação existiam, XlsxCopyStreamAll fazia isto por acidente
O que ainda corre em série, e a flag que desliga tudo
A Fase A sobrevive, menos a extração. Ainda cria cada folha de trabalho e lê as suas relações na thread de chamada, o que é o que deixa cada mapa partilhado imutável assim que os workers começam. A Fase C ainda percorre as folhas em série depois disso para comentários, desenhos, gráficos e tabelas, e a sua proteção mudou de uma verificação nula no antigo array de preparação para zip.Exists contra o nome da parte. As entradas partilhadas só de leitura que os workers tocam, a tabela de strings partilhada e os mapas cellXf, estão completas antes de a Fase B começar e nunca são escritas durante ela
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;
Definir ParallelParse como False antes de Open despacha o mesmo procedimento de tarefa com uma contagem de threads de um, e RunParallelJobs degenera para um ciclo simples na thread de chamada. Vale a pena saber isso por duas razões: é a resposta de uma linha se alguma vez surgir uma preocupação de threading em campo, e significa que os percursos série e paralelo partilham um único corpo de código de análise em vez de divergirem. As exceções dos workers são capturadas, o índice de tarefa mais baixo vence, e o erro é relançado na thread de chamada depois de todos os workers se juntarem, pelo que uma folha de trabalho corrompida ainda surge como uma exceção no local esperado. A afinação geral do percurso de abertura circundante está coberta em o guia sobre desempenho de livros de trabalho grandes em Delphi
O read gate, a fase de abertura paralela e o acesso de entrada em streaming aqui descritos fazem parte do componente HotXLS Excel standard para Delphi e C++Builder, com código-fonte completo; a página do produto contém a referência completa de TXLSXWorkbook, incluindo as propriedades de abertura paralela