Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -268,10 +268,11 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const

int64_t startTime = uv_hrtime();
int64_t startSize = totalSizeCompressed;
auto skipInvalidRead = [&](o2::header::DataOrigin const& origin, InvalidAODReadError const& e) {
auto skipInvalidRead = [&](ConcreteDataMatcher const& concrete, InvalidAODReadError const& e) {
auto skippedTimeframes = ++totalInvalidReadSkipped;
LOGP(error, "Invalid AOD read for table {}: fileCounter {}, timeFrame {}. Skipping timeframe (skipped timeframes: {}). Reason: {}",
origin.as<std::string>(), fcnt, ntf, skippedTimeframes, describeException(e));
concrete.origin.as<std::string>(), fcnt, ntf, skippedTimeframes, describeException(e));
didir->markTimeFrameSkipped(header::DataHeader(concrete.description, concrete.origin, concrete.subSpec), ntf);
arrowContext.clear();
messageContext.discard();
stringContext.clear();
Expand Down Expand Up @@ -334,7 +335,7 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const
if (!skipInvalidReads) {
throw;
}
skipInvalidRead(concrete.origin, e);
skipInvalidRead(concrete, e);
return TFReaderState::INVALID_TIMEFRAME;
}

Expand Down
22 changes: 20 additions & 2 deletions Framework/AnalysisSupport/src/DataInputDirector.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -288,6 +288,15 @@ arrow::dataset::FileSource DataInputDescriptor::getFileFolder(int counter, int n
return {fmt::format("DF_{}", mfilenames[counter].listOfTimeFrameNumbers[numTF]), mCurrentFilesystem};
}

uint64_t DataInputDescriptor::markTimeFrameSkipped(int numTF)
{
if (mCurrentFileID >= 0 && numTF >= 0 && numTF < mfilenames[mCurrentFileID].numberOfTimeFrames) {
mfilenames[mCurrentFileID].alreadyRead[numTF] = false;
return ++mfilenames[mCurrentFileID].invalidReadSkipped;
}
return 0;
}

std::shared_ptr<DataInputDescriptor> DataInputDescriptor::getParentFile(int counter, int numTF, std::string treename, int wantedParentLevel, std::string_view wantedOrigin)
{
if (!mParentFileMap) {
Expand Down Expand Up @@ -364,8 +373,8 @@ void DataInputDescriptor::printFileStatistics()
}
auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(mCurrentFilesystem);
auto f = dynamic_cast<TFile*>(rootFS->GetFile());
std::string monitoringInfo(fmt::format("lfn={},size={},total_df={},read_df={},read_bytes={},read_calls={},io_time={:.1f},wait_time={:.1f},level={}", f->GetName(),
f->GetSize(), getTimeFramesInFile(mCurrentFileID), getReadTimeFramesInFile(mCurrentFileID), f->GetBytesRead(), f->GetReadCalls(),
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(),
f->GetSize(), getTimeFramesInFile(mCurrentFileID), getReadTimeFramesInFile(mCurrentFileID), mfilenames.at(mCurrentFileID).invalidReadSkipped, f->GetBytesRead(), f->GetReadCalls(),
((float)mIOTime / 1e9), ((float)wait_time / 1e9), mLevel));
#if __has_include(<TJAlienFile.h>)
auto alienFile = dynamic_cast<TJAlienFile*>(f);
Expand Down Expand Up @@ -879,6 +888,15 @@ arrow::dataset::FileSource DataInputDirector::getFileFolder(header::DataHeader d
return didesc->getFileFolder(counter, numTF, wantedLevel, origin);
}

void DataInputDirector::markTimeFrameSkipped(header::DataHeader dh, int numTF)
{
auto didesc = getDataInputDescriptor(dh);
if (!didesc) {
didesc = mdefaultDataInputDescriptor.get();
}
didesc->markTimeFrameSkipped(numTF);
}

int DataInputDirector::getTimeFramesInFile(header::DataHeader dh, int counter)
{
auto didesc = getDataInputDescriptor(dh);
Expand Down
3 changes: 3 additions & 0 deletions Framework/AnalysisSupport/src/DataInputDirector.h
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@ struct FileNameHolder {
int numberOfTimeFrames = 0;
std::vector<uint64_t> listOfTimeFrameNumbers;
std::vector<bool> alreadyRead;
uint64_t invalidReadSkipped = 0;
};

FileNameHolder makeFileNameHolder(std::string fileName);
Expand Down Expand Up @@ -106,6 +107,7 @@ class DataInputDescriptor

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

uint64_t getTotalSizeCompressed();
Expand Down
1 change: 1 addition & 0 deletions Framework/Core/src/runDataProcessing.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -1253,6 +1253,7 @@ std::vector<std::regex> getDumpableMetrics()
dumpableMetrics.emplace_back("^aod-bytes-read-compressed$");
dumpableMetrics.emplace_back("^aod-file-read-info$");
dumpableMetrics.emplace_back("^aod-largest-object-written$");
dumpableMetrics.emplace_back("^aod-invalid-read-skipped-timeframes$");
dumpableMetrics.emplace_back("^table-bytes-.*");
dumpableMetrics.emplace_back("^total-timeframes.*");
dumpableMetrics.emplace_back("^device_state.*");
Expand Down