File indexing completed on 2026-10-04 09:23:53
0001
0002
0003
0004
0005
0006
0007
0008
0009
0010
0011
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
0039
0040
0041
0042
0043
0044
0045
0046
0047
0048
0049
0050
0051
0052
0053 class RNTupleFillContext {
0054 friend class ROOT::RNTupleWriter;
0055 friend class RNTupleParallelWriter;
0056 friend class ROOT::Experimental::RNTupleAttrSetWriter;
0057
0058 private:
0059
0060
0061 std::unique_ptr<ROOT::Internal::RPageStorage::RTaskScheduler> fZipTasks;
0062 std::unique_ptr<ROOT::Internal::RPageSink> fSink;
0063
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
0071 std::size_t fUnzippedClusterSize = 0;
0072
0073 std::uint64_t fNBytesFlushed = 0;
0074
0075
0076 std::uint64_t fNBytesFilled = 0;
0077
0078 std::size_t fMaxUnzippedClusterSize;
0079
0080 std::size_t fUnzippedClusterSizeEst;
0081
0082
0083
0084 bool fStagedClusterCommitting = false;
0085
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
0124
0125
0126
0127
0128 void FillNoFlush(ROOT::REntry &entry, ROOT::RNTupleFillStatus &status) { FillNoFlushImpl(entry, status); }
0129
0130
0131
0132 std::size_t Fill(ROOT::REntry &entry) { return FillImpl(entry); }
0133
0134
0135
0136
0137
0138
0139 void FillNoFlush(ROOT::Detail::RRawPtrWriteEntry &entry, ROOT::RNTupleFillStatus &status)
0140 {
0141 FillNoFlushImpl(entry, status);
0142 }
0143
0144
0145
0146 std::size_t Fill(ROOT::Detail::RRawPtrWriteEntry &entry) { return FillImpl(entry); }
0147
0148
0149
0150 void FlushColumns();
0151
0152 void FlushCluster();
0153
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
0164 ROOT::NTupleSize_t GetLastFlushed() const { return fLastFlushed; }
0165
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 }
0182
0183 #endif