ifcparse: build the lazy index in parallel

On a file large enough the DATA section is split at the same boundaries
the parallel parse uses (chunk_bounds(), now shared) and each chunk is
indexed by its own worker with its own paged reader, lexer, shells,
offsets, GlobalIds and inverse records; the results are merged in file
order, so instance order, GlobalId precedence and inverse records are
identical to the serial index. The serial index is the same code run on
one chunk. The name table is reserved before the merge, which also
helps the serial case. The default thread count is the one the full
parse uses.

TXG 58 MB / 210_King 147 MB / OKgate22 231 MB, lazy open: 12 threads
0.37 / 1.01 / 1.34 s against 0.55 / 1.59 / 2.11 s on one thread and
0.44 / 1.17 / 1.91 s for the default (parallel full) open; memory after
the open +15 to +50 MB at 12 threads for the workers' page caches. The
equality test now opens the 12 MB replicated fixture lazily with five
workers and compares it instance by instance with the serial parse.

This commit was written by an AI coding tool and has not been verified by
a human.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013wcN7XquTfUi4vsKQ4KchL
This commit is contained in:
Dion Moult
2026-09-14 18:47:06 +10:00
parent 97c83f0292
commit f24e9a6bca
2 changed files with 250 additions and 162 deletions
+242 -161
View File
@@ -1946,6 +1946,7 @@ bool ifcopenshell::file::initialize(const std::string& path, filetype ty, bool r
if (ty == FT_IFCSPF) { if (ty == FT_IFCSPF) {
storage_.emplace<1>(this, logger_.get()); storage_.emplace<1>(this, logger_.get());
header_.reset(new spf_header(this, &logger_.get())); header_.reset(new spf_header(this, &logger_.get()));
std::get<impl::in_memory_file_storage>(storage_).parse_threads = effective_parse_threads();
bool indexed = false; bool indexed = false;
if (lazy_loading_) { if (lazy_loading_) {
indexed = std::get<impl::in_memory_file_storage>(storage_).index_lazily(path, schema_, max_id_, types_to_bypass_loading_); indexed = std::get<impl::in_memory_file_storage>(storage_).index_lazily(path, schema_, max_id_, types_to_bypass_loading_);
@@ -1959,7 +1960,6 @@ bool ifcopenshell::file::initialize(const std::string& path, filetype ty, bool r
} }
} }
if (!indexed) { if (!indexed) {
std::get<impl::in_memory_file_storage>(storage_).parse_threads = effective_parse_threads();
if (paged_reading_) { if (paged_reading_) {
file_reader<paged_file_impl> s(path, 64 << 10, 64); file_reader<paged_file_impl> s(path, 64 << 10, 64);
std::get<impl::in_memory_file_storage>(storage_).read_from_stream(&s, schema_, max_id_, types_to_bypass_loading_); std::get<impl::in_memory_file_storage>(storage_).read_from_stream(&s, schema_, max_id_, types_to_bypass_loading_);
@@ -2693,6 +2693,104 @@ void for_each_instance_header(Reader& reader, spf_lexer<Reader>& lexer, size_t e
} }
namespace {
// Matches a literal byte by byte, across spans.
class literal_matcher {
const char* literal_;
size_t length_;
size_t matched_ = 0;
public:
explicit literal_matcher(const char* literal)
: literal_(literal), length_(std::strlen(literal)) {}
// True on the byte that completes the literal.
bool feed(char c) {
if (c == literal_[matched_]) {
if (++matched_ == length_) {
matched_ = 0;
return true;
}
} else {
matched_ = c == literal_[0] ? 1 : 0;
}
return false;
}
};
// The split points for a parallel pass over the DATA section: the start
// of the section, then nearest each nominal split point an instance
// boundary, then the end of the section. Fewer than three entries means
// the file cannot be split. Shared by the parallel parse and the
// parallel lazy index.
template <typename Reader>
std::vector<size_t> chunk_bounds(const Reader& source, unsigned threads) {
const size_t n = source.size();
// One pass over the bytes finds the DATA section and, nearest each
// nominal split point, an instance boundary: a '#' that starts a
// line outside any string and any comment. This is the one place
// that looks at raw bytes instead of tokens, because tokenizing the
// file serially to find the split points would leave nothing to
// parallelise. It applies three rules only: a string starts and
// ends at a quote (a doubled quote closes and reopens, which comes
// to the same thing) and cannot span a line; a comment runs from
// /* to */. Getting a string's end wrong can only lose a candidate
// boundary, never accept a wrong one, since no string contains a
// newline.
literal_matcher data_matcher("\nDATA;"), endsec_matcher("\nENDSEC");
size_t data_begin = 0, data_end = 0;
bool in_string = false, in_comment = false, newline = false;
char previous = 0;
std::vector<size_t> bounds;
size_t next_split = 0;
source.for_each_span(0, n, [&](const char* data, size_t length, size_t offset) {
for (size_t i = 0; i < length; ++i) {
const char c = data[i];
const size_t at = offset + i;
if (data_begin == 0) {
if (data_matcher.feed(c)) {
data_begin = at + 1;
bounds.push_back(data_begin);
next_split = data_begin + (n - data_begin) / threads;
}
continue;
}
if (data_end != 0) {
return;
}
if (in_comment) {
if (previous == '*' && c == '/') {
in_comment = false;
}
} else if (in_string) {
if (c == '\'' || c == '\n') {
in_string = false;
}
} else if (c == '\'') {
in_string = true;
} else if (previous == '/' && c == '*') {
in_comment = true;
} else if (endsec_matcher.feed(c)) {
data_end = at + 1 - 7;
return;
} else if (newline && c == '#' && at >= next_split && bounds.size() < threads) {
bounds.push_back(at);
next_split = data_begin + (n - data_begin) * bounds.size() / threads;
}
newline = c == '\n';
previous = c;
}
});
if (data_begin == 0) {
return {};
}
if (data_end == 0) {
data_end = n;
}
bounds.push_back(data_end);
return bounds;
}
}
struct ifcopenshell::impl::in_memory_file_storage::lazy_source { struct ifcopenshell::impl::in_memory_file_storage::lazy_source {
file_reader<paged_file_impl> reader; file_reader<paged_file_impl> reader;
spf_lexer<file_reader<paged_file_impl>> lexer; spf_lexer<file_reader<paged_file_impl>> lexer;
@@ -2781,96 +2879,161 @@ bool ifcopenshell::impl::in_memory_file_storage::index_lazily(const std::string&
this->schema = schema; this->schema = schema;
resolve_references_in_place = true; resolve_references_in_place = true;
lazy_ = true; lazy_ = true;
byref_excl_.reserve(reader.size() / 32);
// One pass over the DATA section with the tokenizer's index policy: // One pass over the DATA section with the tokenizer's index policy: the
// the instance headers through the shared loop, then the attribute list // instance headers through the shared loop, then the attribute list as
// as tokens with only the parentheses, commas and names looked at. // tokens with only the parentheses, commas and names looked at. Nothing
// Nothing is decoded. A token the tokenizer rejects, or a structure the // is decoded. On a file large enough the section is split at the same
// loop below does not expect, stops the index and the caller parses in // boundaries the parallel parse uses and each chunk is indexed by its
// full. // own worker with its own reader, lexer, shells and inverse records,
const char* failure = nullptr; // merged in file order; the serial case is the same code on one chunk.
size_t failure_offset = 0; // A token the tokenizer rejects, or a structure the loop does not
try { // expect, stops the index and the caller parses in full.
for_each_instance_header<index_tokens>(reader, lexer, reader.size(), schema, bypassed_types, lazy_bypassed_, logger_.get(), [&](uint32_t name, const ifcopenshell::declaration* declaration, size_t) { struct index_output {
const uint64_t attributes_offset = reader.tell(); std::vector<shared_pointer_type> shells;
const uint16_t type_index = (uint16_t)declaration->index_in_schema(); std::vector<std::pair<uint32_t, uint64_t>> offsets;
int depth = 1; std::vector<std::pair<size_t, std::string>> guids; // shell index, raw text
int attribute = 0; std::vector<unsigned> bypassed;
bool first_value = true; entities_by_ref inverses;
size_t guid_begin = 0, guid_end = 0; const char* failure = nullptr;
while (depth > 0) { size_t failure_offset = 0;
token t = lexer.next<attribute_tokens>(); std::exception_ptr error;
if (!t) { };
failure = "file ends inside an instance"; const auto index_chunk = [&](file_reader<paged_file_impl>& chunk_reader, spf_lexer<file_reader<paged_file_impl>>& chunk_lexer, size_t end, index_output& out) {
failure_offset = attributes_offset; try {
return false; for_each_instance_header<index_tokens>(chunk_reader, chunk_lexer, end, schema, bypassed_types, out.bypassed, logger_.get(), [&](uint32_t name, const ifcopenshell::declaration* declaration, size_t) {
} const uint64_t attributes_offset = chunk_reader.tell();
if (t.is_operator()) { const uint16_t type_index = (uint16_t)declaration->index_in_schema();
if (t.value_char == '(') { int depth = 1;
++depth; int attribute = 0;
} else if (t.value_char == ')') { bool first_value = true;
--depth; size_t guid_begin = 0, guid_end = 0;
} else if (t.value_char == ',' && depth == 1) { while (depth > 0) {
++attribute; token t = chunk_lexer.next<attribute_tokens>();
} else if (t.value_char == ';') { if (!t) {
failure = "; inside an instance"; out.failure = "file ends inside an instance";
failure_offset = t.start_pos; out.failure_offset = attributes_offset;
return false; return false;
} }
} else if (t.is_identifier()) { if (t.is_operator()) {
byref_excl_.add((uint32_t)t.as_identifier(), name, type_index, attribute); if (t.value_char == '(') {
} else if (t.type == token::Token_STRING && depth == 1 && attribute == 0 && first_value) { ++depth;
guid_begin = t.start_pos + 1; } else if (t.value_char == ')') {
guid_end = reader.tell() - 1; --depth;
} else if (t.value_char == ',' && depth == 1) {
++attribute;
} else if (t.value_char == ';') {
out.failure = "; inside an instance";
out.failure_offset = t.start_pos;
return false;
}
} else if (t.is_identifier()) {
out.inverses.add((uint32_t)t.as_identifier(), name, type_index, attribute);
} else if (t.type == token::Token_STRING && depth == 1 && attribute == 0 && first_value) {
guid_begin = t.start_pos + 1;
guid_end = chunk_reader.tell() - 1;
}
if (depth == 1) {
first_value = false;
}
chunk_lexer.reset_pool();
} }
if (depth == 1) { if (!chunk_lexer.next<attribute_tokens>().is_operator(';')) {
first_value = false; out.failure = "expected ; after )";
out.failure_offset = chunk_reader.tell();
return false;
} }
lexer.reset_pool(); out.shells.push_back(ifcopenshell::make_pointer_type<instance_data>(file, declaration, name, instance_data::lazy_tag{}));
} out.offsets.push_back({name, attributes_offset});
if (!lexer.next<attribute_tokens>().is_operator(';')) { if (guid_end > guid_begin && declaration->is(*ifcroot)) {
failure = "expected ; after )"; std::string guid;
failure_offset = reader.tell(); guid.reserve(guid_end - guid_begin);
return false; for (size_t at = guid_begin; at < guid_end; ++at) {
} guid.push_back(chunk_reader.get(at));
auto data = ifcopenshell::make_pointer_type<instance_data>(file, declaration, name, instance_data::lazy_tag{}); }
out.guids.push_back({out.shells.size() - 1, std::move(guid)});
}
return true;
});
} catch (const invalid_token_exception&) {
out.failure = "invalid token";
out.failure_offset = chunk_reader.tell();
} catch (...) {
out.error = std::current_exception();
}
};
std::vector<index_output> outputs;
constexpr size_t min_bytes_per_thread = 2u << 20;
const unsigned threads = (unsigned)std::min<size_t>(parse_threads, std::max<size_t>(1, reader.size() / min_bytes_per_thread));
const std::vector<size_t> bounds = threads > 1 ? chunk_bounds(reader, threads) : std::vector<size_t>();
if (bounds.size() >= 3) {
outputs.resize(bounds.size() - 1);
std::vector<std::thread> workers;
for (size_t k = 0; k + 1 < bounds.size(); ++k) {
outputs[k].inverses.reserve((bounds[k + 1] - bounds[k]) / 32);
workers.emplace_back([&, k]() {
file_reader<paged_file_impl> chunk_reader = reader.reopen();
chunk_reader.seek(bounds[k]);
spf_lexer<file_reader<paged_file_impl>> chunk_lexer(&chunk_reader, logger_.get());
index_chunk(chunk_reader, chunk_lexer, bounds[k + 1], outputs[k]);
});
}
for (auto& worker : workers) {
worker.join();
}
} else {
outputs.resize(1);
outputs[0].inverses.reserve(reader.size() / 32);
index_chunk(reader, lexer, reader.size(), outputs[0]);
}
for (const auto& out : outputs) {
if (out.error) {
std::rethrow_exception(out.error);
}
if (out.failure != nullptr) {
logger_.get().message(ifcopenshell::logger::LOG_NOTICE, std::string("Lazy loading not possible (") + out.failure + " at offset " + std::to_string(out.failure_offset) + "), parsing the file in full");
return false;
}
}
// Merge in file order: what the serial index did per instance.
size_t shell_count = 0;
for (const auto& out : outputs) {
shell_count += out.shells.size();
}
lazy_offsets_.reserve(shell_count);
byid_.reserve(byid_.size() + shell_count);
for (auto& out : outputs) {
for (const auto& data : out.shells) {
const uint32_t name = data->id();
if (!byid_.insert({name, data}).second) { if (!byid_.insert({name, data}).second) {
logger_.get().message(ifcopenshell::logger::LOG_WARNING, "Overwriting instance with name #" + std::to_string(name)); logger_.get().message(ifcopenshell::logger::LOG_WARNING, "Overwriting instance with name #" + std::to_string(name));
byid_.erase(name); byid_.erase(name);
byid_.insert({name, data}); byid_.insert({name, data});
} }
lazy_offsets_.push_back({name, attributes_offset}); bytype_excl_[data->declaration()].push_back(express::base(data));
express::base instance(data);
bytype_excl_[declaration].push_back(instance);
max_id = (std::max)(max_id, (unsigned int)name); max_id = (std::max)(max_id, (unsigned int)name);
if (guid_end > guid_begin && declaration->is(*ifcroot)) { }
std::string guid; lazy_offsets_.insert(lazy_offsets_.end(), out.offsets.begin(), out.offsets.end());
guid.reserve(guid_end - guid_begin); for (auto& entry : out.guids) {
for (size_t at = guid_begin; at < guid_end; ++at) { std::string& guid = entry.second;
guid.push_back(reader.get(at)); if (guid.find('\\') != std::string::npos || guid.find("''") != std::string::npos) {
} guid = ifcopenshell::decode_spf_string(guid);
if (guid.find('\\') != std::string::npos || guid.find("''") != std::string::npos) {
guid = ifcopenshell::decode_spf_string(guid);
}
std::array<char, 22> key;
if (guid_key(guid, key)) {
if (byguid_.count(key) != 0) {
logger_.get().message(ifcopenshell::logger::LOG_WARNING, "Instance encountered with non-unique GlobalId " + guid);
}
byguid_[key] = instance;
}
} }
return true; std::array<char, 22> key;
}); if (guid_key(guid, key)) {
} catch (const invalid_token_exception&) { if (byguid_.count(key) != 0) {
failure = "invalid token"; logger_.get().message(ifcopenshell::logger::LOG_WARNING, "Instance encountered with non-unique GlobalId " + guid);
failure_offset = reader.tell(); }
} byguid_[key] = express::base(out.shells[entry.first]);
if (failure != nullptr) { }
logger_.get().message(ifcopenshell::logger::LOG_NOTICE, std::string("Lazy loading not possible (") + failure + " at offset " + std::to_string(failure_offset) + "), parsing the file in full"); }
return false; lazy_bypassed_.insert(lazy_bypassed_.end(), out.bypassed.begin(), out.bypassed.end());
byref_excl_.append(std::move(out.inverses));
std::vector<shared_pointer_type>().swap(out.shells);
} }
outputs.clear();
std::sort(lazy_bypassed_.begin(), lazy_bypassed_.end()); std::sort(lazy_bypassed_.begin(), lazy_bypassed_.end());
std::sort(lazy_offsets_.begin(), lazy_offsets_.end(), [](const auto& a, const auto& b) { return a.first < b.first; }); std::sort(lazy_offsets_.begin(), lazy_offsets_.end(), [](const auto& a, const auto& b) { return a.first < b.first; });
@@ -2919,27 +3082,6 @@ void parse_chunk(const Reader& source, size_t begin, size_t end, const ifcopensh
} }
} }
// Matches a literal byte by byte, across spans.
class literal_matcher {
const char* literal_;
size_t length_;
size_t matched_ = 0;
public:
explicit literal_matcher(const char* literal)
: literal_(literal), length_(std::strlen(literal)) {}
// True on the byte that completes the literal.
bool feed(char c) {
if (c == literal_[matched_]) {
if (++matched_ == length_) {
matched_ = 0;
return true;
}
} else {
matched_ = c == literal_[0] ? 1 : 0;
}
return false;
}
};
} }
@@ -2955,68 +3097,7 @@ bool ifcopenshell::impl::in_memory_file_storage::read_instances_parallel(Reader*
return false; return false;
} }
// One pass over the bytes finds the DATA section and, nearest each const std::vector<size_t> bounds = chunk_bounds(*s, threads);
// nominal split point, an instance boundary: a '#' that starts a
// line outside any string and any comment. This is the one place
// that looks at raw bytes instead of tokens, because tokenizing the
// file serially to find the split points would leave nothing to
// parallelise. It applies three rules only: a string starts and
// ends at a quote (a doubled quote closes and reopens, which comes
// to the same thing) and cannot span a line; a comment runs from
// /* to */. Getting a string's end wrong can only lose a candidate
// boundary, never accept a wrong one, since no string contains a
// newline.
literal_matcher data_matcher("\nDATA;"), endsec_matcher("\nENDSEC");
size_t data_begin = 0, data_end = 0;
bool in_string = false, in_comment = false, newline = false;
char previous = 0;
std::vector<size_t> bounds;
size_t next_split = 0;
s->for_each_span(0, n, [&](const char* data, size_t length, size_t offset) {
for (size_t i = 0; i < length; ++i) {
const char c = data[i];
const size_t at = offset + i;
if (data_begin == 0) {
if (data_matcher.feed(c)) {
data_begin = at + 1;
bounds.push_back(data_begin);
next_split = data_begin + (n - data_begin) / threads;
}
continue;
}
if (data_end != 0) {
return;
}
if (in_comment) {
if (previous == '*' && c == '/') {
in_comment = false;
}
} else if (in_string) {
if (c == '\'' || c == '\n') {
in_string = false;
}
} else if (c == '\'') {
in_string = true;
} else if (previous == '/' && c == '*') {
in_comment = true;
} else if (endsec_matcher.feed(c)) {
data_end = at + 1 - 7;
return;
} else if (newline && c == '#' && at >= next_split && bounds.size() < threads) {
bounds.push_back(at);
next_split = data_begin + (n - data_begin) * bounds.size() / threads;
}
newline = c == '\n';
previous = c;
}
});
if (data_begin == 0) {
return false;
}
if (data_end == 0) {
data_end = n;
}
bounds.push_back(data_end);
if (bounds.size() < 3) { if (bounds.size() < 3) {
return false; return false;
} }
@@ -570,6 +570,12 @@ TEST_CASE("Parallel and paged parsing yield the same instances, attributes, inve
paged_parallel.paged_reading(true); paged_parallel.paged_reading(true);
paged_parallel.parse_threads(5); paged_parallel.parse_threads(5);
REQUIRE(paged_parallel.initialize(path.string())); REQUIRE(paged_parallel.initialize(path.string()));
// And the lazy index built by 5 workers.
ifcopenshell::file lazy_parallel(ifcopenshell::uninitialized_tag{});
lazy_parallel.lazy_loading(true);
lazy_parallel.parse_threads(5);
REQUIRE(lazy_parallel.initialize(path.string()));
REQUIRE(lazy_parallel.lazy_loading());
std::filesystem::remove(path); std::filesystem::remove(path);
size_t count = 0; size_t count = 0;
@@ -583,7 +589,7 @@ TEST_CASE("Parallel and paged parsing yield the same instances, attributes, inve
a.to_string(sa); a.to_string(sa);
b.to_string(sb); b.to_string(sb);
REQUIRE(sb.str() == sa.str()); REQUIRE(sb.str() == sa.str());
for (ifcopenshell::file* other : {&paged, &paged_parallel}) { for (ifcopenshell::file* other : {&paged, &paged_parallel, &lazy_parallel}) {
const express::base c = other->instance_by_id((int)a.id()); const express::base c = other->instance_by_id((int)a.id());
REQUIRE(c); REQUIRE(c);
std::ostringstream sc; std::ostringstream sc;
@@ -603,5 +609,6 @@ TEST_CASE("Parallel and paged parsing yield the same instances, attributes, inve
for (const auto& rooted : serial.instances_by_type("IfcRoot")) { for (const auto& rooted : serial.instances_by_type("IfcRoot")) {
const std::string guid = rooted.get_attribute_value(0); const std::string guid = rooted.get_attribute_value(0);
REQUIRE(parallel.instance_by_guid(guid).id() == serial.instance_by_guid(guid).id()); REQUIRE(parallel.instance_by_guid(guid).id() == serial.instance_by_guid(guid).id());
REQUIRE(lazy_parallel.instance_by_guid(guid).id() == serial.instance_by_guid(guid).id());
} }
} }