28 if (
const std::string *storagePath = std::get_if<std::string>(&
fStorage))
31 auto dir = std::get<TDirectory *>(
fStorage);
36std::unique_ptr<ROOT::Experimental::RNTupleProcessor>
42std::unique_ptr<ROOT::Experimental::RNTupleProcessor>
49 std::vector<std::unique_ptr<RNTupleProcessor>> innerProcessors;
50 innerProcessors.reserve(ntuples.size());
52 for (
auto &ntuple : ntuples) {
53 innerProcessors.emplace_back(Create(std::move(ntuple)));
56 return CreateChain(std::move(innerProcessors), options);
59std::unique_ptr<ROOT::Experimental::RNTupleProcessor>
63 if (innerProcessors.empty())
66 return std::unique_ptr<RNTupleChainProcessor>(
new RNTupleChainProcessor(std::move(innerProcessors), options));
69std::unique_ptr<ROOT::Experimental::RNTupleProcessor>
71 const std::vector<std::string> &joinFields,
74 if (joinFields.size() > 4) {
78 if (std::unordered_set(joinFields.begin(), joinFields.end()).size() < joinFields.size()) {
82 std::unique_ptr<RNTupleProcessor> primaryProcessor = Create(std::move(primaryNTuple), options);
84 std::unique_ptr<RNTupleProcessor> auxProcessor = Create(std::move(auxNTuple));
86 return CreateJoin(std::move(primaryProcessor), std::move(auxProcessor), joinFields, options);
89std::unique_ptr<ROOT::Experimental::RNTupleProcessor>
91 std::unique_ptr<RNTupleProcessor> auxProcessor,
92 const std::vector<std::string> &joinFields,
95 if (joinFields.size() > 4) {
99 if (std::unordered_set(joinFields.begin(), joinFields.end()).size() < joinFields.size()) {
103 return std::unique_ptr<RNTupleJoinProcessor>(
104 new RNTupleJoinProcessor(std::move(primaryProcessor), std::move(auxProcessor), joinFields, options));
119 std::shared_ptr<ROOT::Experimental::Internal::RNTupleProcessorEntry> entry)
122 if (IsInitialized() && fPageSource)
126 fEntry = std::make_shared<Internal::RNTupleProcessorEntry>();
128 fEntry = std::move(entry);
130 fPageSource = fNTupleSpec.CreatePageSource();
131 fPageSource->Attach();
133 fNEntries = fPageSource->GetNEntries();
139 auto desc = fPageSource->GetSharedDescriptorGuard();
140 auto fieldZeroId = desc->GetFieldZeroId();
146std::unique_ptr<ROOT::RFieldBase>
148 const std::string &typeName)
153 const std::string onDiskFieldName =
154 qualifiedFieldName.find(
"R_rntproc_join_") == 0 ? qualifiedFieldName.substr(15) : qualifiedFieldName;
156 auto descGuard = fPageSource->GetSharedDescriptorGuard();
157 const auto &desc = descGuard.GetRef();
161 const auto onDiskFieldId = desc.FindFieldId(onDiskFieldName);
167 std::unique_ptr<ROOT::RFieldBase> field;
168 if (typeName.empty()) {
169 const auto &fieldDesc = desc.GetFieldDescriptor(onDiskFieldId);
170 field = fieldDesc.CreateField(desc);
173 std::string subfieldName = onDiskFieldName;
174 auto posDot = onDiskFieldName.find_last_of(
'.');
175 if (posDot != std::string::npos)
176 subfieldName = onDiskFieldName.substr(posDot + 1);
181 field->SetOnDiskId(onDiskFieldId);
182 fieldZero.
Attach(std::move(field));
192 auto fieldIdx = fEntry->FindFieldIndex(fieldName, typeName);
195 std::string qualifiedFieldName = fieldName;
197 qualifiedFieldName = qualifiedFieldName.substr(provenance.
Get().size() + 1);
200 auto field = CreateAndConnectField(qualifiedFieldName, typeName);
203 throw RException(
R__FAIL(
"cannot register field with name \"" + qualifiedFieldName +
204 "\" because it is not present in the on-disk information of the RNTuple(s) this "
205 "processor is created from"));
208 fieldIdx = fEntry->AddField(qualifiedFieldName, std::move(field), valuePtr, provenance);
216 if (entryNumber >= fNEntries || !fEntry)
219 for (
auto fieldIdx : fFieldIdxs) {
220 fEntry->ReadValue(fieldIdx, entryNumber);
223 fNEntriesProcessed++;
228 const std::unordered_set<ROOT::Experimental::Internal::RNTupleProcessorEntry::FieldIndex_t> &fieldIdxs,
233 fFieldIdxs = fieldIdxs;
236 for (
const auto &fieldIdx : fFieldIdxs) {
237 const auto &currField = fEntry->GetValue(fieldIdx).GetField();
238 auto newField = CreateAndConnectField(fEntry->GetQualifiedFieldName(fieldIdx), currField.GetTypeName());
240 fEntry->UpdateField(fieldIdx, std::move(newField));
248 fEntry->ResetFields(fFieldIdxs);
262 static constexpr int width = 32;
264 std::string ntupleNameTrunc = fNTupleSpec.fNTupleName.substr(0,
width - 4);
265 if (ntupleNameTrunc.size() < fNTupleSpec.fNTupleName.size())
266 ntupleNameTrunc = fNTupleSpec.fNTupleName.substr(0,
width - 6) +
"..";
268 output <<
"+" << std::setfill(
'-') << std::setw(
width - 1) <<
"+\n";
269 output << std::setfill(
' ') <<
"| " << ntupleNameTrunc << std::setw(
width - 2 - ntupleNameTrunc.size()) <<
" |\n";
271 if (
const std::string *storage = std::get_if<std::string>(&fNTupleSpec.fStorage)) {
272 std::string storageTrunc = storage->substr(0,
width - 5);
273 if (storageTrunc.size() < storage->size())
274 storageTrunc = storage->substr(0,
width - 8) +
"...";
276 output << std::setfill(
' ') <<
"| " << storageTrunc << std::setw(
width - 2 - storageTrunc.size()) <<
" |\n";
278 output <<
"| " << std::setw(
width - 2) <<
" |\n";
281 output <<
"+" << std::setfill(
'-') << std::setw(
width - 1) <<
"+\n";
299 std::shared_ptr<ROOT::Experimental::Internal::RNTupleProcessorEntry> entry)
305 fEntry = std::make_shared<Internal::RNTupleProcessorEntry>();
307 fEntry = std::move(entry);
309 fInnerProcessors[0]->Initialize(fEntry);
317 for (
unsigned i = 0; i < fInnerProcessors.size(); ++i) {
319 fInnerNEntries[i] = fInnerProcessors[i]->GetNEntries();
322 fNEntries += fInnerNEntries[i];
330 const std::unordered_set<ROOT::Experimental::Internal::RNTupleProcessorEntry::FieldIndex_t> &fieldIdxs,
334 fFieldIdxs = fieldIdxs;
335 fProvenance = provenance;
336 ConnectInnerProcessor(fCurrentProcessorNumber);
341 for (
const auto &innerProc : fInnerProcessors) {
342 innerProc->Disconnect();
348 if (fCurrentProcessorNumber != processorNumber) {
349 fInnerProcessors[fCurrentProcessorNumber]->Disconnect();
350 fCurrentProcessorNumber = processorNumber;
353 auto &innerProc = fInnerProcessors[processorNumber];
354 innerProc->Initialize(fEntry);
355 innerProc->Connect(fFieldIdxs, fProvenance,
true);
363 return fInnerProcessors[fCurrentProcessorNumber]->AddFieldToEntry(fieldName, typeName, valuePtr, provenance);
372 fCurrentProcessorNumber = 0;
373 ConnectInnerProcessor(fCurrentProcessorNumber);
376 std::size_t currProcessorNumber = fCurrentProcessorNumber;
378 for (
unsigned i = 0; i < currProcessorNumber; ++i) {
380 fInnerNEntries[i] = fInnerProcessors[i]->GetNEntries();
382 entriesSeen += fInnerNEntries[i];
388 while (fInnerProcessors[currProcessorNumber]->LoadEntry(localEntryNumber) ==
kInvalidNTupleIndex) {
390 fInnerNEntries[currProcessorNumber] = fInnerProcessors[currProcessorNumber]->GetNEntries();
393 localEntryNumber -= fInnerNEntries[currProcessorNumber];
396 if (++currProcessorNumber >= fInnerProcessors.size())
399 ConnectInnerProcessor(currProcessorNumber);
402 fCurrentProcessorNumber = currProcessorNumber;
403 fNEntriesProcessed++;
404 fLastLoadedEntry = entryNumber;
411 for (
unsigned i = 0; i < fInnerProcessors.size(); ++i) {
412 const auto &innerProc = fInnerProcessors[i];
414 innerProc->Initialize(fEntry);
415 innerProc->AddEntriesToJoinTable(joinTable, entryOffset);
416 entryOffset += innerProc->GetNEntries();
422 for (
const auto &innerProc : fInnerProcessors) {
423 innerProc->PrintStructure(output);
430 std::unique_ptr<RNTupleProcessor> auxProcessor,
431 const std::vector<std::string> &joinFields,
434 fPrimaryProcessor(std::move(primaryProcessor)),
435 fAuxiliaryProcessor(std::move(auxProcessor)),
436 fJoinFieldNames(joinFields)
444 std::shared_ptr<ROOT::Experimental::Internal::RNTupleProcessorEntry> entry)
450 fEntry = std::make_shared<Internal::RNTupleProcessorEntry>();
452 fEntry = std::move(entry);
454 fPrimaryProcessor->Initialize(fEntry);
455 fAuxiliaryProcessor->Initialize(fEntry);
457 if (!fJoinFieldNames.empty()) {
458 for (
const auto &joinField : fJoinFieldNames) {
459 if (!fPrimaryProcessor->CanReadFieldFromDisk(joinField)) {
460 throw RException(
R__FAIL(
"could not find join field \"" + joinField +
"\" in primary processor \"" +
461 fPrimaryProcessor->fOptions.GetProcessorName() +
"\""));
463 if (!fAuxiliaryProcessor->CanReadFieldFromDisk(joinField)) {
464 throw RException(
R__FAIL(
"could not find join field \"" + joinField +
"\" in auxiliary processor \"" +
465 fAuxiliaryProcessor->fOptions.GetProcessorName() +
"\""));
470 auto fieldIdx = AddFieldToEntry(fOptions.GetProcessorName() +
".R_rntproc_join_" + joinField,
"std::uint64_t",
472 fJoinFieldIdxs.insert(fieldIdx);
480 const std::unordered_set<ROOT::Experimental::Internal::RNTupleProcessorEntry::FieldIndex_t> &fieldIdxs,
485 auto auxProvenance = provenance.
Evolve(fAuxiliaryProcessor->fOptions.GetProcessorName());
486 for (
const auto &fieldIdx : fieldIdxs) {
487 const auto &fieldProvenance = fEntry->GetFieldProvenance(fieldIdx);
488 if (fieldProvenance.Contains(auxProvenance))
489 fAuxiliaryFieldIdxs.insert(fieldIdx);
491 fFieldIdxs.insert(fieldIdx);
494 fPrimaryProcessor->Connect(fFieldIdxs, provenance, updateFields);
495 fAuxiliaryProcessor->Connect(fAuxiliaryFieldIdxs, auxProvenance, updateFields);
500 fPrimaryProcessor->Disconnect();
501 fAuxiliaryProcessor->Disconnect();
509 auto auxProvenance = provenance.
Evolve(fAuxiliaryProcessor->fOptions.GetProcessorName());
510 if (auxProvenance.IsPresentInFieldName(fieldName)) {
514 if (fPrimaryProcessor->CanReadFieldFromDisk(fieldName)) {
516 "\" is present in the primary RNTupleProcessor \"" +
517 fPrimaryProcessor->fOptions.GetProcessorName() +
518 "\", but may also refer to a field in the auxiliary RNTupleProcessor named \"" +
519 fAuxiliaryProcessor->fOptions.GetProcessorName() +
520 "\". To avoid this ambiguity, rename the auxiliary RNTupleProcessor."));
523 auto fieldIdx = fAuxiliaryProcessor->AddFieldToEntry(fieldName, typeName, valuePtr, auxProvenance);
525 fAuxiliaryFieldIdxs.insert(fieldIdx);
528 auto fieldIdx = fPrimaryProcessor->AddFieldToEntry(fieldName, typeName, valuePtr, provenance);
530 fFieldIdxs.insert(fieldIdx);
537 for (
const auto &fieldIdx : fAuxiliaryFieldIdxs) {
538 fEntry->SetFieldValidity(fieldIdx, isValid);
545 for (
auto fieldIdx : fFieldIdxs) {
546 fEntry->SetFieldValidity(fieldIdx,
false);
548 SetAuxiliaryFieldValidity(
false);
552 fNEntriesProcessed++;
556 fAuxiliaryProcessor->LoadEntry(entryNumber);
560 if (!fJoinTableIsBuilt) {
561 fAuxiliaryProcessor->AddEntriesToJoinTable(*fJoinTable);
562 fJoinTableIsBuilt =
true;
566 std::vector<ROOT::Experimental::Internal::RNTupleJoinTable::JoinValue_t> values;
567 values.reserve(fJoinFieldIdxs.size());
568 for (
const auto &fieldIdx : fJoinFieldIdxs) {
570 values.push_back(val);
575 const auto entryIdx = fJoinTable->GetEntryIndex(values);
578 SetAuxiliaryFieldValidity(
false);
580 SetAuxiliaryFieldValidity(
true);
581 fAuxiliaryProcessor->LoadEntry(entryIdx);
590 fNEntries = fPrimaryProcessor->GetNEntries();
597 fPrimaryProcessor->AddEntriesToJoinTable(joinTable, entryOffset);
602 std::ostringstream primaryStructureStr;
603 fPrimaryProcessor->PrintStructure(primaryStructureStr);
604 const auto primaryStructure =
ROOT::Split(primaryStructureStr.str(),
"\n",
true);
605 const auto primaryStructureWidth = primaryStructure.front().size();
607 std::ostringstream auxStructureStr;
608 fAuxiliaryProcessor->PrintStructure(auxStructureStr);
609 const auto auxStructure =
ROOT::Split(auxStructureStr.str(),
"\n",
true);
611 const auto maxLength = std::max(primaryStructure.size(), auxStructure.size());
612 for (
unsigned i = 0; i < maxLength; i++) {
613 if (i < primaryStructure.size())
614 output << primaryStructure[i];
616 output << std::setw(primaryStructureWidth) <<
"";
618 if (i < auxStructure.size())
619 output <<
" " << auxStructure[i];
#define R__FAIL(msg)
Short-hand to return an RResult<T> in an error state; the RError is implicitly converted into RResult...
Builds a join table on one or several fields of an RNTuple so it can be joined onto other RNTuples.
static std::unique_ptr< RNTupleJoinTable > Create(const std::vector< std::string > &joinFieldNames)
Create an RNTupleJoinTable from an existing RNTuple.
std::uint64_t JoinValue_t
RNTupleJoinTable & Add(ROOT::Internal::RPageSource &pageSource, PartitionKey_t partitionKey=kDefaultPartitionKey, ROOT::NTupleSize_t entryOffset=0)
Add an entry mapping to the join table.
static constexpr PartitionKey_t kDefaultPartitionKey
std::uint64_t FieldIndex_t
std::string Get() const
Get the full processor provenance, in the form of "x.y.z".
bool IsPresentInFieldName(std::string_view fieldName) const
Check whether the provided field name contains this provenance.
RNTupleProcessorProvenance Evolve(const std::string &processorName) const
Add a new processor to the provenance.
Processor specialization for vertically combined (chained) RNTupleProcessors.
void PrintStructureImpl(std::ostream &output) const final
Processor-specific implementation for printing its structure, called by PrintStructure().
void AddEntriesToJoinTable(Internal::RNTupleJoinTable &joinTable, ROOT::NTupleSize_t entryOffset=0) final
Add the entry mappings for this processor to the provided join table.
void ConnectInnerProcessor(std::size_t processorNumber)
Update the entry to reflect any missing fields in the current inner processor.
Internal::RNTupleProcessorEntry::FieldIndex_t AddFieldToEntry(const std::string &fieldName, const std::string &typeName, void *valuePtr=nullptr, const Internal::RNTupleProcessorProvenance &provenance=Internal::RNTupleProcessorProvenance()) final
Add a field to the entry.
ROOT::NTupleSize_t GetNEntries() final
Get the total number of entries in this processor.
void Initialize(std::shared_ptr< Internal::RNTupleProcessorEntry > entry=nullptr) final
Initialize the processor by creating an (initially empty) fEntry, or setting an existing one.
std::vector< ROOT::NTupleSize_t > fInnerNEntries
void Connect(const std::unordered_set< Internal::RNTupleProcessorEntry::FieldIndex_t > &fieldIdxs, const Internal::RNTupleProcessorProvenance &provenance=Internal::RNTupleProcessorProvenance(), bool updateFields=false) final
Connect the provided fields indices in the entry to their on-disk fields.
ROOT::NTupleSize_t LoadEntry(ROOT::NTupleSize_t entryNumber) final
Load the entry identified by the provided (global) entry number (i.e., considering all RNTuples in th...
void Disconnect() final
Disconnect the processor from associated physical storage.
std::vector< std::unique_ptr< RNTupleProcessor > > fInnerProcessors
Processor specialization for horizontally combined (joined) RNTupleProcessors.
void PrintStructureImpl(std::ostream &output) const final
Processor-specific implementation for printing its structure, called by PrintStructure().
ROOT::NTupleSize_t LoadEntry(ROOT::NTupleSize_t entryNumber) final
Load the entry identified by the provided entry number of the primary processor.
void AddEntriesToJoinTable(Internal::RNTupleJoinTable &joinTable, ROOT::NTupleSize_t entryOffset=0) final
Add the entry mappings for this processor to the provided join table.
void Disconnect() final
Disconnect the processor from associated physical storage.
ROOT::NTupleSize_t GetNEntries() final
Get the total number of entries in this processor.
void SetAuxiliaryFieldValidity(bool validity)
Set the validity for all fields in the auxiliary processor at once.
void Connect(const std::unordered_set< Internal::RNTupleProcessorEntry::FieldIndex_t > &fieldIdxs, const Internal::RNTupleProcessorProvenance &provenance=Internal::RNTupleProcessorProvenance(), bool updateFields=false) final
Connect the provided fields indices in the entry to their on-disk fields.
std::unique_ptr< RNTupleProcessor > fPrimaryProcessor
void Initialize(std::shared_ptr< Internal::RNTupleProcessorEntry > entry=nullptr) final
Initialize the processor by creating an (initially empty) fEntry, or setting an existing one.
Internal::RNTupleProcessorEntry::FieldIndex_t AddFieldToEntry(const std::string &fieldName, const std::string &typeName, void *valuePtr=nullptr, const Internal::RNTupleProcessorProvenance &provenance=Internal::RNTupleProcessorProvenance()) final
Add a field to the entry.
Specification of the name and location of an RNTuple, used for creating a new RNTupleProcessor.
std::variant< std::string, TDirectory * > fStorage
std::unique_ptr< ROOT::Internal::RPageSource > CreatePageSource() const
const std::string & GetProcessorName() const
void SetProcessorName(std::string_view name)
Interface for iterating over entries of vertically ("chained") and/or horizontally ("joined") combine...
static std::unique_ptr< RNTupleProcessor > CreateChain(std::vector< RNTupleOpenSpec > ntuples, const RNTupleProcessorOptions &opts=RNTupleProcessorOptions())
Create an RNTupleProcessor for a chain (i.e., a vertical combination) of RNTuples.
RNTupleProcessorOptions fOptions
friend class RNTupleJoinProcessor
static std::unique_ptr< RNTupleProcessor > CreateJoin(RNTupleOpenSpec primaryNTuple, RNTupleOpenSpec auxNTuple, const std::vector< std::string > &joinFields, const RNTupleProcessorOptions &opts=RNTupleProcessorOptions())
Create an RNTupleProcessor for a join (i.e., a horizontal combination) of RNTuples.
friend class RNTupleChainProcessor
friend class RNTupleSingleProcessor
static std::unique_ptr< RNTupleProcessor > Create(RNTupleOpenSpec ntuple, const RNTupleProcessorOptions &opts=RNTupleProcessorOptions())
Create an RNTupleProcessor for a single RNTuple.
Processor specialization for processing a single RNTuple.
void AddEntriesToJoinTable(Internal::RNTupleJoinTable &joinTable, ROOT::NTupleSize_t entryOffset=0) final
Add the entry mappings for this processor to the provided join table.
void Connect(const std::unordered_set< Internal::RNTupleProcessorEntry::FieldIndex_t > &fieldIdxs, const Internal::RNTupleProcessorProvenance &provenance=Internal::RNTupleProcessorProvenance(), bool updateFields=false) final
Connect the provided fields indices in the entry to their on-disk fields.
void Initialize(std::shared_ptr< Internal::RNTupleProcessorEntry > entry=nullptr) final
Initialize the processor by creating an (initially empty) fEntry, or setting an existing one.
void PrintStructureImpl(std::ostream &output) const final
Processor-specific implementation for printing its structure, called by PrintStructure().
void Disconnect() final
Disconnect the processor from associated physical storage.
RNTupleOpenSpec fNTupleSpec
bool CanReadFieldFromDisk(std::string_view fieldName) final
Check if a field exists on-disk and can be read by the processor.
ROOT::NTupleSize_t LoadEntry(ROOT::NTupleSize_t entryNumber) final
Load the entry identified by the provided (global) entry number (i.e., considering all RNTuples in th...
Internal::RNTupleProcessorEntry::FieldIndex_t AddFieldToEntry(const std::string &fieldName, const std::string &typeName, void *valuePtr=nullptr, const Internal::RNTupleProcessorProvenance &provenance=Internal::RNTupleProcessorProvenance()) final
Add a field to the entry.
std::unique_ptr< ROOT::RFieldBase > CreateAndConnectField(const std::string &qualifiedFieldName, const std::string &typeName)
Create a new field and connect it to the processor's page source.
static std::unique_ptr< RPageSourceFile > CreateFromAnchor(const RNTuple &anchor, const ROOT::RNTupleReadOptions &options=ROOT::RNTupleReadOptions())
Used from the RNTuple class to build a datasource if the anchor is already available.
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)
Base class for all ROOT issued exceptions.
static RResult< std::unique_ptr< RFieldBase > > Create(const std::string &fieldName, const std::string &typeName, const ROOT::RCreateFieldOptions &options, const ROOT::RNTupleDescriptor *desc, ROOT::DescriptorId_t fieldId)
Factory method to resurrect a field from the stored on-disk type information.
The container field for an ntuple model, which itself has no physical representation.
std::vector< std::unique_ptr< RFieldBase > > ReleaseSubfields()
Moves all subfields into the returned vector.
void Attach(std::unique_ptr< RFieldBase > child)
A public version of the Attach method that allows piece-wise construction of the zero field.
Representation of an RNTuple data set in a ROOT file.
void SetAllowFieldSubstitutions(RFieldZero &fieldZero, bool val)
void CallConnectPageSourceOnField(RFieldBase &, ROOT::Internal::RPageSource &)
constexpr NTupleSize_t kInvalidNTupleIndex
std::vector< std::string > Split(std::string_view str, std::string_view delims, bool skipEmpty=false)
Splits a string at each character in delims.
std::uint64_t NTupleSize_t
Integer type long enough to hold the maximum number of entries in a column.
constexpr DescriptorId_t kInvalidDescriptorId