47 "wall clock time spent in critical sections"),
49 "CPU time spent compressing"),
51 "timeCpuCriticalSection",
"ns",
"CPU time spent in critical section")});
75 f.SetOnDiskId(fNFields);
78 for (
auto *
f : fields) {
80 for (
auto &descendant : *
f) {
81 connectField(descendant);
84 fBufferedColumns.resize(fNColumns);
89 return fInnerSink->GetDescriptor();
96 fInnerModel = model.Clone();
97 fInnerSink->Init(*fInnerModel);
108 auto cloned = field->Clone(field->GetFieldName());
111 auto parent = field->GetParent();
114 auto &innerParent = fInnerModel->GetMutableField(parent->GetQualifiedFieldName());
118 fInnerModel->AddField(std::move(cloned));
123 auto cloned = field->
Clone(field->GetFieldName());
127 fieldMap[
p] = &fInnerModel->GetConstField(projectedFields.GetSourceField(field)->GetQualifiedFieldName());
128 auto targetIt = cloned->begin();
129 for (
auto &
f : *field)
130 fieldMap[&(*targetIt++)] =
131 &fInnerModel->GetConstField(projectedFields.GetSourceField(&
f)->GetQualifiedFieldName());
136 fInnerModel->Unfreeze();
138 std::back_inserter(innerChangeset.fAddedFields), cloneAddField);
140 std::back_inserter(innerChangeset.fAddedProjectedFields), cloneAddProjectedField);
141 fInnerModel->Freeze();
142 fInnerSink->UpdateSchema(innerChangeset, firstEntry);
148 RNTuplePlainTimer timer(fCounters->fTimeWallCriticalSection, fCounters->fTimeCpuCriticalSection);
149 fInnerSink->UpdateExtraTypeInfo(extraTypeInfo);
154 fSuppressedColumns.emplace_back(columnHandle);
164 auto &zipItem = fBufferedColumns.at(colId).BufferPage(columnHandle);
165 std::size_t maxSealedPageBytes = page.
GetNBytes() + GetWriteOptions().GetEnablePageChecksums() * kNBytesPageChecksum;
167 auto &sealedPage = fBufferedColumns.at(colId).RegisterSealedPage();
169 auto allocateBuf = [&zipItem, maxSealedPageBytes]() {
170 zipItem.fBuf = MakeUninitArray<unsigned char>(maxSealedPageBytes);
173 auto shrinkSealedPage = [&zipItem, maxSealedPageBytes, &sealedPage]() {
176 auto sealedBufferSize = sealedPage.GetBufferSize();
177 if (sealedBufferSize < maxSealedPageBytes) {
178 auto buf = MakeUninitArray<unsigned char>(sealedBufferSize);
179 memcpy(buf.get(), sealedPage.GetBuffer(), sealedBufferSize);
180 zipItem.fBuf = std::move(buf);
181 sealedPage.SetBuffer(zipItem.fBuf.get());
188 std::size_t bufferedUncompressed = fBufferedUncompressed.load();
189 bool enoughWork = bufferedUncompressed > GetWriteOptions().GetApproxZippedClusterSize();
191 if (!fTaskScheduler || enoughWork) {
195 config.
fPage = &page;
198 config.
fWriteChecksum = GetWriteOptions().GetEnablePageChecksums();
200 config.
fBuffer = zipItem.fBuf.get();
202 RNTupleAtomicTimer timer(fCounters->fTimeWallZip, fCounters->fTimeCpuZip);
203 sealedPage = SealPage(config);
206 zipItem.fSealedPage = &sealedPage;
212 fBufferedUncompressed += page.
GetNBytes();
218 assert(zipItem.fPage.GetNBytes() == page.
GetNBytes());
221 fCounters->fParallelZip.SetValue(1);
224 fTaskScheduler->AddTask([
this, &zipItem, &sealedPage, &element, allocateBuf, shrinkSealedPage] {
227 fBufferedUncompressed -= zipItem.fPage.GetNBytes();
231 config.
fPage = &zipItem.fPage;
234 config.
fWriteChecksum = GetWriteOptions().GetEnablePageChecksums();
237 config.
fBuffer = zipItem.fBuf.get();
240 sealedPage = SealPage(config);
242 zipItem.fSealedPage = &sealedPage;
244 zipItem.fPage =
RPage();
255 std::span<ROOT::Internal::RPageStorage::RSealedPageGroup> )
266 assert(fBufferedUncompressed == 0 &&
"all buffered pages should have been processed");
268 std::vector<RSealedPageGroup> toCommit;
269 toCommit.reserve(fBufferedColumns.size());
270 for (
auto &bufColumn : fBufferedColumns) {
271 R__ASSERT(bufColumn.HasSealedPagesOnly());
272 const auto &sealedPages = bufColumn.GetSealedPages();
273 toCommit.emplace_back(bufColumn.GetHandle().fPhysicalId, sealedPages.cbegin(), sealedPages.cend());
278 RNTuplePlainTimer timer(fCounters->fTimeWallCriticalSection, fCounters->fTimeCpuCriticalSection);
279 fInnerSink->CommitSealedPageV(toCommit);
281 for (
auto handle : fSuppressedColumns)
282 fInnerSink->CommitSuppressedColumn(handle);
283 fSuppressedColumns.clear();
288 for (
auto &bufColumn : fBufferedColumns)
289 bufColumn.DropBufferedPages();
294 std::uint64_t nbytes;
295 FlushClusterImpl([&] { nbytes = fInnerSink->CommitCluster(nNewEntries); });
302 FlushClusterImpl([&] { stagedCluster = fInnerSink->StageCluster(nNewEntries); });
303 return stagedCluster;
309 RNTuplePlainTimer timer(fCounters->fTimeWallCriticalSection, fCounters->fTimeCpuCriticalSection);
310 fInnerSink->CommitStagedClusters(clusters);
316 RNTuplePlainTimer timer(fCounters->fTimeWallCriticalSection, fCounters->fTimeCpuCriticalSection);
317 fInnerSink->CommitClusterGroup();
323 RNTuplePlainTimer timer(fCounters->fTimeWallCriticalSection, fCounters->fTimeCpuCriticalSection);
324 return fInnerSink->CommitDataset();
329 return fInnerSink->ReservePage(columnHandle, nElements);
332std::unique_ptr<ROOT::Internal::RPageSink>
335 return fInnerSink->CloneAsHidden(
name, opts);
340 fInnerSink->CommitAttributeSet(attrSetName, attrAnchorInfo);
#define R__FAIL(msg)
Short-hand to return an RResult<T> in an error state; the RError is implicitly converted into RResult...
#define R__ASSERT(e)
Checks condition e and reports a fatal error if it's false.
winID h TVirtualViewer3D TVirtualGLPainter p
A thread-safe integral performance counter.
A collection of Counter objects with a name, a unit, and a description.
void ObserveMetrics(RNTupleMetrics &observee)
CounterPtrT MakeCounter(const std::string &name, Args &&... args)
A non thread-safe integral performance counter.
Record wall time and CPU time between construction and destruction.
A column is a storage-backed array of a simple, fixed-size type, from which pages can be mapped into ...
ROOT::Internal::RColumnElementBase * GetElement() const
std::deque< RPageZipItem > fBufferedPages
Using a deque guarantees that element iterators are never invalidated by appends to the end of the it...
RPageStorage::SealedPageSequence_t fSealedPages
Pages that have been already sealed by a concurrent task.
RPage ReservePage(ColumnHandle_t columnHandle, std::size_t nElements) final
Get a new, empty page for the given column that can be filled with up to nElements; nElements must be...
std::unique_ptr< RPageSink > CloneAsHidden(std::string_view name, const RNTupleWriteOptions &opts) const final
Creates a new sink with the same underlying storage as this but writing to a different RNTuple named ...
void CommitStagedClusters(std::span< RStagedCluster > clusters) final
Commit staged clusters, logically appending them to the ntuple descriptor.
std::unique_ptr< RCounters > fCounters
std::uint64_t CommitCluster(ROOT::NTupleSize_t nNewEntries) final
Finalize the current cluster and create a new one for the following data.
void UpdateSchema(const RNTupleModelChangeset &changeset, ROOT::NTupleSize_t firstEntry) final
Incorporate incremental changes to the model into the ntuple descriptor.
void FlushClusterImpl(const std::function< void(void)> &FlushClusterFn)
void CommitSealedPage(ROOT::DescriptorId_t physicalColumnId, const RSealedPage &sealedPage) final
Write a preprocessed page to storage. The column must have been added before.
RPageSinkBuf(std::unique_ptr< RPageSink > inner)
void CommitAttributeSet(std::string_view attrSetName, const RNTupleLink &attrAnchorInfo) final
Adds the given anchor information (name + locator) into the main RNTuple's descriptor as an attribute...
RNTupleLink CommitDatasetImpl() final
const ROOT::RNTupleDescriptor & GetDescriptor() const final
Return the RNTupleDescriptor being constructed.
void UpdateExtraTypeInfo(const ROOT::RExtraTypeInfoDescriptor &extraTypeInfo) final
Adds an extra type information record to schema.
void CommitSealedPageV(std::span< RPageStorage::RSealedPageGroup > ranges) final
Write a vector of preprocessed pages to storage. The corresponding columns must have been added befor...
RStagedCluster StageCluster(ROOT::NTupleSize_t nNewEntries) final
Stage the current cluster and create a new one for the following data.
void CommitPage(ColumnHandle_t columnHandle, const RPage &page) final
Write a page to the storage. The column must have been added before.
void CommitSuppressedColumn(ColumnHandle_t columnHandle) final
Commits a suppressed column for the current cluster.
void CommitClusterGroup() final
Write out the page locations (page list envelope) for all the committed clusters since the last call ...
void InitImpl(ROOT::RNTupleModel &model) final
std::unique_ptr< RPageSink > fInnerSink
The inner sink, responsible for actually performing I/O.
void ConnectFields(const std::vector< ROOT::RFieldBase * > &fields, ROOT::NTupleSize_t firstEntry)
ColumnHandle_t AddColumn(ROOT::DescriptorId_t fieldId, RColumn &column) final
Register a new column.
An RAII wrapper used to synchronize a page sink. See GetSinkGuard().
Abstract interface to write data into an ntuple.
const ROOT::RNTupleWriteOptions & GetWriteOptions() const
Returns the sink's write options.
ROOT::Experimental::Detail::RNTupleMetrics fMetrics
const std::string & GetNTupleName() const
Returns the NTuple name.
A page is a slice of a column that is mapped into memory.
std::uint32_t GetNElements() const
std::size_t GetNBytes() const
The space taken by column elements in the buffer.
std::uint32_t GetElementSize() const
std::unordered_map< const ROOT::RFieldBase *, const ROOT::RFieldBase * > FieldMap_t
The map keys are the projected target fields, the map values are the backing source fields Note that ...
RResult< void > Add(std::unique_ptr< ROOT::RFieldBase > field, const FieldMap_t &fieldMap)
Adds a new projected field.
Base class for all ROOT issued exceptions.
A field translates read and write calls from/to underlying columns to/from tree values.
The container field for an ntuple model, which itself has no physical representation.
The on-storage metadata of an RNTuple.
The RNTupleModel encapulates the schema of an RNTuple.
Common user-tunable settings for storing RNTuples.
The field for an untyped record.
virtual TObject * Clone(const char *newname="") const
Make a clone of an object using the Streamer facility.
ROOT::RFieldZero & GetFieldZeroOfModel(RNTupleModel &model)
void AddItemToRecord(RRecordField &record, std::unique_ptr< RFieldBase > newItem)
RProjectedFields & GetProjectedFieldsOfModel(RNTupleModel &model)
void CallConnectPageSinkOnField(RFieldBase &, ROOT::Internal::RPageSink &, ROOT::NTupleSize_t firstEntry=0)
std::uint64_t DescriptorId_t
Distriniguishes elements of the same type within a descriptor, e.g. different fields.
std::uint64_t NTupleSize_t
Integer type long enough to hold the maximum number of entries in a column.
The incremental changes to a RNTupleModel
std::vector< ROOT::RFieldBase * > fAddedProjectedFields
Points to the projected fields in fModel that were added as part of an updater transaction.
std::vector< ROOT::RFieldBase * > fAddedFields
Points to the fields in fModel that were added as part of an updater transaction.
I/O performance counters that get registered in fMetrics.
Parameters for the SealPage() method.
bool fWriteChecksum
Adds a 8 byte little-endian xxhash3 checksum to the page payload.
std::uint32_t fCompressionSettings
Compression algorithm and level to apply.
void * fBuffer
Location for sealed output. The memory buffer has to be large enough.
const ROOT::Internal::RPage * fPage
Input page to be sealed.
bool fAllowAlias
If false, the output buffer must not point to the input page buffer, which would otherwise be an opti...
const ROOT::Internal::RColumnElementBase * fElement
Corresponds to the page's elements, for size calculation etc.
Cluster that was staged, but not yet logically appended to the RNTuple.
ROOT::Internal::RColumn * fColumn
ROOT::DescriptorId_t fPhysicalId
A sealed page contains the bytes of a page as written to storage (packed & compressed).