Skip to content

Commit 522a524

Browse files
authored
Track skipped read timeframes in file statistics & dumping (#15735)
* Track skipped read timeframes in file statistics, dump * Use concrete when recording skipped timeframe
1 parent 6747eef commit 522a524

4 files changed

Lines changed: 28 additions & 5 deletions

File tree

Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -271,10 +271,11 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const
271271

272272
int64_t startTime = uv_hrtime();
273273
int64_t startSize = totalSizeCompressed;
274-
auto skipInvalidRead = [&](o2::header::DataOrigin const& origin, InvalidAODReadError const& e) {
274+
auto skipInvalidRead = [&](ConcreteDataMatcher const& concrete, InvalidAODReadError const& e) {
275275
auto skippedTimeframes = ++totalInvalidReadSkipped;
276276
LOGP(error, "Invalid AOD read for table {}: fileCounter {}, timeFrame {}. Skipping timeframe (skipped timeframes: {}). Reason: {}",
277-
origin.as<std::string>(), fcnt, ntf, skippedTimeframes, describeException(e));
277+
concrete.origin.as<std::string>(), fcnt, ntf, skippedTimeframes, describeException(e));
278+
didir->markTimeFrameSkipped(header::DataHeader(concrete.description, concrete.origin, concrete.subSpec), ntf);
278279
arrowContext.clear();
279280
messageContext.discard();
280281
stringContext.clear();
@@ -361,7 +362,7 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const
361362
if (!skipInvalidReads) {
362363
throw;
363364
}
364-
skipInvalidRead(concrete.origin, e);
365+
skipInvalidRead(concrete, e);
365366
return TFReaderState::INVALID_TIMEFRAME;
366367
}
367368

Framework/AnalysisSupport/src/DataInputDirector.cxx

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -288,6 +288,15 @@ arrow::dataset::FileSource DataInputDescriptor::getFileFolder(int counter, int n
288288
return {fmt::format("DF_{}", mfilenames[counter].listOfTimeFrameNumbers[numTF]), mCurrentFilesystem};
289289
}
290290

291+
uint64_t DataInputDescriptor::markTimeFrameSkipped(int numTF)
292+
{
293+
if (mCurrentFileID >= 0 && numTF >= 0 && numTF < mfilenames[mCurrentFileID].numberOfTimeFrames) {
294+
mfilenames[mCurrentFileID].alreadyRead[numTF] = false;
295+
return ++mfilenames[mCurrentFileID].invalidReadSkipped;
296+
}
297+
return 0;
298+
}
299+
291300
std::shared_ptr<DataInputDescriptor> DataInputDescriptor::getParentFile(int counter, int numTF, std::string treename, int wantedParentLevel, std::string_view wantedOrigin)
292301
{
293302
if (!mParentFileMap) {
@@ -364,8 +373,8 @@ void DataInputDescriptor::printFileStatistics()
364373
}
365374
auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
366375
auto f = dynamic_cast<TFile*>(rootFS->GetFile());
367-
std::string monitoringInfo(fmt::format("lfn={},size={},total_df={},read_df={},read_bytes={},read_calls={},io_time={:.1f},wait_time={:.1f},level={}", f->GetName(),
368-
f->GetSize(), getTimeFramesInFile(mCurrentFileID), getReadTimeFramesInFile(mCurrentFileID), f->GetBytesRead(), f->GetReadCalls(),
376+
std::string monitoringInfo(fmt::format("lfn={},size={},total_df={},read_df={},skipped_df={},read_bytes={},read_calls={},io_time={:.1f},wait_time={:.1f},level={}", f->GetName(),
377+
f->GetSize(), getTimeFramesInFile(mCurrentFileID), getReadTimeFramesInFile(mCurrentFileID), mfilenames.at(mCurrentFileID).invalidReadSkipped, f->GetBytesRead(), f->GetReadCalls(),
369378
((float)mIOTime / 1e9), ((float)wait_time / 1e9), mLevel));
370379
#if __has_include(<TJAlienFile.h>)
371380
auto alienFile = dynamic_cast<TJAlienFile*>(f);
@@ -879,6 +888,15 @@ arrow::dataset::FileSource DataInputDirector::getFileFolder(header::DataHeader d
879888
return didesc->getFileFolder(counter, numTF, wantedLevel, origin);
880889
}
881890

891+
void DataInputDirector::markTimeFrameSkipped(header::DataHeader dh, int numTF)
892+
{
893+
auto didesc = getDataInputDescriptor(dh);
894+
if (!didesc) {
895+
didesc = mdefaultDataInputDescriptor.get();
896+
}
897+
didesc->markTimeFrameSkipped(numTF);
898+
}
899+
882900
int DataInputDirector::getTimeFramesInFile(header::DataHeader dh, int counter)
883901
{
884902
auto didesc = getDataInputDescriptor(dh);

Framework/AnalysisSupport/src/DataInputDirector.h

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,7 @@ struct FileNameHolder {
4444
int numberOfTimeFrames = 0;
4545
std::vector<uint64_t> listOfTimeFrameNumbers;
4646
std::vector<bool> alreadyRead;
47+
uint64_t invalidReadSkipped = 0;
4748
};
4849

4950
FileNameHolder makeFileNameHolder(std::string fileName);
@@ -106,6 +107,7 @@ class DataInputDescriptor
106107

107108
uint64_t getTimeFrameNumber(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin);
108109
arrow::dataset::FileSource getFileFolder(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin);
110+
uint64_t markTimeFrameSkipped(int numTF);
109111
// Open the current file to populate the parent map, then return the parent descriptor and
110112
// the TF index within it that corresponds to numTF at this level. Returns {nullptr, -1} on failure.
111113
std::pair<std::shared_ptr<DataInputDescriptor>, int> navigateToLevel(int counter, int numTF, int wantedParentLevel, std::string_view wantedOrigin);
@@ -171,6 +173,7 @@ class DataInputDirector
171173
bool readTree(DataAllocator& outputs, header::DataHeader dh, int counter, int numTF, size_t& totalSizeCompressed, size_t& totalSizeUncompressed, bool wasAOD);
172174
uint64_t getTimeFrameNumber(header::DataHeader dh, int counter, int numTF);
173175
arrow::dataset::FileSource getFileFolder(header::DataHeader dh, int counter, int numTF);
176+
void markTimeFrameSkipped(header::DataHeader dh, int numTF);
174177
int getTimeFramesInFile(header::DataHeader dh, int counter);
175178

176179
uint64_t getTotalSizeCompressed();

Framework/Core/src/runDataProcessing.cxx

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1253,6 +1253,7 @@ std::vector<std::regex> getDumpableMetrics()
12531253
dumpableMetrics.emplace_back("^aod-bytes-read-compressed$");
12541254
dumpableMetrics.emplace_back("^aod-file-read-info$");
12551255
dumpableMetrics.emplace_back("^aod-largest-object-written$");
1256+
dumpableMetrics.emplace_back("^aod-invalid-read-skipped-timeframes$");
12561257
dumpableMetrics.emplace_back("^table-bytes-.*");
12571258
dumpableMetrics.emplace_back("^total-timeframes.*");
12581259
dumpableMetrics.emplace_back("^device_state.*");

0 commit comments

Comments
 (0)