62 StopBackgroundThread();
67 if (fThreadIo.joinable())
75 if (!fThreadIo.joinable())
80 std::unique_lock<std::mutex> lock(fLockWorkQueue);
82 fCvHasReadWork.notify_one();
89 std::deque<RReadItem> readItems;
92 std::unique_lock<std::mutex> lock(fLockWorkQueue);
93 fCvHasReadWork.wait(lock, [&]{
return !fReadQueue.empty(); });
94 std::swap(readItems, fReadQueue);
97 while (!readItems.empty()) {
98 std::vector<RCluster::RKey> clusterKeys;
99 std::int64_t bunchId = -1;
100 for (
unsigned i = 0; i < readItems.size(); ++i) {
101 const auto &item = readItems[i];
108 if ((bunchId >= 0) && (item.fBunchId != bunchId))
110 bunchId = item.fBunchId;
111 clusterKeys.emplace_back(item.fClusterKey);
114 auto clusters = fPageSource.LoadClusters(clusterKeys);
115 for (std::size_t i = 0; i < clusters.size(); ++i) {
116 readItems[i].fPromise.set_value(std::move(clusters[i]));
118 readItems.erase(readItems.begin(), readItems.begin() + clusters.size());
132 std::int64_t fBunchId = -1;
134 ColumnSet_t fPhysicalColumnSet;
137 static constexpr std::int64_t kFlagRequired = 0x01;
138 static constexpr std::int64_t kFlagLast = 0x02;
141 std::map<DescriptorId_t, RInfo> fMap;
146 fMap.emplace(clusterId, info);
149 bool Contains(DescriptorId_t clusterId) {
150 return fMap.count(clusterId) > 0;
153 std::size_t GetSize()
const {
return fMap.size(); }
155 void Erase(DescriptorId_t clusterId,
const ColumnSet_t &physicalColumns)
157 auto itr = fMap.find(clusterId);
158 if (itr == fMap.end())
161 std::copy_if(itr->second.fPhysicalColumnSet.begin(), itr->second.fPhysicalColumnSet.end(),
162 std::inserter(
d,
d.end()),
163 [&physicalColumns](DescriptorId_t needle) { return physicalColumns.count(needle) == 0; });
167 itr->second.fPhysicalColumnSet =
d;
171 decltype(fMap)::iterator begin() {
return fMap.begin(); }
172 decltype(fMap)::iterator end() {
return fMap.end(); }
180 StartBackgroundThread();
182 std::unordered_set<ROOT::DescriptorId_t> keep{fPageSource.GetPinnedClusters()};
183 for (
auto cid : fPageSource.GetPinnedClusters()) {
185 const auto currentId = next;
186 auto descriptorGuard = fPageSource.FindNextClusterId(currentId, next);
188 !fPageSource.GetEntryRange().IntersectsWith(descriptorGuard->GetClusterDescriptor(next))) {
198 RProvides::RInfo provideInfo;
199 provideInfo.fPhysicalColumnSet = physicalColumns;
200 provideInfo.fBunchId = fBunchId;
201 provideInfo.fFlags = RProvides::kFlagRequired;
203 if (i == fClusterBunchSize)
204 provideInfo.fBunchId = ++fBunchId;
207 auto descriptorGuard = fPageSource.FindNextClusterId(cid, next);
209 if (!fPageSource.GetEntryRange().IntersectsWith(descriptorGuard->GetClusterDescriptor(next)))
213 provideInfo.fFlags |= RProvides::kFlagLast;
215 provide.Insert(cid, provideInfo);
219 provideInfo.fFlags = 0;
223 for (
auto itr = fPool.begin(); itr != fPool.end();) {
224 if (provide.Contains(itr->first)) {
228 if (keep.count(itr->first) > 0) {
232 itr = fPool.erase(itr);
233 fCounters->fNCluster.Dec();
243 std::lock_guard<std::mutex> lockGuard(fLockWorkQueue);
245 for (
auto itr = fInFlightClusters.begin(); itr != fInFlightClusters.end(); ) {
247 if (itr->fFuture.wait_for(std::chrono::seconds(0)) != std::future_status::ready) {
249 provide.Erase(itr->fClusterKey.fClusterId, itr->fClusterKey.fPhysicalColumnSet);
254 auto cptr = itr->fFuture.get();
257 const bool isExpired =
258 !provide.Contains(itr->fClusterKey.fClusterId) && (keep.count(itr->fClusterKey.fClusterId) == 0);
261 itr = fInFlightClusters.erase(itr);
266 fPageSource.UnzipCluster(cptr.get());
269 auto existingCluster = fPool.find(cptr->GetId());
270 if (existingCluster != fPool.end()) {
271 existingCluster->second->Adopt(std::move(*cptr));
273 const auto cid = cptr->GetId();
274 fPool.emplace(cid, std::move(cptr));
275 fCounters->fNCluster.Inc();
277 itr = fInFlightClusters.erase(itr);
281 for (
const auto &[
_, cptr] : fPool) {
282 provide.Erase(cptr->GetId(), cptr->GetAvailPhysicalColumns());
286 bool skipPrefetch =
false;
287 if (provide.GetSize() < fClusterBunchSize) {
289 for (
const auto &kv : provide) {
290 if ((kv.second.fFlags & (RProvides::kFlagRequired | RProvides::kFlagLast)) == 0)
292 skipPrefetch =
false;
302 for (
const auto &kv : provide) {
303 R__ASSERT(!kv.second.fPhysicalColumnSet.empty());
307 readItem.
fBunchId = kv.second.fBunchId;
314 fInFlightClusters.emplace_back(std::move(inFlightCluster));
316 fReadQueue.emplace_back(std::move(readItem));
318 if (!fReadQueue.empty())
319 fCvHasReadWork.notify_one();
323 return WaitFor(clusterId, physicalColumns);
331 auto result = fPool.find(clusterId);
332 if (
result != fPool.end()) {
333 bool hasMissingColumn =
false;
334 for (
auto cid : physicalColumns) {
335 if (
result->second->ContainsColumn(cid))
338 hasMissingColumn =
true;
341 if (!hasMissingColumn)
342 return result->second.get();
346 decltype(fInFlightClusters)::iterator itr;
348 std::lock_guard<std::mutex> lockGuardInFlightClusters(fLockWorkQueue);
349 itr = fInFlightClusters.begin();
350 for (; itr != fInFlightClusters.end(); ++itr) {
351 if (itr->fClusterKey.fClusterId == clusterId)
354 R__ASSERT(itr != fInFlightClusters.end());
361 auto cptr = itr->fFuture.get();
366 fPageSource.UnzipCluster(cptr.get());
368 if (
result != fPool.end()) {
369 result->second->Adopt(std::move(*cptr));
371 const auto cid = cptr->GetId();
372 fPool.emplace(cid, std::move(cptr));
373 fCounters->fNCluster.Inc();
376 std::lock_guard<std::mutex> lockGuardInFlightClusters(fLockWorkQueue);
377 fInFlightClusters.erase(itr);
384 decltype(fInFlightClusters)::iterator itr;
386 std::lock_guard<std::mutex> lockGuardInFlightClusters(fLockWorkQueue);
387 itr = fInFlightClusters.begin();
388 while (itr != fInFlightClusters.end() &&
389 itr->fFuture.wait_for(std::chrono::seconds(0)) == std::future_status::ready) {
392 if (itr == fInFlightClusters.end())
#define R__unlikely(expr)
#define R__ASSERT(e)
Checks condition e and reports a fatal error if it's false.
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
A thread-safe integral performance counter.
CounterPtrT MakeCounter(const std::string &name, Args &&... args)
RCluster * WaitFor(ROOT::DescriptorId_t clusterId, const RCluster::ColumnSet_t &physicalColumns)
Returns the given cluster from the pool, which needs to contain at least the columns physicalColumns.
ROOT::Experimental::Detail::RNTupleMetrics fMetrics
The cluster pool counters are observed by the page source.
unsigned int fClusterBunchSize
The number of clusters that are being read in a single vector read.
void WaitForInFlightClusters()
Used by the unit tests to drain the queue of clusters to be preloaded.
std::unique_ptr< RCounters > fCounters
void StopBackgroundThread()
Stop the I/O background thread. No-op if already stopped. Called by the destructor.
void ExecReadClusters()
The I/O thread routine, there is exactly one I/O thread in-flight for every cluster pool.
RCluster * GetCluster(ROOT::DescriptorId_t clusterId, const RCluster::ColumnSet_t &physicalColumns)
Returns the requested cluster either from the pool or, in case of a cache miss, lets the I/O thread l...
RClusterPool(ROOT::Internal::RPageSource &pageSource, unsigned int clusterBunchSize)
void StartBackgroundThread()
Spawn the I/O background thread. No-op if already started.
ROOT::Internal::RPageSource & fPageSource
Every cluster pool is responsible for exactly one page source that triggers loading of the clusters (...
An in-memory subset of the packed and compressed pages of a cluster.
std::unordered_set< ROOT::DescriptorId_t > ColumnSet_t
Abstract interface to read data from an ntuple.
void Erase(const T &that, std::vector< T > &v)
Erase that element from vector v
std::uint64_t DescriptorId_t
Distriniguishes elements of the same type within a descriptor, e.g. different fields.
constexpr NTupleSize_t kInvalidNTupleIndex
constexpr DescriptorId_t kInvalidDescriptorId
Performance counters that get registered in fMetrics.
Clusters that are currently being processed by the pipeline.
bool operator<(const RInFlightCluster &other) const
First order by cluster id, then by number of columns, than by the column ids in fColumns.
RCluster::RKey fClusterKey
std::future< std::unique_ptr< RCluster > > fFuture
Request to load a subset of the columns of a particular cluster.
RCluster::RKey fClusterKey
std::int64_t fBunchId
Items with different bunch ids are scheduled for different vector reads.
std::promise< std::unique_ptr< RCluster > > fPromise
ROOT::DescriptorId_t fClusterId
ColumnSet_t fPhysicalColumnSet