Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-10-04 09:23:53

0001 /// \file ROOT/RNTupleFillContext.hxx
0002 /// \ingroup NTuple
0003 /// \author Jakob Blomer <jblomer@cern.ch>
0004 /// \date 2024-02-22
0005 
0006 /*************************************************************************
0007  * Copyright (C) 1995-2024, Rene Brun and Fons Rademakers.               *
0008  * All rights reserved.                                                  *
0009  *                                                                       *
0010  * For the licensing terms see $ROOTSYS/LICENSE.                         *
0011  * For the list of contributors see $ROOTSYS/README/CREDITS.             *
0012  *************************************************************************/
0013 
0014 #ifndef ROOT_RNTupleFillContext
0015 #define ROOT_RNTupleFillContext
0016 
0017 #include <ROOT/RConfig.hxx> // for R__unlikely
0018 #include <ROOT/REntry.hxx>
0019 #include <ROOT/RError.hxx>
0020 #include <ROOT/RPageStorage.hxx>
0021 #include <ROOT/RRawPtrWriteEntry.hxx>
0022 #include <ROOT/RNTupleFillStatus.hxx>
0023 #include <ROOT/RNTupleMetrics.hxx>
0024 #include <ROOT/RNTupleModel.hxx>
0025 #include <ROOT/RNTupleTypes.hxx>
0026 
0027 #include <cstddef>
0028 #include <cstdint>
0029 #include <memory>
0030 #include <vector>
0031 
0032 namespace ROOT {
0033 
0034 namespace Experimental {
0035 class RNTupleAttrSetWriter;
0036 }
0037 
0038 // clang-format off
0039 /**
0040 \class ROOT::RNTupleFillContext
0041 \ingroup NTuple
0042 \brief A context for filling entries (data) into clusters of an RNTuple
0043 
0044 An output cluster can be filled with entries. The caller has to make sure that the data that gets filled into a cluster
0045 is not modified for the time of the Fill() call. The fill call serializes the C++ object into the column format and
0046 writes data into the corresponding column page buffers.  Writing of the buffers to storage is deferred and can be
0047 triggered by FlushCluster() or by destructing the context.  On I/O errors, an exception is thrown.
0048 
0049 Instances of this class are not meant to be used in isolation and can be created from an RNTupleParallelWriter. For
0050 sequential writing, please refer to RNTupleWriter.
0051 */
0052 // clang-format on
0053 class RNTupleFillContext {
0054    friend class ROOT::RNTupleWriter;
0055    friend class RNTupleParallelWriter;
0056    friend class ROOT::Experimental::RNTupleAttrSetWriter;
0057 
0058 private:
0059    /// The page sink's parallel page compression scheduler if IMT is on.
0060    /// Needs to be destructed after the page sink is destructed and so declared before.
0061    std::unique_ptr<ROOT::Internal::RPageStorage::RTaskScheduler> fZipTasks;
0062    std::unique_ptr<ROOT::Internal::RPageSink> fSink;
0063    /// Needs to be destructed before fSink
0064    std::unique_ptr<ROOT::RNTupleModel> fModel;
0065 
0066    Experimental::Detail::RNTupleMetrics fMetrics;
0067 
0068    ROOT::NTupleSize_t fLastFlushed = 0;
0069    ROOT::NTupleSize_t fNEntries = 0;
0070    /// Keeps track of the number of bytes written into the current cluster
0071    std::size_t fUnzippedClusterSize = 0;
0072    /// The total number of bytes written to storage (i.e., after compression)
0073    std::uint64_t fNBytesFlushed = 0;
0074    /// The total number of bytes filled into all the so far committed clusters,
0075    /// i.e. the uncompressed size of the written clusters
0076    std::uint64_t fNBytesFilled = 0;
0077    /// Limit for committing cluster no matter the other tunables
0078    std::size_t fMaxUnzippedClusterSize;
0079    /// Estimator of uncompressed cluster size, taking into account the estimated compression ratio
0080    std::size_t fUnzippedClusterSizeEst;
0081 
0082    /// Whether to enable staged cluster committing, where only an explicit call to CommitStagedClusters() will logically
0083    /// append the clusters to the RNTuple.
0084    bool fStagedClusterCommitting = false;
0085    /// Vector of currently staged clusters.
0086    std::vector<ROOT::Internal::RPageSink::RStagedCluster> fStagedClusters;
0087 
0088    template <typename Entry>
0089    void FillNoFlushImpl(Entry &entry, ROOT::RNTupleFillStatus &status)
0090    {
0091       if (R__unlikely(entry.GetModelId() != fModel->GetModelId()))
0092          throw RException(R__FAIL("mismatch between entry and model"));
0093 
0094       const std::size_t bytesWritten = entry.Append();
0095       fUnzippedClusterSize += bytesWritten;
0096       fNEntries++;
0097 
0098       status.fNEntriesSinceLastFlush = fNEntries - fLastFlushed;
0099       status.fUnzippedClusterSize = fUnzippedClusterSize;
0100       status.fLastEntrySize = bytesWritten;
0101       status.fShouldFlushCluster =
0102          (fUnzippedClusterSize >= fMaxUnzippedClusterSize) || (fUnzippedClusterSize >= fUnzippedClusterSizeEst);
0103    }
0104    template <typename Entry>
0105    std::size_t FillImpl(Entry &entry)
0106    {
0107       ROOT::RNTupleFillStatus status;
0108       FillNoFlushImpl(entry, status);
0109       if (status.ShouldFlushCluster())
0110          FlushCluster();
0111       return status.GetLastEntrySize();
0112    }
0113 
0114    RNTupleFillContext(std::unique_ptr<ROOT::RNTupleModel> model, std::unique_ptr<ROOT::Internal::RPageSink> sink);
0115    RNTupleFillContext(const RNTupleFillContext &) = delete;
0116    RNTupleFillContext &operator=(const RNTupleFillContext &) = delete;
0117    RNTupleFillContext(RNTupleFillContext &&) = delete;
0118    RNTupleFillContext &operator=(RNTupleFillContext &&) = delete;
0119 
0120 public:
0121    ~RNTupleFillContext();
0122 
0123    /// Fill an entry into this context, but don't commit the cluster. The calling code must pass an RNTupleFillStatus
0124    /// and check RNTupleFillStatus::ShouldFlushCluster.
0125    ///
0126    /// This method will check the entry's model ID to ensure it comes from the context's own model or throw an exception
0127    /// otherwise.
0128    void FillNoFlush(ROOT::REntry &entry, ROOT::RNTupleFillStatus &status) { FillNoFlushImpl(entry, status); }
0129    /// Fill an entry into this context.  This method will check the entry's model ID to ensure it comes from the
0130    /// context's own model or throw an exception otherwise.
0131    /// \return The number of uncompressed bytes written.
0132    std::size_t Fill(ROOT::REntry &entry) { return FillImpl(entry); }
0133 
0134    /// Fill an RRawPtrWriteEntry into this context, but don't commit the cluster. The calling code must pass an
0135    /// RNTupleFillStatus and check RNTupleFillStatus::ShouldFlushCluster.
0136    ///
0137    /// This method will check the entry's model ID to ensure it comes from the context's own model or throw an exception
0138    /// otherwise.
0139    void FillNoFlush(ROOT::Detail::RRawPtrWriteEntry &entry, ROOT::RNTupleFillStatus &status)
0140    {
0141       FillNoFlushImpl(entry, status);
0142    }
0143    /// Fill an RRawPtrWriteEntry into this context.  This method will check the entry's model ID to ensure it comes
0144    /// from the context's own model or throw an exception otherwise.
0145    /// \return The number of uncompressed bytes written.
0146    std::size_t Fill(ROOT::Detail::RRawPtrWriteEntry &entry) { return FillImpl(entry); }
0147 
0148    /// Flush column data, preparing for CommitCluster or to reduce memory usage. This will trigger compression of pages,
0149    /// but not actually write to storage.
0150    void FlushColumns();
0151    /// Flush so far filled entries to storage
0152    void FlushCluster();
0153    /// Logically append staged clusters to the RNTuple.
0154    void CommitStagedClusters();
0155 
0156    const ROOT::RNTupleModel &GetModel() const { return *fModel; }
0157    std::unique_ptr<ROOT::REntry> CreateEntry() const { return fModel->CreateEntry(); }
0158    std::unique_ptr<ROOT::Detail::RRawPtrWriteEntry> CreateRawPtrWriteEntry() const
0159    {
0160       return fModel->CreateRawPtrWriteEntry();
0161    }
0162 
0163    /// Return the entry number that was last flushed in a cluster.
0164    ROOT::NTupleSize_t GetLastFlushed() const { return fLastFlushed; }
0165    /// Return the number of entries filled so far.
0166    ROOT::NTupleSize_t GetNEntries() const { return fNEntries; }
0167 
0168    void EnableStagedClusterCommitting(bool val = true)
0169    {
0170       if (!val && !fStagedClusters.empty()) {
0171          throw RException(R__FAIL("cannot disable staged committing with pending clusters"));
0172       }
0173       fStagedClusterCommitting = val;
0174    }
0175    bool IsStagedClusterCommittingEnabled() const { return fStagedClusterCommitting; }
0176 
0177    void EnableMetrics() { fMetrics.Enable(); }
0178    const Experimental::Detail::RNTupleMetrics &GetMetrics() const { return fMetrics; }
0179 };
0180 
0181 } // namespace ROOT
0182 
0183 #endif // ROOT_RNTupleFillContext