Logo ROOT  
Reference Guide
 
Loading...
Searching...
No Matches
RPageStorage.cxx
Go to the documentation of this file.
1/// \file RPageStorage.cxx
2/// \author Jakob Blomer <jblomer@cern.ch>
3/// \date 2018-10-04
4
5/*************************************************************************
6 * Copyright (C) 1995-2019, Rene Brun and Fons Rademakers. *
7 * All rights reserved. *
8 * *
9 * For the licensing terms see $ROOTSYS/LICENSE. *
10 * For the list of contributors see $ROOTSYS/README/CREDITS. *
11 *************************************************************************/
12
13#include <ROOT/RPageStorage.hxx>
15#include <ROOT/RColumn.hxx>
16#include <ROOT/RFieldBase.hxx>
20#include <ROOT/RNTupleModel.hxx>
22#include <ROOT/RNTupleUtils.hxx>
23#include <ROOT/RNTupleZip.hxx>
25#include <ROOT/RPageSinkBuf.hxx>
26#include <ROOT/StringUtils.hxx>
27#ifdef R__ENABLE_DAOS
29#endif
30#ifdef R__ENABLE_S3
32#endif
33
34#include <Compression.h>
35#include <TError.h>
36
37#include <algorithm>
38#include <atomic>
39#include <cassert>
40#include <cstring>
41#include <functional>
42#include <memory>
43#include <string_view>
44#include <unordered_map>
45#include <utility>
46
53
55
59
65
67 : fMetrics(""), fPageAllocator(std::make_unique<ROOT::Internal::RPageAllocatorHeap>()), fNTupleName(name)
68{
69}
70
72
74{
75 if (!fHasChecksum)
76 return;
77
78 auto charBuf = reinterpret_cast<const unsigned char *>(fBuffer);
79 auto checksumBuf = const_cast<unsigned char *>(charBuf) + GetDataSize();
80 std::uint64_t xxhash3;
82}
83
85{
86 if (!fHasChecksum)
88
89 auto success = RNTupleSerializer::VerifyXxHash3(reinterpret_cast<const unsigned char *>(fBuffer), GetDataSize());
90 if (!success)
91 return R__FAIL("page checksum verification failed, data corruption detected");
93}
94
96{
97 if (!fHasChecksum)
98 return R__FAIL("invalid attempt to extract non-existing page checksum");
99
100 assert(fBufferSize >= kNBytesPageChecksum);
101 std::uint64_t checksum;
103 reinterpret_cast<const unsigned char *>(fBuffer) + fBufferSize - kNBytesPageChecksum, checksum);
104 return checksum;
105}
106
107//------------------------------------------------------------------------------
108
111{
112 auto [itr, _] = fColumnInfos.emplace(physicalColumnId, std::vector<RColumnInfo>());
113 for (auto &columnInfo : itr->second) {
114 if (columnInfo.fElementId == elementId) {
115 columnInfo.fRefCounter++;
116 return;
117 }
118 }
119 itr->second.emplace_back(RColumnInfo{elementId, 1});
120}
121
124{
125 auto itr = fColumnInfos.find(physicalColumnId);
126 R__ASSERT(itr != fColumnInfos.end());
127 for (std::size_t i = 0; i < itr->second.size(); ++i) {
128 if (itr->second[i].fElementId != elementId)
129 continue;
130
131 itr->second[i].fRefCounter--;
132 if (itr->second[i].fRefCounter == 0) {
133 itr->second.erase(itr->second.begin() + i);
134 if (itr->second.empty()) {
135 fColumnInfos.erase(itr);
136 }
137 }
138 break;
139 }
140}
141
149
151{
152 if (fFirstEntry == ROOT::kInvalidNTupleIndex) {
153 /// Entry range unset, we assume that the entry range covers the complete source
154 return true;
155 }
156
157 if (clusterDesc.GetNEntries() == 0)
158 return true;
159 if ((clusterDesc.GetFirstEntryIndex() + clusterDesc.GetNEntries()) <= fFirstEntry)
160 return false;
161 if (clusterDesc.GetFirstEntryIndex() >= (fFirstEntry + fNEntries))
162 return false;
163 return true;
164}
165
168 fClusterPool(*this, ROOT::Internal::RNTupleReadOptionsManip::GetClusterBunchSize(options)),
169 fPagePool(*this),
170 fOptions(options)
171{
172}
173
175
176std::unique_ptr<ROOT::Internal::RPageSource>
177ROOT::Internal::RPageSource::Create(std::string_view ntupleName, std::string_view location,
178 const ROOT::RNTupleReadOptions &options)
179{
180 if (ntupleName.empty()) {
181 throw RException(R__FAIL("empty RNTuple name"));
182 }
183 if (location.empty()) {
184 throw RException(R__FAIL("empty storage location"));
185 }
186 if (location.find("daos://") == 0)
187#ifdef R__ENABLE_DAOS
188 return std::make_unique<ROOT::Experimental::Internal::RPageSourceDaos>(ntupleName, location, options);
189#else
190 throw RException(R__FAIL("This RNTuple build does not support DAOS."));
191#endif
192
193 if (ROOT::StartsWith(location, "ntpl+s3+http://") || ROOT::StartsWith(location, "ntpl+s3+https://"))
194 throw RException(R__FAIL("S3 read support is not yet implemented."));
195
196 return std::make_unique<ROOT::Internal::RPageSourceFile>(ntupleName, location, options);
197}
198
201{
203 auto physicalId =
204 GetSharedDescriptorGuard()->FindPhysicalColumnId(fieldId, column.GetIndex(), column.GetRepresentationIndex());
206 fActivePhysicalColumns.Insert(physicalId, column.GetElement()->GetIdentifier());
207 return ColumnHandle_t{physicalId, &column};
208}
209
211{
212 fActivePhysicalColumns.Erase(columnHandle.fPhysicalId, columnHandle.fColumn->GetElement()->GetIdentifier());
213}
214
216{
217 if ((range.fFirstEntry + range.fNEntries) > GetNEntries()) {
218 throw RException(R__FAIL("invalid entry range"));
219 }
220 fEntryRange = range;
221}
222
224{
225 if (!fHasStructure)
226 LoadStructureImpl();
227 fHasStructure = true;
228}
229
231{
232 if (fIsAttached)
233 return;
234
235 LoadStructure();
236
237 auto descGuard = GetExclDescriptorGuard();
238 descGuard.MoveIn(AttachImpl());
239 fStructureBuffer.Reset();
240
241 std::vector<unsigned char> buffer;
242 for (const auto &cgDesc : descGuard->GetClusterGroupIterable()) {
243 buffer.resize(cgDesc.GetPageListLength() + cgDesc.GetPageListLocator().GetNBytesOnStorage());
244 auto zipBuffer = buffer.data() + cgDesc.GetPageListLength();
245
246 LoadPageListImpl(cgDesc.GetPageListLocator(), zipBuffer);
247 RNTupleDecompressor::Unzip(zipBuffer, cgDesc.GetPageListLocator().GetNBytesOnStorage(),
248 cgDesc.GetPageListLength(), buffer.data());
249 RNTupleSerializer::DeserializePageList(buffer.data(), cgDesc.GetPageListLength(), cgDesc.GetId(), *descGuard,
250 mode);
251 }
252
253 fIsAttached = true;
254}
255
256std::unique_ptr<ROOT::Internal::RPageSource> ROOT::Internal::RPageSource::Clone() const
257{
258 auto clone = CloneImpl();
259 if (fIsAttached) {
260 clone->GetExclDescriptorGuard().MoveIn(GetSharedDescriptorGuard()->Clone());
261 clone->fHasStructure = true;
262 clone->fIsAttached = true;
263 }
264 return clone;
265}
266
268{
269 return GetSharedDescriptorGuard()->GetNEntries();
270}
271
273{
274 auto descGuard = GetSharedDescriptorGuard();
275 if (descGuard->GetNClusters() == 0)
276 return 0;
277
278 auto itr = descGuard->GetClusterGroupIterable().begin();
279 itr += descGuard->GetNClusterGroups() - 1;
280 R__ASSERT(itr->HasClusterDetails());
281 const auto &cd = descGuard->GetClusterDescriptor(itr->GetClusterIds().back());
282 R__ASSERT(cd.ContainsColumn(physicalColumnId));
283 const auto &columnRange = cd.GetColumnRange(physicalColumnId);
284 return columnRange.GetFirstElementIndex() + columnRange.GetNElements();
285}
286
289{
291 {
292 auto descriptorGuard = GetSharedDescriptorGuard();
293 const auto &clusterDesc = descriptorGuard->GetClusterDescriptor(clusterId);
294 firstEntryInNextCluster = clusterDesc.GetFirstEntryIndex() + clusterDesc.GetNEntries();
295 }
296 return FindClusterId(firstEntryInNextCluster, nextId);
297}
298
301{
303 auto descGuard = GetSharedDescriptorGuard();
304 const auto &desc = descGuard.GetRef();
305
306 if (desc.GetNClusterGroups() == 0)
307 return descGuard;
308
309 // Binary search in the cluster group list, followed by a binary search in the clusters of that cluster group
310
311 auto cgIter = desc.GetClusterGroupIterable().begin();
312 std::size_t cgLeft = 0;
313 std::size_t cgRight = desc.GetNClusterGroups() - 1;
314 while (cgLeft <= cgRight) {
315 const std::size_t cgMidpoint = (cgLeft + cgRight) / 2;
316 const auto &cgDesc = *(cgIter + cgMidpoint);
317
318 if (cgDesc.GetMinEntry() > entryIdx) {
320 cgRight = cgMidpoint - 1;
321 continue;
322 }
323
324 if (cgDesc.GetMinEntry() + cgDesc.GetEntrySpan() <= entryIdx) {
325 cgLeft = cgMidpoint + 1;
326 continue;
327 }
328
329 // Binary search in the current cluster group; since we already checked the element range boundaries,
330 // the element must be in that cluster group.
331 const auto &clusterIds = cgDesc.GetClusterIds();
332 R__ASSERT(!clusterIds.empty());
333 std::size_t clusterLeft = 0;
334 std::size_t clusterRight = clusterIds.size() - 1;
335 while (clusterLeft <= clusterRight) {
336 const std::size_t clusterMidpoint = (clusterLeft + clusterRight) / 2;
337 const auto &clusterDesc = desc.GetClusterDescriptor(clusterIds[clusterMidpoint]);
338
339 if (clusterDesc.GetFirstEntryIndex() > entryIdx) {
342 continue;
343 }
344
345 if (clusterDesc.GetFirstEntryIndex() + clusterDesc.GetNEntries() <= entryIdx) {
347 continue;
348 }
349
351 return descGuard;
352 }
353 R__ASSERT(false);
354 }
355 return descGuard;
356}
357
360{
362 auto descGuard = GetSharedDescriptorGuard();
363 const auto &desc = descGuard.GetRef();
364
365 if (desc.GetNClusterGroups() == 0)
366 return descGuard;
367
368 // Binary search in the cluster group list, followed by a binary search in the clusters of that cluster group
369
370 auto cgIter = desc.GetClusterGroupIterable().begin();
371 std::size_t cgLeft = 0;
372 std::size_t cgRight = desc.GetNClusterGroups() - 1;
373 while (cgLeft <= cgRight) {
374 const std::size_t cgMidpoint = (cgLeft + cgRight) / 2;
375 const auto &clusterIds = (cgIter + cgMidpoint)->GetClusterIds();
376 R__ASSERT(!clusterIds.empty());
377
378 const auto &clusterDesc = desc.GetClusterDescriptor(clusterIds.front());
379 // this may happen if the RNTuple has an empty schema
380 if (!clusterDesc.ContainsColumn(physicalColumnId))
381 return descGuard;
382
383 const auto firstElementInGroup = clusterDesc.GetColumnRange(physicalColumnId).GetFirstElementIndex();
385 // Look into the lower half of cluster groups
387 cgRight = cgMidpoint - 1;
388 continue;
389 }
390
391 const auto &lastColumnRange = desc.GetClusterDescriptor(clusterIds.back()).GetColumnRange(physicalColumnId);
392 if ((lastColumnRange.GetFirstElementIndex() + lastColumnRange.GetNElements()) <= index) {
393 // Look into the upper half of cluster groups
394 cgLeft = cgMidpoint + 1;
395 continue;
396 }
397
398 // Binary search in the current cluster group; since we already checked the element range boundaries,
399 // the element must be in that cluster group.
400 std::size_t clusterLeft = 0;
401 std::size_t clusterRight = clusterIds.size() - 1;
402 while (clusterLeft <= clusterRight) {
403 const std::size_t clusterMidpoint = (clusterLeft + clusterRight) / 2;
405 const auto &columnRange = desc.GetClusterDescriptor(clusterId).GetColumnRange(physicalColumnId);
406
407 if (columnRange.Contains(index)) {
408 cid = clusterId;
409 return descGuard;
410 }
411
412 if (columnRange.GetFirstElementIndex() > index) {
415 continue;
416 }
417
418 if (columnRange.GetFirstElementIndex() + columnRange.GetNElements() <= index) {
420 continue;
421 }
422 }
423 R__ASSERT(false);
424 }
425 return descGuard;
426}
427
429{
430 if (fTaskScheduler)
431 UnzipClusterImpl(cluster);
432}
433
435{
436 RNTupleAtomicTimer timer(fCounters->fTimeWallUnzip, fCounters->fTimeCpuUnzip);
437
438 const auto clusterId = cluster->GetId();
439 auto descriptorGuard = GetSharedDescriptorGuard();
440 const auto &clusterDescriptor = descriptorGuard->GetClusterDescriptor(clusterId);
441
442 fPreloadedClusters[clusterDescriptor.GetFirstEntryIndex()] = clusterId;
443
444 std::atomic<bool> foundChecksumFailure{false};
445
446 std::vector<std::unique_ptr<RColumnElementBase>> allElements;
447 const auto &columnsInCluster = cluster->GetAvailPhysicalColumns();
448 for (const auto columnId : columnsInCluster) {
449 // By the time we unzip a cluster, the set of active columns may have already changed wrt. to the moment when
450 // we requested reading the cluster. That doesn't matter much, we simply decompress what is now in the list
451 // of active columns.
452 if (!fActivePhysicalColumns.HasColumnInfos(columnId))
453 continue;
454 const auto &columnInfos = fActivePhysicalColumns.GetColumnInfos(columnId);
455
456 allElements.reserve(allElements.size() + columnInfos.size());
457 for (const auto &info : columnInfos) {
458 allElements.emplace_back(GenerateColumnElement(info.fElementId));
459
460 const auto &pageRange = clusterDescriptor.GetPageRange(columnId);
461 std::uint64_t pageNo = 0;
462 std::uint64_t firstInPage = 0;
463 for (const auto &pi : pageRange.GetPageInfos()) {
464 auto onDiskPage = cluster->GetOnDiskPage(ROnDiskPage::Key{columnId, pageNo});
466 sealedPage.SetNElements(pi.GetNElements());
467 sealedPage.SetHasChecksum(pi.HasChecksum());
468 sealedPage.SetBufferSize(pi.GetLocator().GetNBytesOnStorage() + pi.HasChecksum() * kNBytesPageChecksum);
469 sealedPage.SetBuffer(onDiskPage->GetAddress());
470 R__ASSERT(onDiskPage && (onDiskPage->GetSize() == sealedPage.GetBufferSize()));
471
472 auto taskFunc = [this, columnId, clusterId, firstInPage, sealedPage, element = allElements.back().get(),
474 indexOffset = clusterDescriptor.GetColumnRange(columnId).GetFirstElementIndex()]() {
475 const ROOT::Internal::RPagePool::RKey keyPagePool{columnId, element->GetIdentifier().fInMemoryType};
476 auto rv = UnsealPage(sealedPage, *element);
477 if (!rv) {
479 return;
480 }
481 auto newPage = rv.Unwrap();
482 fCounters->fSzUnzip.Add(element->GetSize() * sealedPage.GetNElements());
483
484 newPage.SetWindow(indexOffset + firstInPage,
486 fPagePool.PreloadPage(std::move(newPage), keyPagePool);
487 };
489 fTaskScheduler->AddTask(taskFunc);
490
491 firstInPage += pi.GetNElements();
492 pageNo++;
493 } // for all pages in column
494
495 fCounters->fNPageUnsealed.Add(pageNo);
496 } // for all in-memory types of the column
497 } // for all columns in cluster
498
499 fTaskScheduler->Wait();
500
502 throw RException(R__FAIL("page checksum verification failed, data corruption detected"));
503 }
504}
505
511 auto descriptorGuard = GetSharedDescriptorGuard();
512 const auto &clusterDesc = descriptorGuard->GetClusterDescriptor(clusterKey.fClusterId);
513
514 for (auto physicalColumnId : clusterKey.fPhysicalColumnSet) {
515 if (clusterDesc.GetColumnRange(physicalColumnId).IsSuppressed())
516 continue;
517
518 const auto &pageRange = clusterDesc.GetPageRange(physicalColumnId);
520 for (const auto &pageInfo : pageRange.GetPageInfos()) {
521 if (pageInfo.GetLocator().GetType() == RNTupleLocator::kTypePageZero) {
524 pageInfo.GetLocator().GetNBytesOnStorage()));
525 } else {
527 }
528 ++pageNo;
529 }
530 }
531}
532
534{
535 if (fLastUsedCluster == clusterId)
536 return;
537
539 GetSharedDescriptorGuard()->GetClusterDescriptor(clusterId).GetFirstEntryIndex();
540 auto itr = fPreloadedClusters.begin();
541 while ((itr != fPreloadedClusters.end()) && (itr->first < firstEntryIndex)) {
542 if (fPinnedClusters.count(itr->second) > 0) {
543 ++itr;
544 } else {
545 fPagePool.Evict(itr->second);
546 itr = fPreloadedClusters.erase(itr);
547 }
548 }
549 std::size_t poolWindow = 0;
550 while ((itr != fPreloadedClusters.end()) &&
552 ++itr;
553 ++poolWindow;
554 }
555 while (itr != fPreloadedClusters.end()) {
556 if (fPinnedClusters.count(itr->second) > 0) {
557 ++itr;
558 } else {
559 fPagePool.Evict(itr->second);
560 itr = fPreloadedClusters.erase(itr);
561 }
562 }
563
564 fLastUsedCluster = clusterId;
565}
566
569{
570 const auto clusterId = localIndex.GetClusterId();
571
573 {
574 auto descriptorGuard = GetSharedDescriptorGuard();
575 const auto &clusterDescriptor = descriptorGuard->GetClusterDescriptor(clusterId);
576 pageInfo = clusterDescriptor.GetPageRange(physicalColumnId).Find(localIndex.GetIndexInCluster());
577 }
578
579 assert(pageInfo.GetLocator().GetType() != RNTupleLocator::kTypePageZero);
580
581 sealedPage.SetBufferSize(pageInfo.GetLocator().GetNBytesOnStorage() + pageInfo.HasChecksum() * kNBytesPageChecksum);
582 sealedPage.SetNElements(pageInfo.GetNElements());
583 sealedPage.SetHasChecksum(pageInfo.HasChecksum());
584
585 if (!sealedPage.GetBuffer())
586 return;
587
588 LoadSealedPageImpl(pageInfo.GetLocator(), sealedPage);
589 sealedPage.VerifyChecksumIfEnabled().ThrowOnError();
590}
591
594{
595 const auto &pageInfo = pageSummary.fPageInfo;
596 assert(pageInfo.GetLocator().GetType() == RNTupleLocator::kTypePageZero);
597
598 const auto element = columnHandle.fColumn->GetElement();
599 const auto elementSize = element->GetSize();
600 const auto elementInMemoryType = element->GetIdentifier().fInMemoryType;
601
602 auto pageZero = fPageAllocator->NewPage(elementSize, pageInfo.GetNElements());
603 pageZero.GrowUnchecked(pageInfo.GetNElements());
604 std::memset(pageZero.GetBuffer(), 0, pageZero.GetNBytes());
605 pageZero.SetWindow(pageSummary.fColumnOffset + pageInfo.GetFirstElementIndex(),
606 RPage::RClusterInfo(pageSummary.fClusterId, pageSummary.fColumnOffset));
607 return fPagePool.RegisterPage(std::move(pageZero), RPagePool::RKey{columnHandle.fPhysicalId, elementInMemoryType});
608}
609
612{
613 if (pageSummary.fPageInfo.GetLocator().GetType() == RNTupleLocator::kTypeUnknown) {
614 throw RException(R__FAIL("tried to read a page with an unknown locator"));
615 } else if (pageSummary.fPageInfo.GetLocator().GetType() == RNTupleLocator::kTypePageZero) {
616 return LoadZeroPage(columnHandle, pageSummary);
617 }
618
619 const auto &columnId = columnHandle.fPhysicalId;
620 const auto &clusterId = pageSummary.fClusterId;
621 const auto &pageInfo = pageSummary.fPageInfo;
622
623 const auto element = columnHandle.fColumn->GetElement();
624 const auto elementSize = element->GetSize();
625 const auto elementInMemoryType = element->GetIdentifier().fInMemoryType;
626
627 UpdateLastUsedCluster(clusterId);
628
630 sealedPage.SetNElements(pageInfo.GetNElements());
631 sealedPage.SetHasChecksum(pageInfo.HasChecksum());
632 sealedPage.SetBufferSize(pageInfo.GetLocator().GetNBytesOnStorage() + pageInfo.HasChecksum() * kNBytesPageChecksum);
633 std::unique_ptr<unsigned char[]> directReadBuffer; // only used if cluster pool is turned off
634
635 if (fOptions.GetClusterCache() == ROOT::RNTupleReadOptions::EClusterCache::kOff) {
637 sealedPage.SetBuffer(directReadBuffer.get());
638 LoadSealedPageImpl(pageInfo.GetLocator(), sealedPage);
639
640 fCounters->fNPageRead.Inc();
641 fCounters->fNRead.Inc();
642 fCounters->fSzReadPayload.Add(sealedPage.GetBufferSize());
643 } else {
644 if (!fCurrentCluster || (fCurrentCluster->GetId() != clusterId) || !fCurrentCluster->ContainsColumn(columnId))
645 fCurrentCluster = fClusterPool.GetCluster(clusterId, fActivePhysicalColumns.ToColumnSet());
646 R__ASSERT(fCurrentCluster->ContainsColumn(columnId));
647
648 // The cluster pool may have unzipped the required page into the page pool
650 RNTupleLocalIndex(clusterId, pageInfo.GetFirstElementIndex()));
651 if (!cachedPageRef.Get().IsNull())
652 return cachedPageRef;
653
654 ROnDiskPage::Key key(columnId, pageInfo.GetPageNumber());
655 auto onDiskPage = fCurrentCluster->GetOnDiskPage(key);
656 R__ASSERT(onDiskPage && (sealedPage.GetBufferSize() == onDiskPage->GetSize()));
657 sealedPage.SetBuffer(onDiskPage->GetAddress());
658 }
659
661 {
662 RNTupleAtomicTimer timer(fCounters->fTimeWallUnzip, fCounters->fTimeCpuUnzip);
663 newPage = UnsealPage(sealedPage, *element).Unwrap();
664 fCounters->fSzUnzip.Add(elementSize * pageInfo.GetNElements());
665 }
666
667 newPage.SetWindow(pageSummary.fColumnOffset + pageInfo.GetFirstElementIndex(),
669 fCounters->fNPageUnsealed.Inc();
670
671 return fPagePool.RegisterPage(std::move(newPage), RPagePool::RKey{columnId, elementInMemoryType});
672}
673
676{
677 const auto columnId = columnHandle.fPhysicalId;
678 const auto columnElementId = columnHandle.fColumn->GetElement()->GetIdentifier();
679 auto cachedPageRef =
680 fPagePool.GetPage(ROOT::Internal::RPagePool::RKey{columnId, columnElementId.fInMemoryType}, globalIndex);
681 if (!cachedPageRef.Get().IsNull()) {
682 UpdateLastUsedCluster(cachedPageRef.Get().GetClusterInfo().GetId());
683 return cachedPageRef;
684 }
685
687 {
688 auto descriptorGuard = FindClusterId(columnId, globalIndex, pageSummary.fClusterId);
689
690 if (pageSummary.fClusterId == ROOT::kInvalidDescriptorId)
691 throw RException(R__FAIL("entry with index " + std::to_string(globalIndex) + " out of bounds"));
692
693 const auto &clusterDescriptor = descriptorGuard->GetClusterDescriptor(pageSummary.fClusterId);
694 const auto &columnRange = clusterDescriptor.GetColumnRange(columnId);
695 if (columnRange.IsSuppressed())
697
698 pageSummary.fColumnOffset = columnRange.GetFirstElementIndex();
699 R__ASSERT(pageSummary.fColumnOffset <= globalIndex);
700 pageSummary.fPageInfo = clusterDescriptor.GetPageRange(columnId).Find(globalIndex - pageSummary.fColumnOffset);
701 }
702
703 return LoadPageFromSummary(columnHandle, pageSummary);
704}
705
708{
709 const auto clusterId = localIndex.GetClusterId();
710 const auto columnId = columnHandle.fPhysicalId;
711 const auto columnElementId = columnHandle.fColumn->GetElement()->GetIdentifier();
712 auto cachedPageRef =
713 fPagePool.GetPage(ROOT::Internal::RPagePool::RKey{columnId, columnElementId.fInMemoryType}, localIndex);
714 if (!cachedPageRef.Get().IsNull()) {
715 UpdateLastUsedCluster(clusterId);
716 return cachedPageRef;
717 }
718
720 throw RException(R__FAIL("entry out of bounds"));
721
723 {
724 auto descriptorGuard = GetSharedDescriptorGuard();
725 const auto &clusterDescriptor = descriptorGuard->GetClusterDescriptor(clusterId);
726 const auto &columnRange = clusterDescriptor.GetColumnRange(columnId);
727 if (columnRange.IsSuppressed())
729
730 pageSummary.fClusterId = clusterId;
731 pageSummary.fColumnOffset = columnRange.GetFirstElementIndex();
732 pageSummary.fPageInfo = clusterDescriptor.GetPageRange(columnId).Find(localIndex.GetIndexInCluster());
733 }
734
735 return LoadPageFromSummary(columnHandle, pageSummary);
736}
737
739{
740 fMetrics = RNTupleMetrics(prefix);
741 fMetrics.ObserveMetrics(fClusterPool.GetMetrics());
742 fMetrics.ObserveMetrics(fPagePool.GetMetrics());
743 fCounters = std::make_unique<RCounters>(RCounters{
744 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("nReadV", "", "number of vector read requests"),
745 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("nRead", "", "number of byte ranges read"),
746 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("szReadPayload", "B", "volume read from storage (required)"),
747 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("szReadOverhead", "B", "volume read from storage (overhead)"),
748 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("szUnzip", "B", "volume after unzipping"),
749 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("nClusterLoaded", "",
750 "number of partial clusters preloaded from storage"),
751 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("nPageRead", "", "number of pages read from storage"),
752 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("nPageUnsealed", "", "number of pages unzipped and decoded"),
753 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("timeWallRead", "ns", "wall clock time spent reading"),
754 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("timeWallUnzip", "ns", "wall clock time spent decompressing"),
755 *fMetrics.MakeCounter<RNTupleTickCounter<RNTupleAtomicCounter> *>("timeCpuRead", "ns", "CPU time spent reading"),
756 *fMetrics.MakeCounter<RNTupleTickCounter<RNTupleAtomicCounter> *>("timeCpuUnzip", "ns",
757 "CPU time spent decompressing"),
758 *fMetrics.MakeCounter<RNTupleCalcPerf *>(
759 "bwRead", "MB/s", "bandwidth compressed bytes read per second", fMetrics,
760 [](const RNTupleMetrics &metrics) -> std::pair<bool, double> {
761 if (const auto szReadPayload = metrics.GetLocalCounter("szReadPayload")) {
762 if (const auto szReadOverhead = metrics.GetLocalCounter("szReadOverhead")) {
763 if (const auto timeWallRead = metrics.GetLocalCounter("timeWallRead")) {
764 if (auto walltime = timeWallRead->GetValueAsInt()) {
765 double payload = szReadPayload->GetValueAsInt();
766 double overhead = szReadOverhead->GetValueAsInt();
767 // unit: bytes / nanosecond = GB/s
768 return {true, (1000. * (payload + overhead) / walltime)};
769 }
770 }
771 }
772 }
773 return {false, -1.};
774 }),
775 *fMetrics.MakeCounter<RNTupleCalcPerf *>(
776 "bwReadUnzip", "MB/s", "bandwidth uncompressed bytes read per second", fMetrics,
777 [](const RNTupleMetrics &metrics) -> std::pair<bool, double> {
778 if (const auto szUnzip = metrics.GetLocalCounter("szUnzip")) {
779 if (const auto timeWallRead = metrics.GetLocalCounter("timeWallRead")) {
780 if (auto walltime = timeWallRead->GetValueAsInt()) {
781 double unzip = szUnzip->GetValueAsInt();
782 // unit: bytes / nanosecond = GB/s
783 return {true, 1000. * unzip / walltime};
784 }
785 }
786 }
787 return {false, -1.};
788 }),
789 *fMetrics.MakeCounter<RNTupleCalcPerf *>(
790 "bwUnzip", "MB/s", "decompression bandwidth of uncompressed bytes per second", fMetrics,
791 [](const RNTupleMetrics &metrics) -> std::pair<bool, double> {
792 if (const auto szUnzip = metrics.GetLocalCounter("szUnzip")) {
793 if (const auto timeWallUnzip = metrics.GetLocalCounter("timeWallUnzip")) {
794 if (auto walltime = timeWallUnzip->GetValueAsInt()) {
795 double unzip = szUnzip->GetValueAsInt();
796 // unit: bytes / nanosecond = GB/s
797 return {true, 1000. * unzip / walltime};
798 }
799 }
800 }
801 return {false, -1.};
802 }),
803 *fMetrics.MakeCounter<RNTupleCalcPerf *>(
804 "rtReadEfficiency", "", "ratio of payload over all bytes read", fMetrics,
805 [](const RNTupleMetrics &metrics) -> std::pair<bool, double> {
806 if (const auto szReadPayload = metrics.GetLocalCounter("szReadPayload")) {
807 if (const auto szReadOverhead = metrics.GetLocalCounter("szReadOverhead")) {
808 if (auto payload = szReadPayload->GetValueAsInt()) {
809 // r/(r+o) = 1/((r+o)/r) = 1/(1 + o/r)
810 return {true, 1. / (1. + (1. * szReadOverhead->GetValueAsInt()) / payload)};
811 }
812 }
813 }
814 return {false, -1.};
815 }),
816 *fMetrics.MakeCounter<RNTupleCalcPerf *>("rtCompression", "", "ratio of compressed bytes / uncompressed bytes",
817 fMetrics, [](const RNTupleMetrics &metrics) -> std::pair<bool, double> {
818 if (const auto szReadPayload =
819 metrics.GetLocalCounter("szReadPayload")) {
820 if (const auto szUnzip = metrics.GetLocalCounter("szUnzip")) {
821 if (auto unzip = szUnzip->GetValueAsInt()) {
822 return {true, (1. * szReadPayload->GetValueAsInt()) / unzip};
823 }
824 }
825 }
826 return {false, -1.};
827 })});
828}
829
832{
833 return UnsealPage(sealedPage, element, *fPageAllocator);
834}
835
839{
840 // Unsealing a page zero is a no-op. `RPageRange::ExtendToFitColumnRange()` guarantees that the page zero buffer is
841 // large enough to hold `sealedPage.fNElements`
843 auto page = pageAlloc.NewPage(element.GetSize(), sealedPage.GetNElements());
844 page.GrowUnchecked(sealedPage.GetNElements());
845 memset(page.GetBuffer(), 0, page.GetNBytes());
846 return page;
847 }
848
849 auto rv = sealedPage.VerifyChecksumIfEnabled();
850 if (!rv)
851 return R__FORWARD_ERROR(rv);
852
853 const auto bytesPacked = element.GetPackedSize(sealedPage.GetNElements());
854 auto page = pageAlloc.NewPage(element.GetPackedSize(), sealedPage.GetNElements());
855 if (sealedPage.GetDataSize() != bytesPacked) {
857 page.GetBuffer());
858 } else {
859 // We cannot simply map the sealed page as we don't know its life time. Specialized page sources
860 // may decide to implement to not use UnsealPage but to custom mapping / decompression code.
861 // Note that usually pages are compressed.
862 memcpy(page.GetBuffer(), sealedPage.GetBuffer(), bytesPacked);
863 }
864
865 if (!element.IsMappable()) {
866 auto tmp = pageAlloc.NewPage(element.GetSize(), sealedPage.GetNElements());
867 element.Unpack(tmp.GetBuffer(), page.GetBuffer(), sealedPage.GetNElements());
868 page = std::move(tmp);
869 }
870
871 page.GrowUnchecked(sealedPage.GetNElements());
872 return page;
873}
874
876{
877 if (fHasStreamerInfosRegistered)
878 return;
879
880 for (const auto &extraTypeInfo : fDescriptor.GetExtraTypeInfoIterable()) {
882 continue;
883 // We don't need the result, it's enough that during deserialization, BuildCheck() is called for every
884 // streamer info record.
886 }
887
888 fHasStreamerInfosRegistered = true;
889}
890
891//------------------------------------------------------------------------------
892
894{
895 // Make the sort order unique by adding the physical on-disk column id as a secondary key
896 if (fCurrentPageSize == other.fCurrentPageSize)
897 return fColumn->GetOnDiskId() > other.fColumn->GetOnDiskId();
898 return fCurrentPageSize > other.fCurrentPageSize;
899}
900
902{
903 if (fMaxAllocatedBytes - fCurrentAllocatedBytes >= targetAvailableSize)
904 return true;
905
906 auto itr = fColumnsSortedByPageSize.begin();
907 while (itr != fColumnsSortedByPageSize.end()) {
908 if (itr->fCurrentPageSize <= pageSizeLimit)
909 break;
910 if (itr->fCurrentPageSize == itr->fInitialPageSize) {
911 ++itr;
912 continue;
913 }
914
915 // Flushing the current column will invalidate itr
916 auto itrFlush = itr++;
917
918 RColumnInfo next;
919 if (itr != fColumnsSortedByPageSize.end())
920 next = *itr;
921
922 itrFlush->fColumn->Flush();
923 if (fMaxAllocatedBytes - fCurrentAllocatedBytes >= targetAvailableSize)
924 return true;
925
926 if (next.fColumn == nullptr)
927 return false;
928 itr = fColumnsSortedByPageSize.find(next);
929 };
930
931 return false;
932}
933
935{
936 const RColumnInfo key{&column, column.GetWritePageCapacity(), 0};
937 auto itr = fColumnsSortedByPageSize.find(key);
938 if (itr == fColumnsSortedByPageSize.end()) {
939 if (!TryEvict(newWritePageSize, 0))
940 return false;
941 fColumnsSortedByPageSize.insert({&column, newWritePageSize, newWritePageSize});
942 fCurrentAllocatedBytes += newWritePageSize;
943 return true;
944 }
945
947 assert(newWritePageSize >= elem.fInitialPageSize);
948
949 if (newWritePageSize == elem.fCurrentPageSize)
950 return true;
951
952 fColumnsSortedByPageSize.erase(itr);
953
954 if (newWritePageSize < elem.fCurrentPageSize) {
955 // Page got smaller
956 fCurrentAllocatedBytes -= elem.fCurrentPageSize - newWritePageSize;
957 elem.fCurrentPageSize = newWritePageSize;
958 fColumnsSortedByPageSize.insert(elem);
959 return true;
960 }
961
962 // Page got larger, we may need to make space available
963 const auto diffBytes = newWritePageSize - elem.fCurrentPageSize;
964 if (!TryEvict(diffBytes, elem.fCurrentPageSize)) {
965 // Don't change anything, let the calling column flush itself
966 // TODO(jblomer): we may consider skipping the column in TryEvict and thus avoiding erase+insert
967 fColumnsSortedByPageSize.insert(elem);
968 return false;
969 }
970 fCurrentAllocatedBytes += diffBytes;
971 elem.fCurrentPageSize = newWritePageSize;
972 fColumnsSortedByPageSize.insert(elem);
973 return true;
974}
975
976//------------------------------------------------------------------------------
977
979 : RPageStorage(name), fOptions(options.Clone()), fWritePageMemoryManager(options.GetPageBufferBudget())
980{
982}
983
985
987{
988 assert(config.fPage);
989 assert(config.fElement);
990 assert(config.fBuffer);
991
992 unsigned char *pageBuf = reinterpret_cast<unsigned char *>(config.fPage->GetBuffer());
993 bool isAdoptedBuffer = true;
994 auto nBytesPacked = config.fPage->GetNBytes();
995 auto nBytesChecksum = config.fWriteChecksum * kNBytesPageChecksum;
996
997 if (!config.fElement->IsMappable()) {
998 nBytesPacked = config.fElement->GetPackedSize(config.fPage->GetNElements());
999 pageBuf = new unsigned char[nBytesPacked];
1000 isAdoptedBuffer = false;
1001 config.fElement->Pack(pageBuf, config.fPage->GetBuffer(), config.fPage->GetNElements());
1002 }
1004
1005 if ((config.fCompressionSettings != 0) || !config.fElement->IsMappable() || !config.fAllowAlias ||
1006 config.fWriteChecksum) {
1007 nBytesZipped =
1009 if (!isAdoptedBuffer)
1010 delete[] pageBuf;
1011 pageBuf = reinterpret_cast<unsigned char *>(config.fBuffer);
1012 isAdoptedBuffer = true;
1013 }
1014
1016
1017 RSealedPage sealedPage{pageBuf, nBytesZipped + nBytesChecksum, config.fPage->GetNElements(), config.fWriteChecksum};
1018 sealedPage.ChecksumIfEnabled();
1019
1020 return sealedPage;
1021}
1022
1025{
1026 const auto nBytes = page.GetNBytes() + GetWriteOptions().GetEnablePageChecksums() * kNBytesPageChecksum;
1027 if (fSealPageBuffer.size() < nBytes)
1028 fSealPageBuffer.resize(nBytes);
1029
1030 RSealPageConfig config;
1031 config.fPage = &page;
1032 config.fElement = &element;
1033 config.fCompressionSettings = GetWriteOptions().GetCompression();
1034 config.fWriteChecksum = GetWriteOptions().GetEnablePageChecksums();
1035 config.fAllowAlias = true;
1036 config.fBuffer = fSealPageBuffer.data();
1037
1038 return SealPage(config);
1039}
1040
1042{
1043 for (const auto &cb : fOnDatasetCommitCallbacks)
1044 cb(*this);
1045 return CommitDatasetImpl();
1046}
1047
1049{
1050 R__ASSERT(nElements > 0);
1051 const auto elementSize = columnHandle.fColumn->GetElement()->GetSize();
1052 const auto nBytes = elementSize * nElements;
1053 if (!fWritePageMemoryManager.TryUpdate(*columnHandle.fColumn, nBytes))
1054 return ROOT::Internal::RPage();
1055 return fPageAllocator->NewPage(elementSize, nElements);
1056}
1057
1058//------------------------------------------------------------------------------
1059
1060std::unique_ptr<ROOT::Internal::RPageSink>
1061ROOT::Internal::RPagePersistentSink::Create(std::string_view ntupleName, std::string_view location,
1062 const ROOT::RNTupleWriteOptions &options)
1063{
1064 if (ntupleName.empty()) {
1065 throw RException(R__FAIL("empty RNTuple name"));
1066 }
1067 if (location.empty()) {
1068 throw RException(R__FAIL("empty storage location"));
1069 }
1070 if (location.find("daos://") == 0) {
1071#ifdef R__ENABLE_DAOS
1072 return std::make_unique<ROOT::Experimental::Internal::RPageSinkDaos>(ntupleName, location, options);
1073#else
1074 throw RException(R__FAIL("This RNTuple build does not support DAOS."));
1075#endif
1076 }
1077
1078 if (ROOT::StartsWith(location, "ntpl+s3+http://") || ROOT::StartsWith(location, "ntpl+s3+https://")) {
1079#ifdef R__ENABLE_S3
1080 return std::make_unique<ROOT::Experimental::Internal::RPageSinkS3>(ntupleName, location, options);
1081#else
1082 throw RException(R__FAIL("This RNTuple build does not support S3. Rebuild ROOT with the 'curl' "
1083 "cmake option enabled (-Dcurl=ON) to enable the S3 backend."));
1084#endif
1085 }
1086
1087 // Otherwise assume that the user wants us to create a file.
1088 return std::make_unique<ROOT::Internal::RPageSinkFile>(ntupleName, location, options);
1089}
1090
1092 const ROOT::RNTupleWriteOptions &options)
1093 : RPageSink(name, options)
1094{
1095}
1096
1098
1101{
1102 auto columnId = fDescriptorBuilder.GetDescriptor().GetNPhysicalColumns();
1104 columnBuilder.LogicalColumnId(columnId)
1105 .PhysicalColumnId(columnId)
1106 .FieldId(fieldId)
1107 .BitsOnStorage(column.GetBitsOnStorage())
1108 .ValueRange(column.GetValueRange())
1109 .Type(column.GetType())
1110 .Index(column.GetIndex())
1111 .RepresentationIndex(column.GetRepresentationIndex())
1112 .FirstElementIndex(column.GetFirstElementIndex());
1113 // For late model extension, we assume that the primary column representation is the active one for the
1114 // deferred range. All other representations are suppressed.
1115 if (column.GetFirstElementIndex() > 0 && column.GetRepresentationIndex() > 0)
1116 columnBuilder.SetSuppressedDeferred();
1117 fDescriptorBuilder.AddColumn(columnBuilder.MoveDescriptor().Unwrap());
1118 return ColumnHandle_t{columnId, &column};
1119}
1120
1123{
1124 if (fIsInitialized) {
1125 for (const auto &field : changeset.fAddedFields) {
1126 if (field->GetStructure() == ENTupleStructure::kStreamer) {
1127 throw ROOT::RException(R__FAIL("a Model cannot be extended with Streamer fields"));
1128 }
1129 }
1130 }
1131
1132 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1133
1134 if (descriptor.GetNLogicalColumns() > descriptor.GetNPhysicalColumns()) {
1135 // If we already have alias columns, add an offset to the alias columns so that the new physical columns
1136 // of the changeset follow immediately the already existing physical columns
1137 auto getNColumns = [](const ROOT::RFieldBase &f) -> std::size_t {
1138 const auto &reps = f.GetColumnRepresentatives();
1139 if (reps.empty())
1140 return 0;
1141 return reps.size() * reps[0].size();
1142 };
1143 std::uint32_t nNewPhysicalColumns = 0;
1144 for (auto f : changeset.fAddedFields) {
1146 for (const auto &descendant : *f)
1148 }
1149 fDescriptorBuilder.ShiftAliasColumns(nNewPhysicalColumns);
1150 }
1151
1152 auto addField = [&](ROOT::RFieldBase &f) {
1153 auto fieldId = descriptor.GetNFields();
1154 fDescriptorBuilder.AddField(f, fieldId);
1155 fDescriptorBuilder.AddFieldLink(f.GetParent()->GetOnDiskId(), fieldId);
1156 f.SetOnDiskId(fieldId);
1157 ROOT::Internal::CallConnectPageSinkOnField(f, *this, firstEntry); // issues in turn calls to `AddColumn()`
1158 };
1159 auto addProjectedField = [&](ROOT::RFieldBase &f) {
1160 auto fieldId = descriptor.GetNFields();
1161 auto sourceFieldId =
1163 fDescriptorBuilder.AddField(f, fieldId);
1164 fDescriptorBuilder.AddFieldLink(f.GetParent()->GetOnDiskId(), fieldId);
1165 fDescriptorBuilder.AddFieldProjection(sourceFieldId, fieldId);
1166 f.SetOnDiskId(fieldId);
1167 for (const auto &source : descriptor.GetColumnIterable(sourceFieldId)) {
1168 auto targetId = descriptor.GetNLogicalColumns();
1170 columnBuilder.LogicalColumnId(targetId)
1171 .PhysicalColumnId(source.GetLogicalId())
1172 .FieldId(fieldId)
1173 .BitsOnStorage(source.GetBitsOnStorage())
1174 .ValueRange(source.GetValueRange())
1175 .Type(source.GetType())
1176 .Index(source.GetIndex())
1177 .RepresentationIndex(source.GetRepresentationIndex());
1178 fDescriptorBuilder.AddColumn(columnBuilder.MoveDescriptor().Unwrap());
1179 }
1180 };
1181
1182 R__ASSERT(firstEntry >= fPrevClusterNEntries);
1183 const auto nColumnsBeforeUpdate = descriptor.GetNPhysicalColumns();
1184 for (auto f : changeset.fAddedFields) {
1185 addField(*f);
1186 for (auto &descendant : *f)
1188 }
1189 for (auto f : changeset.fAddedProjectedFields) {
1191 for (auto &descendant : *f)
1193 }
1194
1195 const auto nColumns = descriptor.GetNPhysicalColumns();
1196 fOpenColumnRanges.reserve(fOpenColumnRanges.size() + (nColumns - nColumnsBeforeUpdate));
1197 fOpenPageRanges.reserve(fOpenPageRanges.size() + (nColumns - nColumnsBeforeUpdate));
1200 columnRange.SetPhysicalColumnId(i);
1201 // We set the first element index in the current cluster to the first element that is part of a materialized page
1202 // (i.e., that is part of a page list). For columns created during late model extension, however, the column range
1203 // is fixed up as needed by `RClusterDescriptorBuilder::AddExtendedColumnRanges()` on read back.
1204 columnRange.SetFirstElementIndex(descriptor.GetColumnDescriptor(i).GetFirstElementIndex());
1205 columnRange.SetNElements(0);
1206 columnRange.SetCompressionSettings(GetWriteOptions().GetCompression());
1207 fOpenColumnRanges.emplace_back(columnRange);
1209 pageRange.SetPhysicalColumnId(i);
1210 fOpenPageRanges.emplace_back(std::move(pageRange));
1211 }
1212
1213 // Mapping of memory to on-disk column IDs usually happens during serialization of the ntuple header. If the
1214 // header was already serialized, this has to be done manually as it is required for page list serialization.
1215 if (fSerializationContext.GetHeaderSize() > 0)
1216 fSerializationContext.MapSchema(descriptor, /*forHeaderExtension=*/true);
1217}
1218
1220{
1221 if (extraTypeInfo.GetContentId() != EExtraTypeInfoIds::kStreamerInfo)
1222 throw RException(R__FAIL("ROOT bug: unexpected type extra info in UpdateExtraTypeInfo()"));
1223
1224 fInfosOfStreamerFields.merge(RNTupleSerializer::DeserializeStreamerInfos(extraTypeInfo.GetContent()).Unwrap());
1225}
1226
1228{
1229 fDescriptorBuilder.SetNTuple(fNTupleName, model.GetDescription());
1230 fDescriptorBuilder.SetVersionForWriting();
1231 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1232
1234 fDescriptorBuilder.AddField(fieldZero, 0);
1235 fieldZero.SetOnDiskId(0);
1237 projectedFields.GetFieldZero().SetOnDiskId(0);
1238
1240 initialChangeset.fAddedFields.reserve(fieldZero.GetMutableSubfields().size());
1241 for (auto f : fieldZero.GetMutableSubfields())
1242 initialChangeset.fAddedFields.emplace_back(f);
1243 initialChangeset.fAddedProjectedFields.reserve(projectedFields.GetFieldZero().GetMutableSubfields().size());
1244 for (auto f : projectedFields.GetFieldZero().GetMutableSubfields())
1245 initialChangeset.fAddedProjectedFields.emplace_back(f);
1246 UpdateSchema(initialChangeset, 0U);
1247
1248 fSerializationContext = RNTupleSerializer::SerializeHeader(nullptr, descriptor).Unwrap();
1249 auto buffer = MakeUninitArray<unsigned char>(fSerializationContext.GetHeaderSize());
1250 fSerializationContext = RNTupleSerializer::SerializeHeader(buffer.get(), descriptor).Unwrap();
1251 InitImpl(buffer.get(), fSerializationContext.GetHeaderSize());
1252
1253 fDescriptorBuilder.BeginHeaderExtension();
1254}
1255
1256std::unique_ptr<ROOT::RNTupleModel>
1258{
1259 // Create new descriptor
1260 fDescriptorBuilder.SetSchemaFromExisting(srcDescriptor);
1261 // This is needed to be able to use GetTypeNameForComparison()
1262 fDescriptorBuilder.SetVersionForWriting();
1263 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1264
1265 // Create column/page ranges
1266 const auto nColumns = descriptor.GetNPhysicalColumns();
1267 R__ASSERT(fOpenColumnRanges.empty() && fOpenPageRanges.empty());
1268 fOpenColumnRanges.reserve(nColumns);
1269 fOpenPageRanges.reserve(nColumns);
1270 for (ROOT::DescriptorId_t i = 0; i < nColumns; ++i) {
1271 const auto &column = descriptor.GetColumnDescriptor(i);
1273 columnRange.SetPhysicalColumnId(i);
1274 columnRange.SetFirstElementIndex(column.GetFirstElementIndex());
1275 columnRange.SetNElements(0);
1276 columnRange.SetCompressionSettings(GetWriteOptions().GetCompression());
1277 fOpenColumnRanges.emplace_back(columnRange);
1279 pageRange.SetPhysicalColumnId(i);
1280 fOpenPageRanges.emplace_back(std::move(pageRange));
1281 }
1282
1283 if (copyClusters) {
1284 // Clone and add all cluster descriptors
1285 R__ASSERT(srcDescriptor.GetNClusters() == srcDescriptor.GetNActiveClusters());
1286 for (const auto &cluster : srcDescriptor.GetActiveClusterIterable()) {
1287 auto nEntries = cluster.GetNEntries();
1288 for (unsigned int i = 0; i < fOpenColumnRanges.size(); ++i) {
1289 R__ASSERT(fOpenColumnRanges[i].GetPhysicalColumnId() == i);
1290 if (!cluster.ContainsColumn(i)) // a cluster may not contain a column if that column is deferred
1291 break;
1292 const auto &columnRange = cluster.GetColumnRange(i);
1293 R__ASSERT(columnRange.GetPhysicalColumnId() == i);
1294 // TODO: properly handle suppressed columns (check MarkSuppressedColumnRange())
1295 fOpenColumnRanges[i].IncrementFirstElementIndex(columnRange.GetNElements());
1296 }
1297 fDescriptorBuilder.AddCluster(cluster.Clone());
1298 fPrevClusterNEntries += nEntries;
1299 }
1300 }
1301
1302 // Create model
1304 modelOpts.SetReconstructProjections(true);
1305 // We want to emulate unknown types to allow merging RNTuples containing types that we lack dictionaries for.
1306 modelOpts.SetEmulateUnknownTypes(true);
1307 auto model = descriptor.CreateModel(modelOpts);
1308 if (!copyClusters) {
1310 projectedFields.GetFieldZero().SetOnDiskId(model->GetConstFieldZero().GetOnDiskId());
1311 }
1312
1313 // Serialize header and init from it
1314 fSerializationContext = RNTupleSerializer::SerializeHeader(nullptr, descriptor).Unwrap();
1315 auto buffer = MakeUninitArray<unsigned char>(fSerializationContext.GetHeaderSize());
1316 fSerializationContext = RNTupleSerializer::SerializeHeader(buffer.get(), descriptor).Unwrap();
1317 InitImpl(buffer.get(), fSerializationContext.GetHeaderSize());
1318
1319 fDescriptorBuilder.BeginHeaderExtension();
1320
1321 // mark this sink as initialized
1322 fIsInitialized = true;
1323
1324 return model;
1325}
1326
1329 std::span<const RColumnFormat> newRepresentation,
1330 std::uint64_t clusterOffset)
1331{
1332 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1333
1334 assert(&descriptor.GetFieldDescriptor(field.GetId()) == &field);
1335 assert(!field.IsProjectedField());
1336 assert(field.GetColumnCardinality() > 0);
1337 assert(!field.GetLogicalColumnIds().empty());
1338 assert(newRepresentation.size() == field.GetColumnCardinality());
1339
1340 const std::size_t firstPhysicalIndex = fDescriptorBuilder.GetDescriptor().GetNPhysicalColumns();
1341 const std::uint16_t reprIndex = field.GetLogicalColumnIds().size() / field.GetColumnCardinality();
1342
1343 fDescriptorBuilder.ShiftAliasColumns(newRepresentation.size());
1344
1345 std::uint16_t columnIndex = 0; // index into the representation
1346 for (auto columnRepr : newRepresentation) {
1347 std::size_t bitsOnStorage = columnRepr.fBitWidth;
1348 if (!bitsOnStorage) {
1350 if (rangeMin != rangeMax) {
1351 throw ROOT::RException(R__FAIL("bit width must be given for columns of variable bit width"));
1352 }
1354 }
1355
1356 const ROOT::DescriptorId_t firstReprColumnId = field.GetLogicalColumnIds()[columnIndex];
1357 const auto &firstReprColumnRange = fOpenColumnRanges.at(firstReprColumnId);
1359 // NOTE: this is always non-negative because it's the sum of two unsigned integers.
1360 const std::uint64_t newReprFirstElemIndex = firstReprColumnRange.GetFirstElementIndex() + clusterOffset;
1361
1363 columnBuilder.LogicalColumnId(columnId)
1364 .PhysicalColumnId(columnId)
1365 .FieldId(field.GetId())
1366 .BitsOnStorage(bitsOnStorage)
1367 .Type(columnRepr.fType)
1368 .Index(columnIndex)
1369 .FirstElementIndex(newReprFirstElemIndex)
1370 .RepresentationIndex(reprIndex)
1371 .ValueRange(columnRepr.fValueRange);
1373 columnBuilder.SetSuppressedDeferred();
1374 fDescriptorBuilder.AddColumn(columnBuilder.MoveDescriptor().Unwrap());
1375
1376 if (newReprFirstElemIndex != 0) {
1377 for (auto parentId = field.GetParentId(); parentId != ROOT::kInvalidDescriptorId;) {
1378 const ROOT::RFieldDescriptor &parent = descriptor.GetFieldDescriptor(parentId);
1381 fDescriptorBuilder.SetFeature(RNTupleDescriptor::kFeatureFlag_NestedDeferredColumns);
1382 break;
1383 }
1384 parentId = parent.GetParentId();
1385 }
1386 }
1387
1389 columnRange.SetPhysicalColumnId(columnId);
1390 columnRange.SetFirstElementIndex(firstReprColumnRange.GetFirstElementIndex());
1391 columnRange.SetNElements(0);
1392 columnRange.SetCompressionSettings(GetWriteOptions().GetCompression());
1393 fOpenColumnRanges.emplace_back(columnRange);
1394
1396 pageRange.SetPhysicalColumnId(columnId);
1397 fOpenPageRanges.emplace_back(std::move(pageRange));
1398
1399 fSerializationContext.MapPhysicalColumnId(columnId);
1400
1401 ++columnIndex;
1402 }
1403
1404 fDescriptorBuilder.EnsureValidDescriptor().ThrowOnError();
1405
1406 return firstPhysicalIndex;
1407}
1408
1412{
1413 const auto &pointedColumn = desc.GetColumnDescriptor(physicalId);
1414 assert(!pointedColumn.IsAliasColumn());
1415 assert(field.IsProjectedField());
1416
1417 const auto columnId = fDescriptorBuilder.GetDescriptor().GetNLogicalColumns();
1419 columnBuilder.LogicalColumnId(columnId)
1420 .PhysicalColumnId(physicalId)
1421 .FieldId(field.GetId())
1422 .Type(pointedColumn.GetType())
1423 .Index(pointedColumn.GetIndex())
1424 .BitsOnStorage(pointedColumn.GetBitsOnStorage())
1425 .ValueRange(pointedColumn.GetValueRange())
1426 .FirstElementIndex(pointedColumn.GetFirstElementIndex())
1427 .RepresentationIndex(pointedColumn.GetRepresentationIndex());
1428 fDescriptorBuilder.AddColumn(columnBuilder.MoveDescriptor().Unwrap());
1429
1430 fDescriptorBuilder.EnsureValidDescriptor().ThrowOnError();
1431}
1432
1434{
1435 fOpenColumnRanges.at(columnHandle.fPhysicalId).SetIsSuppressed(true);
1436}
1437
1439{
1440 fOpenColumnRanges.at(columnHandle.fPhysicalId).IncrementNElements(page.GetNElements());
1441
1442 auto element = columnHandle.fColumn->GetElement();
1444 {
1445 RNTupleAtomicTimer timer(fCounters->fTimeWallZip, fCounters->fTimeCpuZip);
1446 sealedPage = SealPage(page, *element);
1447 }
1448 fCounters->fSzZip.Add(page.GetNBytes());
1449
1451 pageInfo.SetNElements(page.GetNElements());
1452 pageInfo.SetLocator(CommitSealedPageImpl(columnHandle.fPhysicalId, sealedPage));
1453 pageInfo.SetHasChecksum(GetWriteOptions().GetEnablePageChecksums());
1454 fOpenPageRanges.at(columnHandle.fPhysicalId).GetPageInfos().emplace_back(pageInfo);
1455}
1456
1459{
1460 fOpenColumnRanges.at(physicalColumnId).IncrementNElements(sealedPage.GetNElements());
1461
1463 pageInfo.SetNElements(sealedPage.GetNElements());
1464 pageInfo.SetLocator(CommitSealedPageImpl(physicalColumnId, sealedPage));
1465 pageInfo.SetHasChecksum(sealedPage.GetHasChecksum());
1466 fOpenPageRanges.at(physicalColumnId).GetPageInfos().emplace_back(pageInfo);
1467}
1468
1469std::vector<ROOT::RNTupleLocator>
1470ROOT::Internal::RPagePersistentSink::CommitSealedPageVImpl(std::span<RPageStorage::RSealedPageGroup> ranges,
1471 const std::vector<bool> &mask)
1472{
1473 std::vector<ROOT::RNTupleLocator> locators;
1474 locators.reserve(mask.size());
1475 std::size_t i = 0;
1476 for (auto &range : ranges) {
1477 for (auto sealedPageIt = range.fFirst; sealedPageIt != range.fLast; ++sealedPageIt) {
1478 if (mask[i++])
1479 locators.push_back(CommitSealedPageImpl(range.fPhysicalColumnId, *sealedPageIt));
1480 }
1481 }
1482 locators.shrink_to_fit();
1483 return locators;
1484}
1485
1486void ROOT::Internal::RPagePersistentSink::CommitSealedPageV(std::span<RPageStorage::RSealedPageGroup> ranges)
1487{
1488 /// Used in the `originalPages` map
1489 struct RSealedPageLink {
1490 const RSealedPage *fSealedPage = nullptr; ///< Points to the first occurrence of a page with a specific checksum
1491 std::size_t fLocatorIdx = 0; ///< The index in the locator vector returned by CommitSealedPageVImpl()
1492 };
1493
1494 std::vector<bool> mask;
1495 // For every sealed page, stores the corresponding index in the locator vector returned by CommitSealedPageVImpl()
1496 std::vector<std::size_t> locatorIndexes;
1497 // Maps page checksums to the first sealed page with that checksum
1498 std::unordered_map<std::uint64_t, RSealedPageLink> originalPages;
1499 std::size_t iLocator = 0;
1500 for (auto &range : ranges) {
1501 const auto rangeSize = std::distance(range.fFirst, range.fLast);
1502 mask.reserve(mask.size() + rangeSize);
1503 locatorIndexes.reserve(locatorIndexes.size() + rangeSize);
1504
1505 for (auto sealedPageIt = range.fFirst; sealedPageIt != range.fLast; ++sealedPageIt) {
1506 if (!fFeatures.fCanMergePages || !fOptions->GetEnableSamePageMerging()) {
1507 mask.emplace_back(true);
1508 locatorIndexes.emplace_back(iLocator++);
1509 continue;
1510 }
1511 // Same page merging requires page checksums - this is checked in the write options
1512 R__ASSERT(sealedPageIt->GetHasChecksum());
1513
1514 const auto chk = sealedPageIt->GetChecksum().Unwrap();
1515 auto itr = originalPages.find(chk);
1516 if (itr == originalPages.end()) {
1517 originalPages.insert({chk, {&(*sealedPageIt), iLocator}});
1518 mask.emplace_back(true);
1519 locatorIndexes.emplace_back(iLocator++);
1520 continue;
1521 }
1522
1523 const auto *p = itr->second.fSealedPage;
1524 if ((sealedPageIt->GetDataSize() != p->GetDataSize()) ||
1525 (memcmp(sealedPageIt->GetBuffer(), p->GetBuffer(), p->GetDataSize()) != 0)) {
1526 mask.emplace_back(true);
1527 locatorIndexes.emplace_back(iLocator++);
1528 continue;
1529 }
1530
1531 mask.emplace_back(false);
1532 locatorIndexes.emplace_back(itr->second.fLocatorIdx);
1533 }
1534
1535 mask.shrink_to_fit();
1536 locatorIndexes.shrink_to_fit();
1537 }
1538
1539 auto locators = CommitSealedPageVImpl(ranges, mask);
1540 unsigned i = 0;
1541
1542 for (auto &range : ranges) {
1543 for (auto sealedPageIt = range.fFirst; sealedPageIt != range.fLast; ++sealedPageIt) {
1544 fOpenColumnRanges.at(range.fPhysicalColumnId).IncrementNElements(sealedPageIt->GetNElements());
1545
1547 pageInfo.SetNElements(sealedPageIt->GetNElements());
1548 pageInfo.SetLocator(locators[locatorIndexes[i++]]);
1549 pageInfo.SetHasChecksum(sealedPageIt->GetHasChecksum());
1550 fOpenPageRanges.at(range.fPhysicalColumnId).GetPageInfos().emplace_back(pageInfo);
1551 }
1552 }
1553}
1554
1557{
1559 stagedCluster.fNBytesWritten = StageClusterImpl();
1560 stagedCluster.fNEntries = nNewEntries;
1561
1562 for (unsigned int i = 0; i < fOpenColumnRanges.size(); ++i) {
1563 RStagedCluster::RColumnInfo columnInfo;
1564 columnInfo.fCompressionSettings = fOpenColumnRanges[i].GetCompressionSettings().value();
1565 if (fOpenColumnRanges[i].IsSuppressed()) {
1566 assert(fOpenPageRanges[i].GetPageInfos().empty());
1567 columnInfo.fPageRange.SetPhysicalColumnId(i);
1568 columnInfo.fIsSuppressed = true;
1569 // We reset suppressed columns to the state they would have if they were active (not suppressed).
1570 fOpenColumnRanges[i].SetNElements(0);
1571 fOpenColumnRanges[i].SetIsSuppressed(false);
1572 } else {
1573 std::swap(columnInfo.fPageRange, fOpenPageRanges[i]);
1574 fOpenPageRanges[i].SetPhysicalColumnId(i);
1575
1576 columnInfo.fNElements = fOpenColumnRanges[i].GetNElements();
1577 fOpenColumnRanges[i].SetNElements(0);
1578 }
1579 stagedCluster.fColumnInfos.push_back(std::move(columnInfo));
1580 }
1581
1582 return stagedCluster;
1583}
1584
1586{
1587 for (const auto &cluster : clusters) {
1589 clusterBuilder.ClusterId(fDescriptorBuilder.GetDescriptor().GetNActiveClusters())
1590 .FirstEntryIndex(fPrevClusterNEntries)
1591 .NEntries(cluster.fNEntries);
1592 for (const auto &columnInfo : cluster.fColumnInfos) {
1593 const auto colId = columnInfo.fPageRange.GetPhysicalColumnId();
1594 if (columnInfo.fIsSuppressed) {
1595 assert(columnInfo.fPageRange.GetPageInfos().empty());
1596 clusterBuilder.MarkSuppressedColumnRange(colId);
1597 } else {
1598 clusterBuilder.CommitColumnRange(colId, fOpenColumnRanges[colId].GetFirstElementIndex(),
1599 columnInfo.fCompressionSettings, columnInfo.fPageRange);
1600 fOpenColumnRanges[colId].IncrementFirstElementIndex(columnInfo.fNElements);
1601 }
1602 }
1603
1604 clusterBuilder.CommitSuppressedColumnRanges(fDescriptorBuilder.GetDescriptor()).ThrowOnError();
1605 for (const auto &columnInfo : cluster.fColumnInfos) {
1606 if (!columnInfo.fIsSuppressed)
1607 continue;
1608 const auto colId = columnInfo.fPageRange.GetPhysicalColumnId();
1609 // For suppressed columns, we need to reset the first element index to the first element of the next (upcoming)
1610 // cluster. This information has been determined for the committed cluster descriptor through
1611 // CommitSuppressedColumnRanges(), so we can use the information from the descriptor.
1612 const auto &columnRangeFromDesc = clusterBuilder.GetColumnRange(colId);
1613 fOpenColumnRanges[colId].SetFirstElementIndex(columnRangeFromDesc.GetFirstElementIndex() +
1614 columnRangeFromDesc.GetNElements());
1615 }
1616
1617 fDescriptorBuilder.AddCluster(clusterBuilder.MoveDescriptor().Unwrap());
1618 fPrevClusterNEntries += cluster.fNEntries;
1619 }
1620}
1621
1623{
1624 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1625
1626 const auto nClusters = descriptor.GetNActiveClusters();
1627 std::vector<ROOT::DescriptorId_t> physClusterIDs;
1628 physClusterIDs.reserve(nClusters);
1629 for (auto i = fNextClusterInGroup; i < nClusters; ++i) {
1630 physClusterIDs.emplace_back(fSerializationContext.MapClusterId(i));
1631 }
1632
1633 auto szPageList =
1634 RNTupleSerializer::SerializePageList(nullptr, descriptor, physClusterIDs, fSerializationContext).Unwrap();
1636 RNTupleSerializer::SerializePageList(bufPageList.get(), descriptor, physClusterIDs, fSerializationContext);
1637
1638 const auto clusterGroupId = descriptor.GetNClusterGroups();
1639 const auto locator = CommitClusterGroupImpl(bufPageList.get(), szPageList);
1641 cgBuilder.ClusterGroupId(clusterGroupId).PageListLocator(locator).PageListLength(szPageList);
1642 if (fNextClusterInGroup == nClusters) {
1643 cgBuilder.MinEntry(0).EntrySpan(0).NClusters(0);
1644 } else {
1645 const auto &firstClusterDesc = descriptor.GetClusterDescriptor(fNextClusterInGroup);
1646 const auto &lastClusterDesc = descriptor.GetClusterDescriptor(nClusters - 1);
1647 cgBuilder.MinEntry(firstClusterDesc.GetFirstEntryIndex())
1648 .EntrySpan(lastClusterDesc.GetFirstEntryIndex() + lastClusterDesc.GetNEntries() -
1649 firstClusterDesc.GetFirstEntryIndex())
1650 .NClusters(nClusters - fNextClusterInGroup);
1651 }
1652 std::vector<ROOT::DescriptorId_t> clusterIds;
1653 clusterIds.reserve(nClusters);
1654 for (auto i = fNextClusterInGroup; i < nClusters; ++i) {
1655 clusterIds.emplace_back(i);
1656 }
1657 cgBuilder.AddSortedClusters(clusterIds);
1658 fDescriptorBuilder.AddClusterGroup(cgBuilder.MoveDescriptor().Unwrap());
1659 fSerializationContext.MapClusterGroupId(clusterGroupId);
1660
1661 fNextClusterInGroup = nClusters;
1662}
1663
1666{
1668
1670 auto attrSetDesc = attrSetDescBuilder.SchemaVersion(kSchemaVersionMajor, kSchemaVersionMinor)
1671 .AnchorLength(attrAnchorInfo.fLength)
1672 .AnchorLocator(attrAnchorInfo.fLocator)
1673 .Name(attrSetName)
1674 .MoveDescriptor()
1675 .Unwrap();
1676 fDescriptorBuilder.AddAttributeSet(std::move(attrSetDesc)).ThrowOnError();
1677}
1678
1680{
1681 if (!fInfosOfStreamerFields.empty()) {
1682 // De-duplicate extra type infos before writing. Usually we won't have them already in the descriptor, but
1683 // this may happen when we are writing back an already-existing RNTuple, e.g. when doing incremental merging.
1684 for (const auto &etDesc : fDescriptorBuilder.GetDescriptor().GetExtraTypeInfoIterable()) {
1685 if (etDesc.GetContentId() == EExtraTypeInfoIds::kStreamerInfo) {
1686 // The specification mandates that the type name for a kStreamerInfo should be empty and the type version
1687 // should be zero.
1688 R__ASSERT(etDesc.GetTypeName().empty());
1689 R__ASSERT(etDesc.GetTypeVersion() == 0);
1690 auto etInfo = RNTupleSerializer::DeserializeStreamerInfos(etDesc.GetContent()).Unwrap();
1691 fInfosOfStreamerFields.merge(etInfo);
1692 }
1693 }
1694
1697 .Content(RNTupleSerializer::SerializeStreamerInfos(fInfosOfStreamerFields));
1698 fDescriptorBuilder.ReplaceExtraTypeInfo(extraInfoBuilder.MoveDescriptor().Unwrap());
1699 }
1700
1701 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1702
1703 auto szFooter = RNTupleSerializer::SerializeFooter(nullptr, descriptor, fSerializationContext).Unwrap();
1705 RNTupleSerializer::SerializeFooter(bufFooter.get(), descriptor, fSerializationContext);
1706
1707 return CommitDatasetImpl(bufFooter.get(), szFooter);
1708}
1709
1711{
1712 fMetrics = RNTupleMetrics(prefix);
1713 fCounters = std::make_unique<RCounters>(RCounters{
1714 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("nPageCommitted", "", "number of pages committed to storage"),
1715 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("szWritePayload", "B", "volume written for committed pages"),
1716 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("szZip", "B", "volume before zipping"),
1717 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("timeWallWrite", "ns", "wall clock time spent writing"),
1718 *fMetrics.MakeCounter<RNTupleAtomicCounter *>("timeWallZip", "ns", "wall clock time spent compressing"),
1719 *fMetrics.MakeCounter<RNTupleTickCounter<RNTupleAtomicCounter> *>("timeCpuWrite", "ns", "CPU time spent writing"),
1720 *fMetrics.MakeCounter<RNTupleTickCounter<RNTupleAtomicCounter> *>("timeCpuZip", "ns",
1721 "CPU time spent compressing")});
1722}
fBuffer
#define R__FORWARD_ERROR(res)
Short-hand to return an RResult<T> in an error state (i.e. after checking)
Definition RError.hxx:326
#define R__FAIL(msg)
Short-hand to return an RResult<T> in an error state; the RError is implicitly converted into RResult...
Definition RError.hxx:322
#define f(i)
Definition RSha256.hxx:104
ROOT::Detail::TRangeCast< T, true > TRangeDynCast
TRangeDynCast is an adapter class that allows the typed iteration through a TCollection.
#define R__ASSERT(e)
Checks condition e and reports a fatal error if it's false.
Definition TError.h:130
winID h TVirtualViewer3D TVirtualGLPainter p
Option_t Option_t TPoint TPoint const char GetTextMagnitude GetFillStyle GetLineColor GetLineWidth GetMarkerStyle GetTextAlign GetTextColor GetTextSize void char Point_t Rectangle_t WindowAttributes_t Float_t Float_t Float_t Int_t Int_t UInt_t UInt_t Rectangle_t mask
Option_t Option_t TPoint TPoint const char GetTextMagnitude GetFillStyle GetLineColor GetLineWidth GetMarkerStyle GetTextAlign GetTextColor GetTextSize void char Point_t Rectangle_t WindowAttributes_t Float_t Float_t Float_t Int_t Int_t UInt_t UInt_t Rectangle_t result
Option_t Option_t TPoint TPoint const char GetTextMagnitude GetFillStyle GetLineColor GetLineWidth GetMarkerStyle GetTextAlign GetTextColor GetTextSize void char Point_t Rectangle_t WindowAttributes_t index
Option_t Option_t TPoint TPoint const char mode
char name[80]
Definition TGX11.cxx:142
#define _(A, B)
Definition cfortran.h:108
A thread-safe integral performance counter.
A metric element that computes its floating point value from other counters.
A collection of Counter objects with a name, a unit, and a description.
A helper class for piece-wise construction of an RClusterDescriptor.
A helper class for piece-wise construction of an RClusterGroupDescriptor.
An in-memory subset of the packed and compressed pages of a cluster.
Definition RCluster.hxx:147
std::unordered_set< ROOT::DescriptorId_t > ColumnSet_t
Definition RCluster.hxx:149
A helper class for piece-wise construction of an RColumnDescriptor.
A column element encapsulates the translation between basic C++ types and their column representation...
static std::pair< std::uint16_t, std::uint16_t > GetValidBitRange(ROOT::ENTupleColumnType type)
Most types have a fixed on-disk bit width.
virtual RIdentifier GetIdentifier() const =0
A column is a storage-backed array of a simple, fixed-size type, from which pages can be mapped into ...
Definition RColumn.hxx:37
std::optional< std::pair< double, double > > GetValueRange() const
Definition RColumn.hxx:344
std::uint16_t GetRepresentationIndex() const
Definition RColumn.hxx:350
ROOT::Internal::RColumnElementBase * GetElement() const
Definition RColumn.hxx:337
ROOT::ENTupleColumnType GetType() const
Definition RColumn.hxx:338
ROOT::NTupleSize_t GetFirstElementIndex() const
Definition RColumn.hxx:352
std::size_t GetWritePageCapacity() const
Definition RColumn.hxx:359
std::uint16_t GetBitsOnStorage() const
Definition RColumn.hxx:339
std::uint32_t GetIndex() const
Definition RColumn.hxx:349
A helper class for piece-wise construction of an RExtraTypeInfoDescriptor.
static std::size_t Zip(const void *from, std::size_t nbytes, int compression, void *to)
Returns the size of the compressed data, written into the provided output buffer.
static void Unzip(const void *from, size_t nbytes, size_t dataLen, void *to)
The nbytes parameter provides the size ls of the from buffer.
static unsigned int GetClusterBunchSize(const RNTupleReadOptions &options)
static std::uint32_t SerializeXxHash3(const unsigned char *data, std::uint64_t length, std::uint64_t &xxhash3, void *buffer)
Writes a XxHash-3 64bit checksum of the byte range given by data and length.
static RResult< void > DeserializePageList(const void *buffer, std::uint64_t bufSize, ROOT::DescriptorId_t clusterGroupId, RNTupleDescriptor &desc, EDescriptorDeserializeMode mode)
static RResult< StreamerInfoMap_t > DeserializeStreamerInfos(const std::string &extraTypeInfoContent)
static RResult< void > VerifyXxHash3(const unsigned char *data, std::uint64_t length, std::uint64_t &xxhash3)
Expects an xxhash3 checksum in the 8 bytes following data + length and verifies it.
static RResult< std::uint32_t > SerializePageList(void *buffer, const RNTupleDescriptor &desc, std::span< ROOT::DescriptorId_t > physClusterIDs, const RContext &context)
static RResult< std::uint32_t > SerializeFooter(void *buffer, const RNTupleDescriptor &desc, const RContext &context)
static std::uint32_t DeserializeUInt64(const void *buffer, std::uint64_t &val)
static RResult< RContext > SerializeHeader(void *buffer, const RNTupleDescriptor &desc)
static std::string SerializeStreamerInfos(const StreamerInfoMap_t &infos)
A memory region that contains packed and compressed pages.
Definition RCluster.hxx:98
A page as being stored on disk, that is packed and compressed.
Definition RCluster.hxx:40
Uses standard C++ memory allocation for the column data pages.
Abstract interface to allocate and release pages.
RStagedCluster StageCluster(ROOT::NTupleSize_t nNewEntries) final
Stage the current cluster and create a new one for the following data.
void UpdateSchema(const ROOT::Internal::RNTupleModelChangeset &changeset, ROOT::NTupleSize_t firstEntry) override
Incorporate incremental changes to the model into the ntuple descriptor.
void CommitSealedPage(ROOT::DescriptorId_t physicalColumnId, const RPageStorage::RSealedPage &sealedPage) final
Write a preprocessed page to storage. The column must have been added before.
std::unique_ptr< RNTupleModel > InitFromDescriptor(const ROOT::RNTupleDescriptor &descriptor, bool copyClusters)
Initialize sink based on an existing descriptor and fill into the descriptor builder,...
void UpdateExtraTypeInfo(const ROOT::RExtraTypeInfoDescriptor &extraTypeInfo) final
Adds an extra type information record to schema.
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...
ColumnHandle_t AddColumn(ROOT::DescriptorId_t fieldId, ROOT::Internal::RColumn &column) final
Register a new column.
virtual std::vector< RNTupleLocator > CommitSealedPageVImpl(std::span< RPageStorage::RSealedPageGroup > ranges, const std::vector< bool > &mask)
Vector commit of preprocessed pages.
RPagePersistentSink(std::string_view ntupleName, const ROOT::RNTupleWriteOptions &options)
void CommitSuppressedColumn(ColumnHandle_t columnHandle) final
Commits a suppressed column for the current cluster.
void AddAliasColumn(const ROOT::RNTupleDescriptor &desc, const ROOT::RFieldDescriptor &field, ROOT::DescriptorId_t physicalId)
Adds a new alias column pointing to an existing column with the given physical id to the given field.
void CommitStagedClusters(std::span< RStagedCluster > clusters) final
Commit staged clusters, logically appending them to the ntuple descriptor.
static std::unique_ptr< RPageSink > Create(std::string_view ntupleName, std::string_view location, const ROOT::RNTupleWriteOptions &options=ROOT::RNTupleWriteOptions())
Guess the concrete derived page source from the location.
void CommitPage(ColumnHandle_t columnHandle, const ROOT::Internal::RPage &page) final
Write a page to the storage. The column must have been added before.
virtual void InitImpl(unsigned char *serializedHeader, std::uint32_t length)=0
void CommitClusterGroup() final
Write out the page locations (page list envelope) for all the committed clusters since the last call ...
void CommitSealedPageV(std::span< RPageStorage::RSealedPageGroup > ranges) final
Write a vector of preprocessed pages to storage. The corresponding columns must have been added befor...
void EnableDefaultMetrics(const std::string &prefix)
Enables the default set of metrics provided by RPageSink.
ROOT::DescriptorId_t AddColumnRepresentation(const ROOT::RFieldDescriptor &field, std::span< const ROOT::Internal::RColumnFormat > newRepresentation, std::uint64_t clusterOffset)
Adds a new column representation to the given field.
Reference to a page stored in the page pool.
Abstract interface to write data into an ntuple.
RNTupleLink CommitDataset()
Run the registered callbacks and finalize the current cluster and the entrire data set.
virtual ROOT::Internal::RPage ReservePage(ColumnHandle_t columnHandle, std::size_t nElements)
Get a new, empty page for the given column that can be filled with up to nElements; nElements must be...
RSealedPage SealPage(const ROOT::Internal::RPage &page, const ROOT::Internal::RColumnElementBase &element)
Helper for streaming a page.
RPageSink(std::string_view ntupleName, const ROOT::RNTupleWriteOptions &options)
void Insert(ROOT::DescriptorId_t physicalColumnId, ROOT::Internal::RColumnElementBase::RIdentifier elementId)
ROOT::Internal::RCluster::ColumnSet_t ToColumnSet() const
void Erase(ROOT::DescriptorId_t physicalColumnId, ROOT::Internal::RColumnElementBase::RIdentifier elementId)
An RAII wrapper used for the read-only access to RPageSource::fDescriptor. See GetExclDescriptorGuard...
RSharedDescriptorGuard FindNextClusterId(ROOT::DescriptorId_t clusterId, ROOT::DescriptorId_t &nextId)
Uses FindClusterId to search for the cluster with the entry index following the last entry index of t...
void LoadStructure()
Loads header and footer without decompressing or deserializing them.
virtual ROOT::Internal::RPageRef LoadPage(ColumnHandle_t columnHandle, ROOT::NTupleSize_t globalIndex)
Allocates and fills a page that contains the index-th element.
void RegisterStreamerInfos()
Builds the streamer info records from the descriptor's extra type info section.
ROOT::NTupleSize_t GetNElements(ROOT::DescriptorId_t physicalColumnId)
void Attach(ROOT::Internal::RNTupleSerializer::EDescriptorDeserializeMode mode=ROOT::Internal::RNTupleSerializer::EDescriptorDeserializeMode::kForReading)
Open the physical storage container and deserialize header and footer.
ColumnHandle_t AddColumn(ROOT::DescriptorId_t fieldId, ROOT::Internal::RColumn &column) override
Register a new column.
void UnzipCluster(ROOT::Internal::RCluster *cluster)
Parallel decompression and unpacking of the pages in the given cluster.
void EnableDefaultMetrics(const std::string &prefix)
Enables the default set of metrics provided by RPageSource.
ROOT::NTupleSize_t GetNEntries()
ROOT::Internal::RPageRef LoadZeroPage(ColumnHandle_t columnHandle, const RPageSummary &pageSummary)
void UpdateLastUsedCluster(ROOT::DescriptorId_t clusterId)
Does nothing if fLastUsedCluster == clusterId.
RSharedDescriptorGuard FindClusterId(ROOT::DescriptorId_t physicalColumnId, ROOT::NTupleSize_t index, ROOT::DescriptorId_t &cid)
Returns a shared descriptor guard to ensure that the returned cluster id is useable,...
ROOT::Internal::RPageRef LoadPageFromSummary(ColumnHandle_t columnHandle, const RPageSummary &pageSummary)
void DropColumn(ColumnHandle_t columnHandle) override
Unregisters a column.
void LoadSealedPage(ROOT::DescriptorId_t physicalColumnId, RNTupleLocalIndex localIndex, RSealedPage &sealedPage)
Read the packed and compressed bytes of a page into the memory buffer provided by sealedPage.
virtual void UnzipClusterImpl(ROOT::Internal::RCluster *cluster)
RPageSource(std::string_view ntupleName, const ROOT::RNTupleReadOptions &fOptions)
void PrepareLoadCluster(const ROOT::Internal::RCluster::RKey &clusterKey, ROOT::Internal::ROnDiskPageMap &pageZeroMap, const std::function< void(ROOT::DescriptorId_t, ROOT::NTupleSize_t, const ROOT::RClusterDescriptor::RPageInfo &)> &perPageFunc)
Prepare a page range read for the column set in clusterKey.
void SetEntryRange(const REntryRange &range)
Promise to only read from the given entry range.
std::unique_ptr< RPageSource > Clone() const
Open the same storage multiple time, e.g.
static std::unique_ptr< RPageSource > Create(std::string_view ntupleName, std::string_view location, const ROOT::RNTupleReadOptions &options=ROOT::RNTupleReadOptions())
Guess the concrete derived page source from the file name (location)
static RResult< ROOT::Internal::RPage > UnsealPage(const RSealedPage &sealedPage, const ROOT::Internal::RColumnElementBase &element, ROOT::Internal::RPageAllocator &pageAlloc)
Helper for unstreaming a page.
Common functionality of an ntuple storage for both reading and writing.
RPageStorage(std::string_view name)
Stores information about the cluster in which this page resides.
Definition RPage.hxx:52
A page is a slice of a column that is mapped into memory.
Definition RPage.hxx:43
static const void * GetPageZeroBuffer()
Return a pointer to the page zero buffer used if there is no on-disk data for a particular deferred c...
Definition RPage.cxx:22
const ROOT::RFieldBase * GetSourceField(const ROOT::RFieldBase *target) const
bool TryEvict(std::size_t targetAvailableSize, std::size_t pageSizeLimit)
Flush columns in order of allocated write page size until the sum of all write page allocations leave...
bool TryUpdate(ROOT::Internal::RColumn &column, std::size_t newWritePageSize)
Try to register the new write page size for the given column.
The window of element indexes of a particular column in a particular cluster.
Records the partition of data into pages for a particular column in a particular cluster.
Metadata for RNTuple clusters.
Base class for all ROOT issued exceptions.
Definition RError.hxx:78
Field specific extra type information from the header / extenstion header.
A field translates read and write calls from/to underlying columns to/from tree values.
Metadata stored for every field of an RNTuple.
ROOT::ENTupleStructure GetStructure() const
ROOT::DescriptorId_t GetParentId() const
The on-storage metadata of an RNTuple.
@ kFeatureFlag_NestedDeferredColumns
Signals that the RNTuple contains at least one deferred column that is part of a collection and was e...
Addresses a column element or field item relative to a particular cluster, instead of a global NTuple...
The RNTupleModel encapulates the schema of an RNTuple.
Common user-tunable settings for reading RNTuples.
Common user-tunable settings for storing RNTuples.
const_iterator begin() const
const_iterator end() const
void ThrowOnError()
Short-hand method to throw an exception in the case of errors.
Definition RError.hxx:312
The class is used as a return type for operations that can fail; wraps a value of type T or an RError...
Definition RError.hxx:222
ROOT::RFieldZero & GetFieldZeroOfModel(RNTupleModel &model)
RResult< void > EnsureValidNameForRNTuple(std::string_view name, std::string_view where)
Check whether a given string is a valid name according to the RNTuple specification.
RProjectedFields & GetProjectedFieldsOfModel(RNTupleModel &model)
std::unique_ptr< RColumnElementBase > GenerateColumnElement(std::type_index inMemoryType, ROOT::ENTupleColumnType onDiskType)
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.
constexpr NTupleSize_t kInvalidNTupleIndex
bool StartsWith(std::string_view string, std::string_view prefix)
std::uint64_t NTupleSize_t
Integer type long enough to hold the maximum number of entries in a column.
constexpr DescriptorId_t kInvalidDescriptorId
The identifiers that specifies the content of a (partial) cluster.
Definition RCluster.hxx:151
Every concrete RColumnElement type is identified by its on-disk type (column type) and the in-memory ...
The incremental changes to a RNTupleModel
On-disk pages within a page source are identified by the column and page number.
Definition RCluster.hxx:50
Default 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.
Default I/O performance counters that get registered in fMetrics
Used in SetEntryRange / GetEntryRange.
bool IntersectsWith(const ROOT::RClusterDescriptor &clusterDesc) const
Returns true if the given cluster has entries within the entry range.
Summarizes meta-data necessary to load a certain page. Used by LoadPageFromSummary().
A sealed page contains the bytes of a page as written to storage (packed & compressed).
RResult< void > VerifyChecksumIfEnabled() const
RResult< std::uint64_t > GetChecksum() const
Returns a failure if the sealed page has no checksum.
bool operator>(const RColumnInfo &other) const
Information about a single page in the context of a cluster's page range.