HotXLS, нативная библиотека компонентов Excel для Delphi и C++Builder, распаковывает несколько листов XLSX одновременно из одного открытого ZIP-пакета. Механизм — TZipReadGate, небольшой класс в lxZipArchive.pas, хранящий поток пакета плюс одну критическую секцию и предоставляющий ровно один метод. Он сериализует пару seek-и-read. Всё, что выше этой пары, выполняется параллельно
Проблема, вынудившая к этому решению, знакома каждому разработчику на Delphi, кто открывал большую книгу. 80 МБ xlsx — это 80 МБ сжатого deflate XML, а листовые части внутри него расширяются примерно в пять-десять раз. Если ваш путь открытия извлекает каждый лист в поток памяти перед разбором, вы платите за распакованные байты сверх строящейся книги, и пик приходит до того, как создана хотя бы одна ячейка. Эта статья посвящена именно параллелизму на уровне пакета, устраняющему этот промежуточный этап. Потолок аллокатора, стоящий над ним, разобран в статье о параллельном разборе XLSX и менеджере памяти, а API однократного чтения без материализации разобран в статье о потоковом прямом читателе
Почему старый путь открытия держал каждый лист в ОЗУ
Исходное параллельное открытие в HotXLS было трёхфазным конвейером, и только средняя фаза выполнялась на воркерах. Фаза A последовательно обходила список листов, создавала каждый лист, читала его часть отношений и копировала весь распакованный XML листа в приватный TMemoryStream. Фаза B развёртывала ParseWorksheetXml по пулу. Фаза C возвращалась к архиву на вызывающем потоке за мелкими спутниковыми частями: комментарии, потоковые комментарии, рисунки, диаграммы, таблицы. Эта форма была выбрана по заявленной причине. Комментарий в заголовке lxParallelParse.pas раньше буквально гласил, что архив zip и его состояние распаковки не потокобезопасны, а внутренние заметки шли дальше: не стоит утруждать себя блокировкой архива, потому что как только состояние распаковки сериализовано на элемент, блокировка ничего не покупает. Фаза A существовала, чтобы держать каждое обращение к архиву на одном потоке. Ценой было то, что книга с восемью активными листами держала в памяти одновременно восемь полностью распакованных буферов XML листов, и эти буферы — самые крупные временные объекты во всём пути открытия
Могут ли два потока распаковывать из одного потока ZIP?
Да, и старое суждение было ошибочным конкретным, локализуемым образом: оно свело два разных вида состояния в одно предложение. Состояние распаковки действительно не разделяемо. zlib z_stream несёт скользящее окно, таблицы Хаффмана и битовую позицию для одного сжатого элемента, и два потока, проталкивающие байты через один и тот же поток, производят мусор. Лежащий в основе источник байт — совершенно другой вопрос, и ответ здесь в том, что у файлового потока есть ровно один кусок изменяемого разделяемого состояния, заслуживающий защиты, — его курсор позиции
Контейнер ZIP делает это разделение законным. Каждый элемент в архиве ZIP сжат независимо: собственный локальный заголовок файла, собственный битовый поток deflate по собственному DataOffset, собственные CRC32 и размеры в центральном каталоге. Здесь нет общего словаря, охватывающего элементы, как в цельном блоке 7z, так что элемент N можно распаковать, не касаясь элемента M. Дайте каждому воркеру собственный z_stream над собственным диапазоном байт, и единственное, на чём они столкнутся, — это seek. Именно это столкновение устраняет TZipReadGate, и весь класс достаточно короток, чтобы прочесть на одном экране
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;
Что защищает TZipReadGate и что он намеренно не защищает
TZipReadGate.ReadAt охраняет одну неделимую операцию — позиционирование общего потока и чтение из него — и ничего больше. TZipArchive.OpenArchive строит шлюз над FInputStream, как только центральный каталог успешно разобран, а TZipArchive.Close его освобождает. Архивы, открытые для записи, никогда его не получают. Каждое чтение, которое воркер выполняет из пакета, поэтому проходит через единую критическую секцию, удерживаемую на время одного буферизованного чтения
Всё остальное остаётся вне блокировки, потому что уже приватно или уже неизменяемо. TZipSubStream хранит собственную FPosition, так что каждый воркер отслеживает своё место в своей записи. TZLibStream, который TZipEntry.GetStream строит поверх этого подпотока, отдельный для каждой записи, создаётся с windowBits равным -15 для raw deflate, и никогда не разделяется. Центральный каталог полностью разобран до начала работы любого воркера, включая каждый локальный заголовок, так что GetEntryByName к началу параллелизма — это чисто читающий поиск по хэшу. Сама маршрутизация — три строки в TZipSubStream.Read, и ветка без шлюза — то, что удерживает весь существующий однопоточный вызывающий код на старом пути кода
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;
Сколько стоит шлюз при конкуренции?
Меньше, чем предполагает фраза «глобальная блокировка архива», из-за гранулярности, которую случайно использует TZLibStream. Его входной буфер — BufferSize, определённый как $4000, так что ReadInputBuffer при каждом пополнении вытягивает 16 КБ сжатых байт и передаёт их в zng_inflate. Один захват блокировки поэтому покрывает 16 КБ входа deflate, что для XML листа расширяется примерно до 100 КБ разметки, которую воркер затем декодирует и разбирает, ничего не удерживая. Блокировка удерживается на время позиционированного чтения из кэша операционной системы; работа, которую она стережёт, измеряется миллисекундами
Честная граница — там, где это соотношение переворачивается. Записи, хранящиеся без сжатия, а не в deflate, читаются через шлюз один к одному без работы распаковки, скрывающей задержку, так что пакет, полный хранимых элементов, сериализовался бы гораздо жёстче. Холодный файл на медленном носителе расширяет критическую секцию, потому что чтение внутри неё теперь настоящая передача с диска, а не попадание в кэш. И за пределами горстки воркеров шлюз в любом случае не первое, во что вы упрётесь: разбор листа тяжёл по выделениям, и менеджер памяти Delphi сериализует выделения между потоками задолго до того, как шлюз чтения станет ограничением. Именно поэтому TXLSXWorkbook.ParallelParseThreads по умолчанию использует автоматический потолок вместо одного потока на ядро
Тело воркера и цикл дочитывания, который легко забыть
С шлюзом на месте HotXLS полностью удалил постановку Фазы A. Теперь воркер открывает собственный поток записи и подаёт его прямо в парсер. Два временных поля несут входные данные: FParZip хранит архив на время параллельной фазы, FParSheetPartNames хранит имена частей, и оба очищаются в блоке finally, так что ни один устаревший указатель не переживает неудавшееся открытие. Поток, возвращаемый TZipArchive.OpenFile, — это TZipVerifiedStream, оборачивающий TZLibStream, оборачивающий TZipSubStream, и освобождение внешнего освобождает всю цепочку
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;
Цикл дочитывания — та деталь, которую прямой перенос старого кода упустил бы, а её незаметная потеря молча отключает проверку целостности. TZipVerifiedStream накапливает текущий CRC32 по мере прохождения байт и вызывает VerifyComplete только тогда, когда его позиция достигает несжатого размера, записанного в центральном каталоге; именно оттуда берутся исключения о несовпадении размера и несовпадении CRC32, плюс чтение-зонд в один байт, ловящее запись длиннее заявленной. Читатель XML останавливается на закрывающем элементе и обычно оставляет непрочитанным перевод строки или несколько байт завершающих пробелов, так что без дочитывания позиция никогда не достигает заявленного размера, и проверки никогда не срабатывают. Чтение остатка в черновой буфер ничего не стоит и восстанавливает их. Когда существовали промежуточные потоки, XlsxCopyStreamAll делал это случайно
Что по-прежнему выполняется последовательно и флаг, отключающий всё это
Фаза A выживает, за вычетом извлечения. Она по-прежнему создаёт каждый лист и читает его отношения на вызывающем потоке, что и оставляет каждую разделяемую карту неизменяемой к моменту старта воркеров. Фаза C по-прежнему обходит листы последовательно после этого за комментариями, рисунками, диаграммами и таблицами, а её защита изменилась с проверки на null в старом массиве постановки на zip.Exists по имени части. Разделяемые только для чтения входные данные, которых касаются воркеры, разделяемая таблица строк и карты cellXf, полностью готовы до начала Фазы B и никогда не записываются во время неё
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;
Установка ParallelParse в False перед Open запускает ту же процедуру задания с числом потоков, равным единице, и RunParallelJobs вырождается в обычный цикл на вызывающем потоке. Это стоит знать по двум причинам: это однострочный ответ, если проблема потоков когда-либо всплывёт на практике, и это означает, что последовательный и параллельный пути разделяют единое тело кода разбора, а не расходятся. Исключения воркеров перехватываются, побеждает наименьший индекс задания, и ошибка заново вызывается на вызывающем потоке после присоединения каждого воркера, так что повреждённый лист всё равно проявляется как одно исключение в ожидаемом месте. Общая настройка окружающего пути открытия разобрана в руководстве по производительности больших книг в Delphi
Описанные здесь шлюз чтения, фаза параллельного открытия и потоковый доступ к записям поставляются как часть стандартного компонента HotXLS Excel для Delphi и C++Builder с полным исходным кодом; страница продукта содержит полный справочник TXLSXWorkbook, включая свойства параллельного открытия