44#include <unordered_map>
91 return R__FAIL(
"page checksum verification failed, data corruption detected");
98 return R__FAIL(
"invalid attempt to extract non-existing page checksum");
100 assert(fBufferSize >= kNBytesPageChecksum);
103 reinterpret_cast<const unsigned char *
>(
fBuffer) + fBufferSize - kNBytesPageChecksum,
checksum);
127 for (std::size_t i = 0; i <
itr->second.size(); ++i) {
131 itr->second[i].fRefCounter--;
132 if (
itr->second[i].fRefCounter == 0) {
134 if (
itr->second.empty()) {
135 fColumnInfos.erase(
itr);
161 if (
clusterDesc.GetFirstEntryIndex() >= (fFirstEntry + fNEntries))
176std::unique_ptr<ROOT::Internal::RPageSource>
183 if (location.empty()) {
186 if (location.find(
"daos://") == 0)
188 return std::make_unique<ROOT::Experimental::Internal::RPageSourceDaos>(
ntupleName, location, options);
196 return std::make_unique<ROOT::Internal::RPageSourceFile>(
ntupleName, location, options);
217 if ((
range.fFirstEntry +
range.fNEntries) > GetNEntries()) {
227 fHasStructure =
true;
237 auto descGuard = GetExclDescriptorGuard();
239 fStructureBuffer.Reset();
241 std::vector<unsigned char> buffer;
243 buffer.resize(
cgDesc.GetPageListLength() +
cgDesc.GetPageListLocator().GetNBytesOnStorage());
248 cgDesc.GetPageListLength(), buffer.data());
258 auto clone = CloneImpl();
260 clone->GetExclDescriptorGuard().MoveIn(GetSharedDescriptorGuard()->Clone());
261 clone->fHasStructure =
true;
262 clone->fIsAttached =
true;
269 return GetSharedDescriptorGuard()->GetNEntries();
274 auto descGuard = GetSharedDescriptorGuard();
281 const auto &cd =
descGuard->GetClusterDescriptor(
itr->GetClusterIds().back());
303 auto descGuard = GetSharedDescriptorGuard();
306 if (
desc.GetNClusterGroups() == 0)
313 std::size_t
cgRight =
desc.GetNClusterGroups() - 1;
362 auto descGuard = GetSharedDescriptorGuard();
365 if (
desc.GetNClusterGroups() == 0)
372 std::size_t
cgRight =
desc.GetNClusterGroups() - 1;
436 RNTupleAtomicTimer
timer(fCounters->fTimeWallUnzip, fCounters->fTimeCpuUnzip);
446 std::vector<std::unique_ptr<RColumnElementBase>>
allElements;
452 if (!fActivePhysicalColumns.HasColumnInfos(
columnId))
463 for (
const auto &pi :
pageRange.GetPageInfos()) {
468 sealedPage.SetBufferSize(pi.GetLocator().GetNBytesOnStorage() + pi.HasChecksum() * kNBytesPageChecksum);
495 fCounters->fNPageUnsealed.Add(
pageNo);
499 fTaskScheduler->Wait();
502 throw RException(
R__FAIL(
"page checksum verification failed, data corruption detected"));
524 pageInfo.GetLocator().GetNBytesOnStorage()));
539 GetSharedDescriptorGuard()->GetClusterDescriptor(
clusterId).GetFirstEntryIndex();
540 auto itr = fPreloadedClusters.
begin();
542 if (fPinnedClusters.count(
itr->second) > 0) {
545 fPagePool.Evict(
itr->second);
546 itr = fPreloadedClusters.erase(
itr);
550 while ((
itr != fPreloadedClusters.
end()) &&
555 while (
itr != fPreloadedClusters.
end()) {
556 if (fPinnedClusters.count(
itr->second) > 0) {
559 fPagePool.Evict(
itr->second);
560 itr = fPreloadedClusters.erase(
itr);
589 sealedPage.VerifyChecksumIfEnabled().ThrowOnError();
640 fCounters->fNPageRead.Inc();
641 fCounters->fNRead.Inc();
642 fCounters->fSzReadPayload.Add(
sealedPage.GetBufferSize());
644 if (!fCurrentCluster || (fCurrentCluster->GetId() !=
clusterId) || !fCurrentCluster->ContainsColumn(
columnId))
645 fCurrentCluster = fClusterPool.GetCluster(
clusterId, fActivePhysicalColumns.ToColumnSet());
655 auto onDiskPage = fCurrentCluster->GetOnDiskPage(key);
662 RNTupleAtomicTimer
timer(fCounters->fTimeWallUnzip, fCounters->fTimeCpuUnzip);
669 fCounters->fNPageUnsealed.Inc();
682 UpdateLastUsedCluster(
cachedPageRef.Get().GetClusterInfo().GetId());
741 fMetrics.ObserveMetrics(fClusterPool.GetMetrics());
742 fMetrics.ObserveMetrics(fPagePool.GetMetrics());
743 fCounters = std::make_unique<RCounters>(
RCounters{
746 *fMetrics.MakeCounter<
RNTupleAtomicCounter *>(
"szReadPayload",
"B",
"volume read from storage (required)"),
747 *fMetrics.MakeCounter<
RNTupleAtomicCounter *>(
"szReadOverhead",
"B",
"volume read from storage (overhead)"),
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"),
757 "CPU time spent decompressing"),
759 "bwRead",
"MB/s",
"bandwidth compressed bytes read per second", fMetrics,
776 "bwReadUnzip",
"MB/s",
"bandwidth uncompressed bytes read per second", fMetrics,
790 "bwUnzip",
"MB/s",
"decompression bandwidth of uncompressed bytes per second", fMetrics,
804 "rtReadEfficiency",
"",
"ratio of payload over all bytes read", fMetrics,
810 return {
true, 1. / (1. + (1. *
szReadOverhead->GetValueAsInt()) / payload)};
816 *fMetrics.MakeCounter<
RNTupleCalcPerf *>(
"rtCompression",
"",
"ratio of compressed bytes / uncompressed bytes",
819 metrics.GetLocalCounter(
"szReadPayload")) {
877 if (fHasStreamerInfosRegistered)
880 for (
const auto &
extraTypeInfo : fDescriptor.GetExtraTypeInfoIterable()) {
888 fHasStreamerInfosRegistered =
true;
896 if (fCurrentPageSize ==
other.fCurrentPageSize)
897 return fColumn->GetOnDiskId() >
other.fColumn->GetOnDiskId();
898 return fCurrentPageSize >
other.fCurrentPageSize;
906 auto itr = fColumnsSortedByPageSize.
begin();
907 while (
itr != fColumnsSortedByPageSize.
end()) {
910 if (
itr->fCurrentPageSize ==
itr->fInitialPageSize) {
919 if (
itr != fColumnsSortedByPageSize.
end())
928 itr = fColumnsSortedByPageSize.find(next);
937 auto itr = fColumnsSortedByPageSize.find(key);
938 if (
itr == fColumnsSortedByPageSize.
end()) {
952 fColumnsSortedByPageSize.erase(
itr);
958 fColumnsSortedByPageSize.insert(
elem);
967 fColumnsSortedByPageSize.insert(
elem);
972 fColumnsSortedByPageSize.insert(
elem);
979 :
RPageStorage(
name), fOptions(options.Clone()), fWritePageMemoryManager(options.GetPageBufferBudget())
992 unsigned char *
pageBuf =
reinterpret_cast<unsigned char *
>(config.
fPage->GetBuffer());
997 if (!config.
fElement->IsMappable()) {
1026 const auto nBytes =
page.GetNBytes() + GetWriteOptions().GetEnablePageChecksums() * kNBytesPageChecksum;
1027 if (fSealPageBuffer.size() <
nBytes)
1028 fSealPageBuffer.resize(
nBytes);
1034 config.
fWriteChecksum = GetWriteOptions().GetEnablePageChecksums();
1036 config.
fBuffer = fSealPageBuffer.data();
1038 return SealPage(config);
1043 for (
const auto &
cb : fOnDatasetCommitCallbacks)
1045 return CommitDatasetImpl();
1060std::unique_ptr<ROOT::Internal::RPageSink>
1067 if (location.empty()) {
1070 if (location.find(
"daos://") == 0) {
1071#ifdef R__ENABLE_DAOS
1072 return std::make_unique<ROOT::Experimental::Internal::RPageSinkDaos>(
ntupleName, location, options);
1080 return std::make_unique<ROOT::Experimental::Internal::RPageSinkS3>(
ntupleName, location, options);
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."));
1088 return std::make_unique<ROOT::Internal::RPageSinkFile>(
ntupleName, location, options);
1102 auto columnId = fDescriptorBuilder.GetDescriptor().GetNPhysicalColumns();
1117 fDescriptorBuilder.AddColumn(
columnBuilder.MoveDescriptor().Unwrap());
1124 if (fIsInitialized) {
1132 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1134 if (descriptor.GetNLogicalColumns() > descriptor.GetNPhysicalColumns()) {
1138 const auto &
reps =
f.GetColumnRepresentatives();
1141 return reps.size() *
reps[0].size();
1153 auto fieldId = descriptor.GetNFields();
1154 fDescriptorBuilder.AddField(
f,
fieldId);
1155 fDescriptorBuilder.AddFieldLink(
f.GetParent()->GetOnDiskId(),
fieldId);
1160 auto fieldId = descriptor.GetNFields();
1163 fDescriptorBuilder.AddField(
f,
fieldId);
1164 fDescriptorBuilder.AddFieldLink(
f.GetParent()->GetOnDiskId(),
fieldId);
1168 auto targetId = descriptor.GetNLogicalColumns();
1171 .PhysicalColumnId(
source.GetLogicalId())
1173 .BitsOnStorage(
source.GetBitsOnStorage())
1174 .ValueRange(
source.GetValueRange())
1176 .Index(
source.GetIndex())
1177 .RepresentationIndex(
source.GetRepresentationIndex());
1178 fDescriptorBuilder.AddColumn(
columnBuilder.MoveDescriptor().Unwrap());
1189 for (
auto f :
changeset.fAddedProjectedFields) {
1195 const auto nColumns = descriptor.GetNPhysicalColumns();
1204 columnRange.SetFirstElementIndex(descriptor.GetColumnDescriptor(i).GetFirstElementIndex());
1206 columnRange.SetCompressionSettings(GetWriteOptions().GetCompression());
1210 fOpenPageRanges.emplace_back(std::move(
pageRange));
1215 if (fSerializationContext.GetHeaderSize() > 0)
1216 fSerializationContext.MapSchema(descriptor,
true);
1222 throw RException(
R__FAIL(
"ROOT bug: unexpected type extra info in UpdateExtraTypeInfo()"));
1229 fDescriptorBuilder.SetNTuple(fNTupleName, model.GetDescription());
1230 fDescriptorBuilder.SetVersionForWriting();
1231 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1234 fDescriptorBuilder.AddField(
fieldZero, 0);
1241 for (
auto f :
fieldZero.GetMutableSubfields())
1251 InitImpl(buffer.get(), fSerializationContext.GetHeaderSize());
1253 fDescriptorBuilder.BeginHeaderExtension();
1256std::unique_ptr<ROOT::RNTupleModel>
1262 fDescriptorBuilder.SetVersionForWriting();
1263 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1266 const auto nColumns = descriptor.GetNPhysicalColumns();
1267 R__ASSERT(fOpenColumnRanges.empty() && fOpenPageRanges.empty());
1268 fOpenColumnRanges.reserve(
nColumns);
1271 const auto &column = descriptor.GetColumnDescriptor(i);
1274 columnRange.SetFirstElementIndex(column.GetFirstElementIndex());
1276 columnRange.SetCompressionSettings(GetWriteOptions().GetCompression());
1280 fOpenPageRanges.emplace_back(std::move(
pageRange));
1288 for (
unsigned int i = 0; i < fOpenColumnRanges.size(); ++i) {
1289 R__ASSERT(fOpenColumnRanges[i].GetPhysicalColumnId() == i);
1290 if (!
cluster.ContainsColumn(i))
1295 fOpenColumnRanges[i].IncrementFirstElementIndex(
columnRange.GetNElements());
1297 fDescriptorBuilder.AddCluster(
cluster.Clone());
1304 modelOpts.SetReconstructProjections(
true);
1307 auto model = descriptor.CreateModel(
modelOpts);
1310 projectedFields.GetFieldZero().SetOnDiskId(model->GetConstFieldZero().GetOnDiskId());
1317 InitImpl(buffer.get(), fSerializationContext.GetHeaderSize());
1319 fDescriptorBuilder.BeginHeaderExtension();
1322 fIsInitialized =
true;
1332 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1340 const std::size_t
firstPhysicalIndex = fDescriptorBuilder.GetDescriptor().GetNPhysicalColumns();
1341 const std::uint16_t
reprIndex =
field.GetLogicalColumnIds().size() /
field.GetColumnCardinality();
1365 .FieldId(
field.GetId())
1374 fDescriptorBuilder.AddColumn(
columnBuilder.MoveDescriptor().Unwrap());
1392 columnRange.SetCompressionSettings(GetWriteOptions().GetCompression());
1397 fOpenPageRanges.emplace_back(std::move(
pageRange));
1399 fSerializationContext.MapPhysicalColumnId(
columnId);
1404 fDescriptorBuilder.EnsureValidDescriptor().ThrowOnError();
1417 const auto columnId = fDescriptorBuilder.GetDescriptor().GetNLogicalColumns();
1421 .FieldId(
field.GetId())
1427 .RepresentationIndex(
pointedColumn.GetRepresentationIndex());
1428 fDescriptorBuilder.AddColumn(
columnBuilder.MoveDescriptor().Unwrap());
1430 fDescriptorBuilder.EnsureValidDescriptor().ThrowOnError();
1435 fOpenColumnRanges.at(
columnHandle.fPhysicalId).SetIsSuppressed(
true);
1440 fOpenColumnRanges.at(
columnHandle.fPhysicalId).IncrementNElements(
page.GetNElements());
1445 RNTupleAtomicTimer
timer(fCounters->fTimeWallZip, fCounters->fTimeCpuZip);
1448 fCounters->fSzZip.Add(
page.GetNBytes());
1453 pageInfo.SetHasChecksum(GetWriteOptions().GetEnablePageChecksums());
1469std::vector<ROOT::RNTupleLocator>
1471 const std::vector<bool> &
mask)
1473 std::vector<ROOT::RNTupleLocator>
locators;
1476 for (
auto &
range : ranges) {
1494 std::vector<bool>
mask;
1498 std::unordered_map<std::uint64_t, RSealedPageLink>
originalPages;
1500 for (
auto &
range : ranges) {
1506 if (!fFeatures.fCanMergePages || !fOptions->GetEnableSamePageMerging()) {
1507 mask.emplace_back(
true);
1518 mask.emplace_back(
true);
1523 const auto *
p =
itr->second.fSealedPage;
1526 mask.emplace_back(
true);
1531 mask.emplace_back(
false);
1535 mask.shrink_to_fit();
1542 for (
auto &
range : ranges) {
1544 fOpenColumnRanges.at(
range.fPhysicalColumnId).IncrementNElements(
sealedPageIt->GetNElements());
1550 fOpenPageRanges.at(
range.fPhysicalColumnId).GetPageInfos().emplace_back(
pageInfo);
1562 for (
unsigned int i = 0; i < fOpenColumnRanges.size(); ++i) {
1564 columnInfo.fCompressionSettings = fOpenColumnRanges[i].GetCompressionSettings().value();
1565 if (fOpenColumnRanges[i].IsSuppressed()) {
1566 assert(fOpenPageRanges[i].GetPageInfos().empty());
1567 columnInfo.fPageRange.SetPhysicalColumnId(i);
1570 fOpenColumnRanges[i].SetNElements(0);
1571 fOpenColumnRanges[i].SetIsSuppressed(
false);
1573 std::swap(
columnInfo.fPageRange, fOpenPageRanges[i]);
1574 fOpenPageRanges[i].SetPhysicalColumnId(i);
1576 columnInfo.fNElements = fOpenColumnRanges[i].GetNElements();
1577 fOpenColumnRanges[i].SetNElements(0);
1589 clusterBuilder.ClusterId(fDescriptorBuilder.GetDescriptor().GetNActiveClusters())
1590 .FirstEntryIndex(fPrevClusterNEntries)
1600 fOpenColumnRanges[
colId].IncrementFirstElementIndex(
columnInfo.fNElements);
1604 clusterBuilder.CommitSuppressedColumnRanges(fDescriptorBuilder.GetDescriptor()).ThrowOnError();
1617 fDescriptorBuilder.AddCluster(
clusterBuilder.MoveDescriptor().Unwrap());
1618 fPrevClusterNEntries +=
cluster.fNEntries;
1624 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
1626 const auto nClusters = descriptor.GetNActiveClusters();
1629 for (
auto i = fNextClusterInGroup; i <
nClusters; ++i) {
1630 physClusterIDs.emplace_back(fSerializationContext.MapClusterId(i));
1643 cgBuilder.MinEntry(0).EntrySpan(0).NClusters(0);
1645 const auto &
firstClusterDesc = descriptor.GetClusterDescriptor(fNextClusterInGroup);
1650 .NClusters(
nClusters - fNextClusterInGroup);
1652 std::vector<ROOT::DescriptorId_t>
clusterIds;
1654 for (
auto i = fNextClusterInGroup; i <
nClusters; ++i) {
1658 fDescriptorBuilder.AddClusterGroup(
cgBuilder.MoveDescriptor().Unwrap());
1676 fDescriptorBuilder.AddAttributeSet(std::move(
attrSetDesc)).ThrowOnError();
1681 if (!fInfosOfStreamerFields.empty()) {
1684 for (
const auto &
etDesc : fDescriptorBuilder.GetDescriptor().GetExtraTypeInfoIterable()) {
1691 fInfosOfStreamerFields.merge(
etInfo);
1698 fDescriptorBuilder.ReplaceExtraTypeInfo(
extraInfoBuilder.MoveDescriptor().Unwrap());
1701 const auto &descriptor = fDescriptorBuilder.GetDescriptor();
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"),
1717 *fMetrics.MakeCounter<
RNTupleAtomicCounter *>(
"timeWallWrite",
"ns",
"wall clock time spent writing"),
1718 *fMetrics.MakeCounter<
RNTupleAtomicCounter *>(
"timeWallZip",
"ns",
"wall clock time spent compressing"),
1721 "CPU time spent compressing")});
#define R__FORWARD_ERROR(res)
Short-hand to return an RResult<T> in an error state (i.e. after checking)
#define R__FAIL(msg)
Short-hand to return an RResult<T> in an error state; the RError is implicitly converted into RResult...
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.
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
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.
std::unordered_set< ROOT::DescriptorId_t > ColumnSet_t
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 ...
std::optional< std::pair< double, double > > GetValueRange() const
std::uint16_t GetRepresentationIndex() const
ROOT::Internal::RColumnElementBase * GetElement() const
ROOT::ENTupleColumnType GetType() const
ROOT::NTupleSize_t GetFirstElementIndex() const
std::size_t GetWritePageCapacity() const
std::uint16_t GetBitsOnStorage() const
std::uint32_t GetIndex() const
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.
EDescriptorDeserializeMode
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.
A page as being stored on disk, that is packed and compressed.
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.
RNTupleLink CommitDatasetImpl() final
~RPagePersistentSink() override
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.
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.
A page is a slice of a column that is mapped into memory.
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...
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.
Metadata for RNTuple clusters.
Base class for all ROOT issued exceptions.
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.
The class is used as a return type for operations that can fail; wraps a value of type T or an RError...
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.
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.
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.
ROOT::Internal::RColumn * fColumn
bool operator>(const RColumnInfo &other) const
Information about a single page in the context of a cluster's page range.
Modifiers passed to CreateModel()