Project
Loading...
Searching...
No Matches
DataInputDirector.cxx
Go to the documentation of this file.
1// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
2// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
3// All rights not expressly granted are reserved.
4//
5// This software is distributed under the terms of the GNU General Public
6// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
7//
8// In applying this license CERN does not waive the privileges and immunities
9// granted to it by virtue of its status as an Intergovernmental Organization
10// or submit itself to any jurisdiction.
11#include "DataInputDirector.h"
13#include "Framework/Logger.h"
17#include "Framework/Output.h"
18#include "Framework/Signpost.h"
20#include "Headers/DataHeader.h"
21#include "Monitoring/Tags.h"
22#include "Monitoring/Metric.h"
23#include "Monitoring/Monitoring.h"
24
25#include "rapidjson/document.h"
26#include "rapidjson/prettywriter.h"
27#include "rapidjson/filereadstream.h"
28
29#include "TGrid.h"
30#include "TObjString.h"
31#include "TMap.h"
32#include "TFile.h"
33
34#include <arrow/dataset/file_base.h>
35#include <arrow/dataset/dataset.h>
36#include <uv.h>
37#include <exception>
38#include <memory>
39
40#if __has_include(<TJAlienFile.h>)
41#include <TJAlienFile.h>
42
43#include <utility>
44#endif
45
46#include <dlfcn.h>
47O2_DECLARE_DYNAMIC_LOG(reader_memory_dump);
48
49namespace o2::framework
50{
51using namespace rapidjson;
52
54{
55
56 FileNameHolder holder;
57 holder.fileName = fileName;
58 return holder;
59}
60
62 : mAlienSupport(alienSupport),
63 mContext(context),
64 mLevel(level)
65{
66 std::vector<char const*> capabilitiesSpecs = {
67 "O2Framework:RNTupleObjectReadingCapability",
68 "O2Framework:TTreeObjectReadingCapability",
69 };
70
71 std::vector<LoadablePlugin> plugins;
72 for (auto spec : capabilitiesSpecs) {
73 auto morePlugins = PluginManager::parsePluginSpecString(spec);
74 for (auto& extra : morePlugins) {
75 plugins.push_back(extra);
76 }
77 }
78
79 PluginManager::loadFromPlugin<RootObjectReadingCapability, RootObjectReadingCapabilityPlugin>(plugins, mFactory.capabilities);
80}
81
83{
84 LOGP(info, "DataInputDescriptor");
85 LOGP(info, " Table name : {}", tablename);
86 LOGP(info, " Tree name : {}", treename);
87 LOGP(info, " Input files file : {}", getInputfilesFilename());
88 LOGP(info, " File name regex : {}", getFilenamesRegexString());
89 LOGP(info, " Input files : {}", mfilenames.size());
90 for (auto& fn : mfilenames) {
91 LOGP(info, " {} {}", fn.fileName, fn.numberOfTimeFrames);
92 }
93 LOGP(info, " Total number of TF: {}", getNumberTimeFrames());
94}
95
97{
98 return (minputfilesFile.empty() && minputfilesFilePtr) ? (std::string)*minputfilesFilePtr : minputfilesFile;
99}
100
102{
103 return (mFilenameRegex.empty() && mFilenameRegexPtr) ? (std::string)*mFilenameRegexPtr : mFilenameRegex;
104}
105
107{
108 return std::regex(getFilenamesRegexString());
109}
110
112{
113 // remove leading file:// from file name
114 if (fn.fileName.rfind("file://", 0) == 0) {
115 fn.fileName.erase(0, 7);
116 } else if (!mAlienSupport && fn.fileName.rfind("alien://", 0) == 0 && !gGrid) {
117 LOGP(debug, "AliEn file requested. Enabling support.");
118 TGrid::Connect("alien://");
119 mAlienSupport = true;
120 }
121
122 mtotalNumberTimeFrames += fn.numberOfTimeFrames;
123 mfilenames.emplace_back(fn);
124}
125
126bool DataInputDescriptor::setFile(int counter, int wantedParentLevel, std::string_view origin)
127{
128 // no files left
129 if (counter >= getNumberInputfiles()) {
130 return false;
131 }
132
133 // In case the origin starts with a anything but AOD, we add the origin as the suffix
134 // of the filename. In the future we might expand this for proper rewriting of the
135 // filename based on the origin and the original file information.
136 std::string filename = mfilenames[counter].fileName;
137 // In case we do not need to remap parent levels, the requested origin is what
138 // drives the filename.
139 if (wantedParentLevel == -1 && !origin.starts_with("AOD")) {
140 filename = std::regex_replace(filename, std::regex("[.]root$"), fmt::format("_{}.root", origin));
141 }
142
143 // open file
144 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
145 if (rootFS.get()) {
146 if (rootFS->GetFile()->GetName() == filename) {
147 return true;
148 }
150 }
151
152 TFile* tfile = nullptr;
153 bool externalFile = false;
154 for (auto& [name, f] : mContext.openFiles) {
155 if (name == filename) {
156 tfile = f;
157 externalFile = true;
158 break;
159 }
160 }
161 if (tfile == nullptr) {
162 tfile = TFile::Open(filename.c_str());
163 }
164 if (!tfile || tfile->IsZombie()) {
165 if (tfile && !externalFile) {
166 delete tfile;
167 }
168 throw std::runtime_error(fmt::format("Couldn't open file \"{}\"!", filename));
169 }
170 mCurrentFilesystem = std::make_shared<TFileFileSystem>(tfile, 50 * 1024 * 1024, mFactory, !externalFile);
171 rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
173
174 // get the parent file map if exists
175 mParentFileMap = (TMap*)rootFS->GetFile()->Get("parentFiles"); // folder name (DF_XXX) --> parent file (absolute path)
176 if (mParentFileMap && !mContext.parentFileReplacement.empty()) {
177 auto pos = mContext.parentFileReplacement.find(';');
178 if (pos == std::string::npos) {
179 throw std::runtime_error(fmt::format("Invalid syntax in aod-parent-base-path-replacement: \"{}\"", mContext.parentFileReplacement.c_str()));
180 }
181 auto from = mContext.parentFileReplacement.substr(0, pos);
182 auto to = mContext.parentFileReplacement.substr(pos + 1);
183
184 auto it = mParentFileMap->MakeIterator();
185 while (auto obj = it->Next()) {
186 auto objString = (TObjString*)mParentFileMap->GetValue(obj);
187 objString->String().ReplaceAll(from.c_str(), to.c_str());
188 }
189 delete it;
190 }
191
192 // get the directory names
193 if (mfilenames[counter].numberOfTimeFrames <= 0) {
194 const std::regex TFRegex = std::regex("/?DF_([0-9]+)(|-.*)$");
195 TList* keyList = rootFS->GetFile()->GetListOfKeys();
196 std::vector<std::string> finalList;
197
198 // extract TF numbers and sort accordingly
199 // We use an extra seen set to make sure we preserve the order in which
200 // we instert things in the final list and to make sure we do not have duplicates.
201 // Multiple folder numbers can happen if we use a flat structure /DF_<df>-<tablename>
202 std::unordered_set<size_t> seen;
203 for (auto key : *keyList) {
204 std::smatch matchResult;
205 std::string keyName = ((TObjString*)key)->GetString().Data();
206 bool match = std::regex_match(keyName, matchResult, TFRegex);
207 if (match) {
208 auto folderNumber = std::stoul(matchResult[1].str());
209 if (seen.find(folderNumber) == seen.end()) {
210 seen.insert(folderNumber);
211 mfilenames[counter].listOfTimeFrameNumbers.emplace_back(folderNumber);
212 }
213 }
214 }
215
216 if (mParentFileMap != nullptr) {
217 // If we have a parent map, we should not process in DF alphabetical order but according to parent file to avoid swapping between files
218 std::ranges::sort(mfilenames[counter].listOfTimeFrameNumbers,
219 [this](long const& l1, long const& l2) -> bool {
220 auto p1 = (TObjString*)this->mParentFileMap->GetValue(("DF_" + std::to_string(l1)).c_str());
221 auto p2 = (TObjString*)this->mParentFileMap->GetValue(("DF_" + std::to_string(l2)).c_str());
222 return p1->GetString().CompareTo(p2->GetString()) < 0;
223 });
224 } else {
225 std::sort(mfilenames[counter].listOfTimeFrameNumbers.begin(), mfilenames[counter].listOfTimeFrameNumbers.end());
226 }
227
228 mfilenames[counter].alreadyRead.resize(mfilenames[counter].alreadyRead.size() + mfilenames[counter].listOfTimeFrameNumbers.size(), false);
229 mfilenames[counter].numberOfTimeFrames = mfilenames[counter].listOfTimeFrameNumbers.size();
230 }
231
232 mCurrentFileID = counter;
233 mCurrentFileStartedAt = uv_hrtime();
234 mIOTime = 0;
235
236 return true;
237}
238
239uint64_t DataInputDescriptor::getTimeFrameNumber(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin)
240{
241
242 // open file
243 if (!setFile(counter, wantedParentLevel, wantedOrigin)) {
244 return 0ul;
245 }
246
247 // no TF left
248 if (mfilenames[counter].numberOfTimeFrames > 0 && numTF >= mfilenames[counter].numberOfTimeFrames) {
249 return 0ul;
250 }
251
252 return (mfilenames[counter].listOfTimeFrameNumbers)[numTF];
253}
254
255std::pair<std::shared_ptr<DataInputDescriptor>, int> DataInputDescriptor::navigateToLevel(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin)
256{
257 if (!setFile(counter, wantedParentLevel, wantedOrigin)) {
258 return {nullptr, -1};
259 }
260 if (numTF < 0 || numTF >= mfilenames[counter].numberOfTimeFrames) {
261 return {nullptr, -1};
262 }
263 recordTimeFrameRead(counter, numTF);
264 auto folderName = fmt::format("DF_{}", mfilenames[counter].listOfTimeFrameNumbers[numTF]);
265 auto parentFile = getParentFile(counter, numTF, "", wantedParentLevel, wantedOrigin);
266 if (parentFile == nullptr) {
267 return {nullptr, -1};
268 }
269 return {parentFile, parentFile->findDFNumber(0, folderName)};
270}
271
272arrow::dataset::FileSource DataInputDescriptor::getFileFolder(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin)
273{
274 // If mapped to a parent level deeper than current, skip directly to the right level.
275 if ((wantedParentLevel != -1) && (mLevel < wantedParentLevel)) {
276 auto [parentFile, parentNumTF] = navigateToLevel(counter, numTF, wantedParentLevel, wantedOrigin);
277 if (parentFile == nullptr || parentNumTF == -1) {
278 return {};
279 }
280 return parentFile->getFileFolder(0, parentNumTF, wantedParentLevel, wantedOrigin);
281 }
282
283 // open file
284 if (!setFile(counter, wantedParentLevel, wantedOrigin)) {
285 return {};
286 }
287
288 // no TF left
289 if ((mfilenames[counter].numberOfTimeFrames > 0) && (numTF >= mfilenames[counter].numberOfTimeFrames)) {
290 return {};
291 }
292
293 recordTimeFrameRead(counter, numTF);
294 mfilenames[counter].alreadyRead[numTF] = true;
295
296 return {fmt::format("DF_{}", mfilenames[counter].listOfTimeFrameNumbers[numTF]), mCurrentFilesystem};
297}
298
299void DataInputDescriptor::recordTimeFrameRead(int counter, int numTF)
300{
301 auto read = std::pair{counter, numTF};
302 if (std::find(mTimeFrameReads.begin(), mTimeFrameReads.end(), read) == mTimeFrameReads.end()) {
303 mTimeFrameReads.push_back(read);
304 }
305}
306
308{
309 mTimeFrameReads.clear();
310 mTimeFrameActive = true;
311 if (mParentFile) {
312 mParentFile->beginTimeFrame();
313 }
314}
315
317{
318 mTimeFrameActive = false;
319 if (skipped) {
320 for (auto [file, df] : mTimeFrameReads) {
321 mfilenames[file].alreadyRead[df] = false;
322 ++mfilenames[file].invalidReadSkipped;
323 }
324 }
325 mTimeFrameReads.clear();
326 if (mParentFile) {
327 mParentFile->finishTimeFrame(skipped);
328 }
329 for (auto& parent : mRetainedParents) {
330 parent->finishTimeFrame(skipped);
331 parent->closeInputFile();
332 }
333 mRetainedParents.clear();
334 for (auto& [counter, info] : mPendingFileStatistics) {
335 reportFileStatistics(counter, std::move(info));
336 }
337 mPendingFileStatistics.clear();
338}
339
340void DataInputDescriptor::releaseParentFile()
341{
342 if (!mParentFile) {
343 return;
344 }
345 if (mTimeFrameActive) {
346 // Keep both the read bookkeeping and the open file for final statistics.
347 mRetainedParents.push_back(std::move(mParentFile));
348 } else {
349 mParentFile->closeInputFile();
350 mParentFile.reset();
351 }
352}
353
354std::shared_ptr<DataInputDescriptor> DataInputDescriptor::getParentFile(int counter, int numTF, std::string treename, int wantedParentLevel, std::string_view wantedOrigin)
355{
356 if (!mParentFileMap) {
357 // This file has no parent map
358 return nullptr;
359 }
360
361 auto folderName = fmt::format("DF_{}", mfilenames[counter].listOfTimeFrameNumbers[numTF]);
362 auto parentFileName = (TObjString*)mParentFileMap->GetValue(folderName.c_str());
363 // The current DF is not found in the parent map (this should not happen and is a fatal error)
364 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
365 if (!parentFileName) {
366 throw InvalidAODReadError(fmt::format(R"(parent file map exists but does not contain the current DF "{}" in file "{}")", folderName.c_str(), rootFS->GetFile()->GetName()));
367 }
368
369 if (mParentFile) {
370 // Is this still the corresponding to the correct file?
371 auto parentRootFS = std::dynamic_pointer_cast<TFileFileSystem>(mParentFile->mCurrentFilesystem);
372 if (parentFileName->GetString().CompareTo(parentRootFS->GetFile()->GetName()) == 0) {
373 return mParentFile;
374 } else {
375 releaseParentFile();
376 }
377 }
378
379 if (mLevel == mContext.allowedParentLevel) {
380 throw std::runtime_error(fmt::format(R"(while looking for tree "{}", the parent file was requested but we are already at level {} of maximal allowed level {} for DF "{}" in file "{}")", treename.c_str(), mLevel, mContext.allowedParentLevel, folderName.c_str(),
381 rootFS->GetFile()->GetName()));
382 }
383
384 LOGP(info, "Opening parent file {} for DF {}", parentFileName->GetString().Data(), folderName.c_str());
385 mParentFile = std::make_shared<DataInputDescriptor>(mAlienSupport, mLevel + 1, mContext);
386 mParentFile->mdefaultFilenamesPtr.emplace_back(makeFileNameHolder(parentFileName->GetString().Data()));
387 mParentFile->fillInputfiles();
388 try {
389 mParentFile->setFile(0, wantedParentLevel, wantedOrigin);
390 if (mTimeFrameActive) {
391 mParentFile->beginTimeFrame();
392 }
393 } catch (...) {
394 mParentFile.reset();
395 std::throw_with_nested(InvalidAODReadError(fmt::format("Unable to open parent file \"{}\" for {}", parentFileName->GetString().Data(), folderName)));
396 }
397 return mParentFile;
398}
399
401{
402 return mfilenames.at(counter).numberOfTimeFrames;
403}
404
406{
407 auto& list = mfilenames.at(counter).alreadyRead;
408 return std::count(list.begin(), list.end(), true);
409}
410
412{
413 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
414 auto f = dynamic_cast<TFile*>(rootFS->GetFile());
415 std::string monitoringInfo(fmt::format("lfn={},size={}", f->GetName(), f->GetSize()));
416#if __has_include(<TJAlienFile.h>)
417 auto alienFile = dynamic_cast<TJAlienFile*>(f);
418 if (alienFile) {
419 monitoringInfo += fmt::format(",se={},open_time={:.1f}", alienFile->GetSE(), alienFile->GetElapsed());
420 }
421#endif
422 if (mContext.monitoring) {
423 mContext.monitoring->send(o2::monitoring::Metric{monitoringInfo, "aod-file-open-info"}.addTag(o2::monitoring::tags::Key::Subsystem, o2::monitoring::tags::Value::DPL));
424 }
425 LOGP(info, "Opening file: {}", monitoringInfo);
426}
427
429{
430 int64_t wait_time = (int64_t)uv_hrtime() - (int64_t)mCurrentFileStartedAt - (int64_t)mIOTime;
431 if (wait_time < 0) {
432 wait_time = 0;
433 }
434 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
435 auto f = dynamic_cast<TFile*>(rootFS->GetFile());
436 std::string monitoringInfo(fmt::format("lfn={},size={},total_df={},read_bytes={},read_calls={},io_time={:.1f},wait_time={:.1f},level={}", f->GetName(),
437 f->GetSize(), getTimeFramesInFile(mCurrentFileID), f->GetBytesRead(), f->GetReadCalls(),
438 ((float)mIOTime / 1e9), ((float)wait_time / 1e9), mLevel));
439#if __has_include(<TJAlienFile.h>)
440 auto alienFile = dynamic_cast<TJAlienFile*>(f);
441 if (alienFile) {
442 monitoringInfo += fmt::format(",se={},open_time={:.1f}", alienFile->GetSE(), alienFile->GetElapsed());
443 }
444#endif
445 if (mTimeFrameActive) {
446 // Snapshot I/O statistics now, but publish DF counts after commit or rollback.
447 mPendingFileStatistics.emplace_back(mCurrentFileID, std::move(monitoringInfo));
448 } else {
449 reportFileStatistics(mCurrentFileID, std::move(monitoringInfo));
450 }
451}
452
453void DataInputDescriptor::reportFileStatistics(int counter, std::string monitoringInfo)
454{
455 monitoringInfo += fmt::format(",read_df={},skipped_df={}", getReadTimeFramesInFile(counter), mfilenames.at(counter).invalidReadSkipped);
456 if (mContext.monitoring) {
457 mContext.monitoring->send(o2::monitoring::Metric{monitoringInfo, "aod-file-read-info"}.addTag(o2::monitoring::tags::Key::Subsystem, o2::monitoring::tags::Value::DPL));
458 }
459 LOGP(info, "Read info: {}", monitoringInfo);
460}
461
463{
464 if (mCurrentFilesystem.get()) {
465 releaseParentFile();
466
467 delete mParentFileMap;
468 mParentFileMap = nullptr;
469
471 mCurrentFilesystem.reset();
472 }
473}
474
476{
477 if (getNumberInputfiles() > 0) {
478 // 1. mfilenames
479 return getNumberInputfiles();
480 }
481
482 auto fileName = getInputfilesFilename();
483 if (!fileName.empty()) {
484 // 2. getFilenamesRegex() @ getInputfilesFilename()
485 try {
486 std::ifstream filelist(fileName);
487 if (!filelist.is_open()) {
488 throw std::runtime_error(fmt::format(R"(Couldn't open file "{}")", fileName));
489 }
490 while (std::getline(filelist, fileName)) {
491 // remove white spaces, empty lines are skipped
492 fileName.erase(std::remove_if(fileName.begin(), fileName.end(), ::isspace), fileName.end());
493 if (!fileName.empty() && (getFilenamesRegexString().empty() ||
494 std::regex_match(fileName, getFilenamesRegex()))) {
496 }
497 }
498 } catch (...) {
499 LOGP(error, "Check the input files file! Unable to process \"{}\"!", getInputfilesFilename());
500 return 0;
501 }
502 } else {
503 // 3. getFilenamesRegex() @ mdefaultFilenamesPtr
504 if (!mdefaultFilenamesPtr.empty()) {
505 for (auto& fileNameHolder : mdefaultFilenamesPtr) {
506 if (getFilenamesRegexString().empty() ||
507 std::regex_match(fileNameHolder.fileName, getFilenamesRegex())) {
508 addFileNameHolder(fileNameHolder);
509 }
510 }
511 }
512 }
513
514 return getNumberInputfiles();
515}
516
517int DataInputDescriptor::findDFNumber(int file, std::string dfName)
518{
519 auto dfList = mfilenames[file].listOfTimeFrameNumbers;
520 auto it = std::find_if(dfList.begin(), dfList.end(), [dfName](size_t i) { return fmt::format("DF_{}", i) == dfName; });
521 if (it == dfList.end()) {
522 return -1;
523 }
524 return it - dfList.begin();
525}
526
529 : mTarget(target)
530 {
531 start = uv_hrtime();
532 }
534 {
535 if (!active) {
536 return;
537 }
538 O2_SIGNPOST_ACTION(reader_memory_dump, [](void*) {
539 void (*dump_)(const char*);
540 if (void* sym = dlsym(nullptr, "igprof_dump_now")) {
541 dump_ = __extension__(void (*)(const char*)) sym;
542 if (dump_) {
543 std::string filename = fmt::format("reader-memory-dump-{}.gz", uv_hrtime());
544 dump_(filename.c_str());
545 }
546 }
547 });
548 mTarget += (uv_hrtime() - start);
549 }
550
552 {
553 active = false;
554 }
555
556 bool active = true;
557 uint64_t& mTarget;
558 uint64_t start;
559 uint64_t stop;
560};
561
562bool DataInputDescriptor::readTree(DataAllocator& outputs, header::DataHeader dh, int counter, int numTF, std::string treename, size_t& totalSizeCompressed, size_t& totalSizeUncompressed)
563try {
564 CalculateDelta t(mIOTime);
565 std::string wantedOrigin = dh.dataOrigin.as<std::string>();
566 int wantedLevel = mContext.levelForOrigin(wantedOrigin);
567
568 // If this origin is mapped to a parent level deeper than current, skip directly without
569 // attempting to read from this level.
570 if (wantedLevel != -1 && mLevel < wantedLevel) {
571 auto [parentFile, parentNumTF] = navigateToLevel(counter, numTF, wantedLevel, wantedOrigin);
572 if (counter >= getNumberInputfiles() || numTF < 0 || numTF >= mfilenames[counter].numberOfTimeFrames) {
573 t.deactivate();
574 return false;
575 }
576 if (parentFile == nullptr) {
577 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
578 throw InvalidAODReadError(fmt::format(R"(No parent file found for "{}" while looking for level {} in "{}")", treename, wantedLevel, rootFS->GetFile()->GetName()));
579 }
580 if (parentNumTF == -1) {
581 auto parentRootFS = std::dynamic_pointer_cast<TFileFileSystem>(parentFile->mCurrentFilesystem);
582 throw InvalidAODReadError(fmt::format(R"(DF not found in parent file "{}")", parentRootFS->GetFile()->GetName()));
583 }
584 t.deactivate();
585 return parentFile->readTree(outputs, dh, 0, parentNumTF, treename, totalSizeCompressed, totalSizeUncompressed);
586 }
587
588 auto folder = getFileFolder(counter, numTF, wantedLevel, wantedOrigin);
589 if (!folder.filesystem()) {
590 t.deactivate();
591 return false;
592 }
593
594 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(folder.filesystem());
595
596 if (!rootFS) {
597 t.deactivate();
598 throw std::runtime_error(fmt::format(R"(Not a TFile filesystem!)"));
599 }
600 // FIXME: Ugly. We should detect the format from the treename, good enough for now.
601 std::shared_ptr<arrow::dataset::FileFormat> format;
602 FragmentToBatch::StreamerCreator creator = nullptr;
603
604 auto fullpath = arrow::dataset::FileSource{folder.path() + "/" + treename, folder.filesystem()};
605
606 for (auto& capability : mFactory.capabilities) {
607 auto objectPath = capability.lfn2objectPath(fullpath.path());
608 void* handle = capability.getHandle(rootFS, objectPath);
609 if (handle) {
610 format = capability.factory().format();
611 creator = capability.factory().deferredOutputStreamer;
612 // Account for the bytes we are about to read. This used to sit further down, where
613 // the TTree was opened by hand; moving the reading to the arrow::Dataset API left
614 // the accounting behind, which is why aod-bytes-read-* and the --aod-max-read-rate
615 // pacing that derives from them both read zero. Each format reports its own size,
616 // so we just ask; here is where the object is resolved and its size is known.
617 if (capability.accountBytes) {
618 capability.accountBytes(handle, totalSizeCompressed, totalSizeUncompressed);
619 }
620 break;
621 }
622 }
623
624 // FIXME: we should distinguish between an actually missing object and one which has a non compatible
625 // format.
626 if (!format) {
627 t.deactivate();
628 LOGP(debug, "Could not find tree {}. Trying in parent file.", fullpath.path());
629 auto parentFile = getParentFile(counter, numTF, treename, wantedLevel, wantedOrigin);
630 if (parentFile == nullptr) {
631 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
632 throw std::runtime_error(fmt::format(R"(Couldn't get TTree "{}" from "{}". Please check https://aliceo2group.github.io/analysis-framework/docs/troubleshooting/#tree-not-found for more information.)", fullpath.path(), rootFS->GetFile()->GetName()));
633 }
634 int parentNumTF = parentFile->findDFNumber(0, folder.path());
635 if (parentNumTF == -1) {
636 auto parentRootFS = std::dynamic_pointer_cast<TFileFileSystem>(parentFile->mCurrentFilesystem);
637 throw InvalidAODReadError(fmt::format(R"(DF {} listed in parent file map but not found in the corresponding file "{}")", folder.path(), parentRootFS->GetFile()->GetName()));
638 }
639 // first argument is 0 as the parent file object contains only 1 file
640 return parentFile->readTree(outputs, dh, 0, parentNumTF, treename, totalSizeCompressed, totalSizeUncompressed);
641 }
642
643 auto schemaOpt = format->Inspect(fullpath);
644 if (!schemaOpt.ok()) {
645 throw InvalidAODReadError(fmt::format("Unable to inspect tree {}: {}", treename, schemaOpt.status().ToString()));
646 }
647 auto physicalSchema = schemaOpt;
648 std::vector<std::shared_ptr<arrow::Field>> fields;
649 for (auto& original : (*schemaOpt)->fields()) {
650 if (original->name().ends_with("_size")) {
651 continue;
652 }
653 fields.push_back(original);
654 }
655 auto datasetSchema = std::make_shared<arrow::Schema>(fields);
656
657 auto fragment = format->MakeFragment(fullpath, {}, *physicalSchema);
658
659 // create table output
660 auto o = Output(dh);
661
662 // FIXME: This should allow me to create a memory pool
663 // which I can then use to scan the dataset.
664 auto f2b = outputs.make<FragmentToBatch>(o, creator, *fragment);
665
668 f2b->setLabel(treename.c_str());
669 char const* operation = "read";
670 try {
671 f2b->fill(datasetSchema, format);
672 operation = "finalize";
673 f2b.release();
674 } catch (...) {
675 f2b.discard();
676 std::throw_with_nested(InvalidAODReadError(fmt::format("Unable to {} tree {}", operation, treename)));
677 }
678
679 return true;
680} catch (InvalidAODReadError const&) {
681 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
682 auto filename = rootFS ? std::string(rootFS->GetFile()->GetName()) : mfilenames.at(counter).fileName;
683 auto const& dfs = mfilenames.at(counter).listOfTimeFrameNumbers;
684 auto df = numTF >= 0 && static_cast<size_t>(numTF) < dfs.size() ? fmt::format("DF_{}", dfs[numTF]) : fmt::format("timeframe index {}", numTF);
685 std::throw_with_nested(InvalidAODReadError(fmt::format("Reading tree {} in {} from file \"{}\" at parent level {}", treename, df, filename, mLevel)));
686}
687
688DataInputDirector::DataInputDirector(std::vector<std::string> inputFiles, DataInputDirectorContext&& context)
689 : mContext{context}
690{
691 if (inputFiles.size() == 1 && !inputFiles[0].empty() && inputFiles[0][0] == '@') {
692 setInputfilesFile(inputFiles.back().substr(1, -1));
693 } else {
694 for (auto inputFile : inputFiles) {
695 mdefaultInputFiles.emplace_back(makeFileNameHolder(inputFile));
696 }
697 }
698
700}
701
703{
704 mdefaultInputFiles.clear();
705 mdefaultDataInputDescriptor = nullptr;
706
707 mdataInputDescriptors.clear();
708}
709
711{
712 mdataInputDescriptors.clear();
713 mdefaultInputFiles.clear();
714 mFilenameRegex = std::string("");
715};
716
718{
719 if (mdefaultDataInputDescriptor) {
720 mdefaultDataInputDescriptor.reset();
721 }
722 mdefaultDataInputDescriptor = std::make_shared<DataInputDescriptor>(mAlienSupport, 0, mContext);
723
724 mdefaultDataInputDescriptor->setInputfilesFile(minputfilesFile);
725 mdefaultDataInputDescriptor->setFilenamesRegex(mFilenameRegex);
726 mdefaultDataInputDescriptor->setDefaultInputfiles(mdefaultInputFiles);
727 mdefaultDataInputDescriptor->tablename = "any";
728 mdefaultDataInputDescriptor->treename = "any";
729 mdefaultDataInputDescriptor->fillInputfiles();
730
731 mAlienSupport &= mdefaultDataInputDescriptor->isAlienSupportOn();
732}
733
734bool DataInputDirector::readJson(std::string const& fnjson)
735{
736 // open the file
737 FILE* f = fopen(fnjson.c_str(), "r");
738 if (!f) {
739 LOGP(error, "Could not open JSON file \"{}\"!", fnjson);
740 return false;
741 }
742
743 // create streamer
744 char readBuffer[65536];
745 FileReadStream inputStream(f, readBuffer, sizeof(readBuffer));
746
747 // parse the json file
748 Document jsonDoc;
749 jsonDoc.ParseStream(inputStream);
750 auto status = readJsonDocument(&jsonDoc);
751
752 // clean up
753 fclose(f);
754
755 return status;
756}
757
758bool DataInputDirector::readJsonDocument(Document* jsonDoc)
759{
760 // initialisations
761 std::string fileName("");
762 const char* itemName;
763
764 // is it a proper json document?
765 if (jsonDoc->HasParseError()) {
766 LOGP(error, "Check the JSON document! There is a problem with the format!");
767 return false;
768 }
769
770 // InputDirector
771 itemName = "InputDirector";
772 const Value& didirItem = (*jsonDoc)[itemName];
773 if (!didirItem.IsObject()) {
774 LOGP(info, "No \"{}\" object found in the JSON document!", itemName);
775 return true;
776 }
777
778 // now read various items
779 itemName = "debugmode";
780 if (didirItem.HasMember(itemName)) {
781 if (didirItem[itemName].IsBool()) {
782 mDebugMode = (didirItem[itemName].GetBool());
783 } else {
784 LOGP(error, "Check the JSON document! Item \"{}\" must be a boolean!", itemName);
785 return false;
786 }
787 } else {
788 mDebugMode = false;
789 }
790
791 if (mDebugMode) {
792 StringBuffer buffer;
793 buffer.Clear();
794 PrettyWriter<StringBuffer> writer(buffer);
795 didirItem.Accept(writer);
796 LOGP(info, "InputDirector object: {}", std::string(buffer.GetString()));
797 }
798
799 itemName = "fileregex";
800 if (didirItem.HasMember(itemName)) {
801 if (didirItem[itemName].IsString()) {
802 setFilenamesRegex(didirItem[itemName].GetString());
803 } else {
804 LOGP(error, "Check the JSON document! Item \"{}\" must be a string!", itemName);
805 return false;
806 }
807 }
808
809 itemName = "resfiles";
810 if (didirItem.HasMember(itemName)) {
811 if (didirItem[itemName].IsString()) {
812 fileName = didirItem[itemName].GetString();
813 if (fileName.size() && fileName[0] == '@') {
814 fileName.erase(0, 1);
815 setInputfilesFile(fileName);
816 } else {
818 mdefaultInputFiles.emplace_back(makeFileNameHolder(fileName));
819 }
820 } else if (didirItem[itemName].IsArray()) {
822 auto fns = didirItem[itemName].GetArray();
823 for (auto& fn : fns) {
824 mdefaultInputFiles.emplace_back(makeFileNameHolder(fn.GetString()));
825 }
826 } else {
827 LOGP(error, "Check the JSON document! Item \"{}\" must be a string or an array!", itemName);
828 return false;
829 }
830 }
831
832 itemName = "InputDescriptors";
833 if (didirItem.HasMember(itemName)) {
834 if (!didirItem[itemName].IsArray()) {
835 LOGP(error, "Check the JSON document! Item \"{}\" must be an array!", itemName);
836 return false;
837 }
838
839 // loop over DataInputDescriptors
840 for (auto& didescItem : didirItem[itemName].GetArray()) {
841 if (!didescItem.IsObject()) {
842 LOGP(error, "Check the JSON document! \"{}\" must be objects!", itemName);
843 return false;
844 }
845 // create a new dataInputDescriptor
846 auto didesc = DataInputDescriptor(mAlienSupport, 0, mContext);
847 didesc.setDefaultInputfiles(mdefaultInputFiles);
848
849 itemName = "table";
850 if (didescItem.HasMember(itemName)) {
851 if (didescItem[itemName].IsString()) {
852 didesc.tablename = didescItem[itemName].GetString();
853 didesc.matcher = DataDescriptorQueryBuilder::buildNode(didesc.tablename);
854 } else {
855 LOGP(error, "Check the JSON document! Item \"{}\" must be a string!", itemName);
856 return false;
857 }
858 } else {
859 LOGP(error, "Check the JSON document! Item \"{}\" is missing!", itemName);
860 return false;
861 }
862
863 itemName = "treename";
864 if (didescItem.HasMember(itemName)) {
865 if (didescItem[itemName].IsString()) {
866 didesc.treename = didescItem[itemName].GetString();
867 } else {
868 LOGP(error, "Check the JSON document! Item \"{}\" must be a string!", itemName);
869 return false;
870 }
871 } else {
872 auto m = DataDescriptorQueryBuilder::getTokens(didesc.tablename);
873 didesc.treename = m[2];
874 }
875
876 itemName = "fileregex";
877 if (didescItem.HasMember(itemName)) {
878 if (didescItem[itemName].IsString()) {
879 if (didesc.getNumberInputfiles() == 0) {
880 didesc.setFilenamesRegex(didescItem[itemName].GetString());
881 }
882 } else {
883 LOGP(error, "Check the JSON document! Item \"{}\" must be a string!", itemName);
884 return false;
885 }
886 } else {
887 if (didesc.getNumberInputfiles() == 0) {
888 didesc.setFilenamesRegex(mFilenameRegexPtr);
889 }
890 }
891
892 itemName = "resfiles";
893 if (didescItem.HasMember(itemName)) {
894 if (didescItem[itemName].IsString()) {
895 fileName = didescItem[itemName].GetString();
896 if (fileName.size() && fileName[0] == '@') {
897 didesc.setInputfilesFile(fileName.erase(0, 1));
898 } else {
899 if (didesc.getFilenamesRegexString().empty() ||
900 std::regex_match(fileName, didesc.getFilenamesRegex())) {
901 didesc.addFileNameHolder(makeFileNameHolder(fileName));
902 }
903 }
904 } else if (didescItem[itemName].IsArray()) {
905 auto fns = didescItem[itemName].GetArray();
906 for (auto& fn : fns) {
907 if (didesc.getFilenamesRegexString().empty() ||
908 std::regex_match(fn.GetString(), didesc.getFilenamesRegex())) {
909 didesc.addFileNameHolder(makeFileNameHolder(fn.GetString()));
910 }
911 }
912 } else {
913 LOGP(error, "Check the JSON document! Item \"{}\" must be a string or an array!", itemName);
914 return false;
915 }
916 } else {
917 didesc.setInputfilesFile(minputfilesFilePtr);
918 }
919
920 // fill mfilenames and add InputDescriptor to InputDirector
921 if (didesc.fillInputfiles() > 0) {
922 mdataInputDescriptors.emplace_back(didesc);
923 } else {
924 didesc.printOut();
925 LOGP(info, "This DataInputDescriptor is ignored because its file list is empty!");
926 }
927 mAlienSupport &= didesc.isAlienSupportOn();
928 }
929 }
930
931 // add a default DataInputDescriptor
933
934 // check that all DataInputDescriptors have the same number of input files
935 if (!isValid()) {
936 printOut();
937 return false;
938 }
939
940 // print the DataIputDirector
941 if (mDebugMode) {
942 printOut();
943 }
944
945 return true;
946}
947
949{
950 // compute list of matching outputs
952
953 for (auto& didesc : mdataInputDescriptors) {
954 if (didesc.matcher->match(dh, context)) {
955 return &didesc;
956 }
957 }
958
959 return nullptr;
960}
961
962arrow::dataset::FileSource DataInputDirector::getFileFolder(header::DataHeader dh, int counter, int numTF)
963{
964 auto didesc = getDataInputDescriptor(dh);
965 // if NOT match then use defaultDataInputDescriptor
966 if (!didesc) {
967 didesc = mdefaultDataInputDescriptor.get();
968 }
969 std::string origin = dh.dataOrigin.as<std::string>();
970 int wantedLevel = mContext.levelForOrigin(origin);
971
972 return didesc->getFileFolder(counter, numTF, wantedLevel, origin);
973}
974
976{
977 mdefaultDataInputDescriptor->beginTimeFrame();
978 for (auto& descriptor : mdataInputDescriptors) {
979 descriptor.beginTimeFrame();
980 }
981}
982
984{
985 mdefaultDataInputDescriptor->finishTimeFrame(skipped);
986 for (auto& descriptor : mdataInputDescriptors) {
987 descriptor.finishTimeFrame(skipped);
988 }
989}
990
992{
993 auto didesc = getDataInputDescriptor(dh);
994 // if NOT match then use defaultDataInputDescriptor
995 if (!didesc) {
996 didesc = mdefaultDataInputDescriptor.get();
997 }
998
999 return didesc->getTimeFramesInFile(counter);
1000}
1001
1003{
1004 auto didesc = getDataInputDescriptor(dh);
1005 // if NOT match then use defaultDataInputDescriptor
1006 if (!didesc) {
1007 didesc = mdefaultDataInputDescriptor.get();
1008 }
1009 std::string origin = dh.dataOrigin.as<std::string>();
1010 int wantedLevel = mContext.levelForOrigin(origin);
1011
1012 return didesc->getTimeFrameNumber(counter, numTF, wantedLevel, origin);
1013}
1014
1015bool DataInputDirector::readTree(DataAllocator& outputs, header::DataHeader dh, int counter, int numTF, size_t& totalSizeCompressed, size_t& totalSizeUncompressed, bool wasAOD)
1016{
1017 std::string treename;
1018
1019 auto didesc = getDataInputDescriptor(dh);
1020 if (didesc) {
1021 // if match then use filename and treename from DataInputDescriptor
1022 treename = didesc->treename;
1023 } else {
1024 // if NOT match then use
1025 // . filename from defaultDataInputDescriptor
1026 // . treename from DataHeader
1027 didesc = mdefaultDataInputDescriptor.get();
1028 treename = aod::datamodel::getTreeName(dh, wasAOD);
1029 }
1030 std::string origin = dh.dataOrigin.as<std::string>();
1031
1032 auto result = didesc->readTree(outputs, dh, counter, numTF, treename, totalSizeCompressed, totalSizeUncompressed);
1033 return result;
1034}
1035
1037{
1039 mdefaultDataInputDescriptor->closeInputFile();
1040 for (auto& didesc : mdataInputDescriptors) {
1041 didesc.closeInputFile();
1042 }
1043}
1044
1045bool DataInputDirector::isValid()
1046{
1047 bool status = true;
1048 int numberFiles = mdefaultDataInputDescriptor->getNumberInputfiles();
1049 for (auto& didesc : mdataInputDescriptors) {
1050 status &= didesc.getNumberInputfiles() == numberFiles;
1051 }
1052
1053 return status;
1054}
1055
1057{
1058 bool status = mdefaultDataInputDescriptor->getNumberInputfiles() <= counter;
1059 for (auto& didesc : mdataInputDescriptors) {
1060 status &= (didesc.getNumberInputfiles() <= counter);
1061 }
1062
1063 return status;
1064}
1065
1067{
1068 LOGP(info, "DataInputDirector");
1069 LOGP(info, " Default input files file : {}", minputfilesFile);
1070 LOGP(info, " Default file name regex : {}", mFilenameRegex);
1071 LOGP(info, " Default file names : {}", mdefaultInputFiles.size());
1072 for (auto const& fn : mdefaultInputFiles) {
1073 LOGP(info, " {} {}", fn.fileName, fn.numberOfTimeFrames);
1074 }
1075 LOGP(info, " Default DataInputDescriptor:");
1076 mdefaultDataInputDescriptor->printOut();
1077 LOGP(info, " DataInputDescriptors : {}", getNumberInputDescriptors());
1078 for (auto const& didesc : mdataInputDescriptors) {
1079 didesc.printOut();
1080 }
1081}
1082
1084{
1085 return mContext.levelForOrigin(origin.as<std::string>());
1086}
1087
1088} // namespace o2::framework
header::DataOrigin origin
o2::monitoring::tags::Value Value
BasicOp operation
std::vector< std::shared_ptr< arrow::Field > > fields
std::ostringstream debug
int32_t i
constexpr int p2()
constexpr int p1()
constexpr to accelerate the coordinates changing
uint16_t pos
Definition RawData.h:3
#define O2_DECLARE_DYNAMIC_LOG(name)
Definition Signpost.h:490
#define O2_SIGNPOST_ACTION(log, callback)
Definition Signpost.h:511
StringRef key
decltype(auto) make(const Output &spec, Args... args)
uint64_t getTimeFrameNumber(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin)
std::shared_ptr< DataInputDescriptor > getParentFile(int counter, int numTF, std::string treename, int wantedParentLevel, std::string_view wantedOrigin)
bool readTree(DataAllocator &outputs, header::DataHeader dh, int counter, int numTF, std::string treename, size_t &totalSizeCompressed, size_t &totalSizeUncompressed)
arrow::dataset::FileSource getFileFolder(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin)
DataInputDescriptor(bool alienSupport, int level, DataInputDirectorContext &context)
void addFileNameHolder(FileNameHolder fn)
void finishTimeFrame(bool skipped=false)
bool setFile(int counter, int wantedParentLevel, std::string_view wantedOrigin)
int findDFNumber(int file, std::string dfName)
std::pair< std::shared_ptr< DataInputDescriptor >, int > navigateToLevel(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin)
DataInputDescriptor * getDataInputDescriptor(header::DataHeader dh)
DataInputDirector(std::vector< std::string > inputFiles, DataInputDirectorContext &&context)
arrow::dataset::FileSource getFileFolder(header::DataHeader dh, int counter, int numTF)
void setInputfilesFile(std::string iffn)
int getTimeFramesInFile(header::DataHeader dh, int counter)
uint64_t getTimeFrameNumber(header::DataHeader dh, int counter, int numTF)
int getLevelForOrigin(header::DataOrigin origin) const
void setFilenamesRegex(std::string dfn)
void finishTimeFrame(bool skipped=false)
bool readJson(std::string const &fnjson)
bool readTree(DataAllocator &outputs, header::DataHeader dh, int counter, int numTF, size_t &totalSizeCompressed, size_t &totalSizeUncompressed, bool wasAOD)
std::function< std::shared_ptr< arrow::io::OutputStream >(std::shared_ptr< arrow::dataset::FileFragment >, const std::shared_ptr< arrow::ResizableBuffer > &buffer)> StreamerCreator
void setLabel(const char *label)
bool match(const std::vector< std::string > &queries, const char *pattern)
Definition dcs-ccdb.cxx:229
const GLfloat * m
Definition glcorearb.h:4066
GLuint64EXT * result
Definition glcorearb.h:5662
GLuint buffer
Definition glcorearb.h:655
GLuint const GLchar * name
Definition glcorearb.h:781
GLdouble f
Definition glcorearb.h:310
GLenum target
Definition glcorearb.h:1641
typedef void(APIENTRYP PFNGLCULLFACEPROC)(GLenum mode)
GLint level
Definition glcorearb.h:275
GLuint start
Definition glcorearb.h:469
GLint GLint GLsizei GLint GLenum format
Definition glcorearb.h:275
GLuint counter
Definition glcorearb.h:3987
std::string getTreeName(header::DataHeader dh, bool wasAOD)
Defining ITS Vertex explicitly as messageable.
Definition Cartesian.h:288
FileNameHolder makeFileNameHolder(std::string fileName)
std::string to_string(gsl::span< T, Size > span)
Definition common.h:52
std::string filename()
void empty(int)
static std::unique_ptr< data_matcher::DataDescriptorMatcher > buildNode(std::string const &nodeString)
static std::vector< std::string > getTokens(std::string const &nodeString)
o2::monitoring::Monitoring * monitoring
std::vector< std::pair< std::string, TFile * > > openFiles
int levelForOrigin(std::string_view origin) const
static std::vector< LoadablePlugin > parsePluginSpecString(char const *str)
Parse a comma separated list of <library>:<plugin-name> plugin declarations.
std::vector< RootObjectReadingCapability > capabilities
the main header struct
Definition DataHeader.h:620
std::enable_if_t< std::is_same< T, std::string >::value==true, T > as() const
get the descriptor as std::string
Definition DataHeader.h:301
const std::string str