diff --git a/Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx b/Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx index 1abc4b9ffdd48..39878ce6b1127 100644 --- a/Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx +++ b/Framework/AnalysisSupport/src/AODJAlienReaderHelpers.cxx @@ -11,6 +11,7 @@ #include "AODJAlienReaderHelpers.h" #include +#include #include #include #include @@ -19,6 +20,7 @@ #include "Framework/DataProcessingStats.h" #include "Framework/RootArrowFilesystem.h" #include "Framework/AlgorithmSpec.h" +#include "Framework/ArrowContext.h" #include "Framework/ConfigParamRegistry.h" #include "Framework/ControlService.h" #include "Framework/CallbackService.h" @@ -26,6 +28,8 @@ #include "Framework/DeviceSpec.h" #include "Framework/RawDeviceService.h" #include "Framework/DataSpecUtils.h" +#include "Framework/MessageContext.h" +#include "Framework/StringContext.h" #include "Framework/ConfigContext.h" #include "DataInputDirector.h" #include "Framework/SourceInfoHeader.h" @@ -193,6 +197,10 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const int level = originLevelMapping.empty() ? -1 : 0; auto fileCounter = std::make_shared(0); auto numTF = std::make_shared(-1); + bool const skipInvalidReads = [] { + auto const* envValue = getenv("DPL_AOD_READER_SKIP_INVALID"); + return envValue != nullptr && strcmp(envValue, "0") != 0 && strcmp(envValue, "false") != 0; + }(); return adaptStateless([TFNumberHeader, TFFileNameHeader, requestedTables, @@ -200,7 +208,8 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const numTF, watchdog, maxRate, - didir, reportTFN, reportTFFileName, level](Monitoring& monitoring, DataAllocator& outputs, ControlService& control, DeviceSpec const& device, DataProcessingStats& dpstats) { + skipInvalidReads, + didir, reportTFN, reportTFFileName, level](Monitoring& monitoring, DataAllocator& outputs, ControlService& control, DeviceSpec const& device, DataProcessingStats& dpstats, ArrowContext& arrowContext, MessageContext& messageContext, StringContext& stringContext) { // Each parallel reader device.inputTimesliceId reads the files fileCounter*device.maxInputTimeslices+device.inputTimesliceId // the TF to read is numTF assert(device.inputTimesliceId < device.maxInputTimeslices); @@ -214,10 +223,10 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const } // loop over requested tables - bool first = true; static size_t totalSizeUncompressed = 0; static size_t totalSizeCompressed = 0; static uint64_t totalDFSent = 0; + static uint64_t totalInvalidReadSkipped = 0; // check if RuntimeLimit is reached if (!watchdog->update()) { @@ -232,41 +241,76 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const int64_t startTime = uv_hrtime(); int64_t startSize = totalSizeCompressed; - for (auto& route : requestedTables) { - if ((device.inputTimesliceId % route.maxTimeslices) != route.timeslice) { - continue; + auto skipInvalidRead = [&](o2::header::DataOrigin const& origin, InvalidAODReadError const& e) { + auto skippedTimeframes = ++totalInvalidReadSkipped; + LOGP(error, "Invalid AOD read for table {}: fileCounter {}, timeFrame {}. Skipping timeframe (skipped timeframes: {}). Reason: {}", + origin.as(), fcnt, ntf, skippedTimeframes, e.what()); + arrowContext.clear(); + messageContext.discard(); + stringContext.clear(); + dpstats.updateStats({static_cast(ProcessingStatsId::AOD_INVALID_READ_SKIPPED_TIMEFRAMES), DataProcessingStats::Op::Add, 1}); + *fileCounter = (fcnt - device.inputTimesliceId) / device.maxInputTimeslices; + *numTF = ntf; + }; + enum class TFReaderState { + READ_FIRST_TABLE, + READ_FIRST_TABLE_FROM_NEXT_FILE, + READ_NEXT_TABLE, + TRY_NEXT_FILE, + TIMEFRAME_READ, + INVALID_TIMEFRAME, + END_OF_INPUT, + }; + auto readState = TFReaderState::READ_FIRST_TABLE; + size_t routeIndex = 0; + while (readState == TFReaderState::READ_FIRST_TABLE || + readState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE || + readState == TFReaderState::READ_NEXT_TABLE || + readState == TFReaderState::TRY_NEXT_FILE) { + if (readState == TFReaderState::TRY_NEXT_FILE) { + fcnt += device.maxInputTimeslices; + if (didir->atEnd(fcnt)) { + readState = TFReaderState::END_OF_INPUT; + break; + } + ntf = 0; + routeIndex = 0; + readState = TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE; + } + + while (routeIndex < requestedTables.size() && + (device.inputTimesliceId % requestedTables[routeIndex].maxTimeslices) != requestedTables[routeIndex].timeslice) { + ++routeIndex; + } + if (routeIndex == requestedTables.size()) { + readState = TFReaderState::TIMEFRAME_READ; + break; } - // create header + auto& route = requestedTables[routeIndex]; auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher); auto dh = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec); bool wasAOD = std::ranges::any_of(route.matcher.metadata, [](ConfigParamSpec const& p) { return p.name.starts_with("aod-origin-replaced"); }); - if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) { - if (first) { - // check if there is a next file to read - fcnt += device.maxInputTimeslices; - if (didir->atEnd(fcnt)) { - LOGP(info, "No input files left to read for reader {}!", device.inputTimesliceId); - didir->closeInputFiles(); - monitoring.flushBuffer(); - control.endOfStream(); - control.readyToQuit(QuitRequest::Me); - return; + try { + if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) { + if (readState == TFReaderState::READ_FIRST_TABLE) { + readState = TFReaderState::TRY_NEXT_FILE; + continue; } - // get first folder of next file - ntf = 0; - if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) { - LOGP(fatal, "Can not retrieve tree for table {}: fileCounter {}, timeFrame {}", concrete.origin.as(), fcnt, ntf); - throw std::runtime_error("Processing is stopped!"); - } - } else { LOGP(fatal, "Can not retrieve tree for table {}: fileCounter {}, timeFrame {}", concrete.origin.as(), fcnt, ntf); throw std::runtime_error("Processing is stopped!"); } + } catch (InvalidAODReadError const& e) { + if (!skipInvalidReads) { + throw; + } + skipInvalidRead(concrete.origin, e); + readState = TFReaderState::INVALID_TIMEFRAME; + break; } - if (first) { + if (readState == TFReaderState::READ_FIRST_TABLE || readState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE) { if (reportTFN) { // TF number auto timeFrameNumber = didir->getTimeFrameNumber(dh, fcnt, ntf); @@ -289,7 +333,27 @@ AlgorithmSpec AODJAlienReaderHelpers::rootFileReaderCallback(ConfigContext const outputs.make(o2) = currentFilename; } } - first = false; + ++routeIndex; + readState = TFReaderState::READ_NEXT_TABLE; + } + + switch (readState) { + case TFReaderState::TIMEFRAME_READ: + break; + case TFReaderState::INVALID_TIMEFRAME: + return; + case TFReaderState::END_OF_INPUT: + LOGP(info, "No input files left to read for reader {}!", device.inputTimesliceId); + didir->closeInputFiles(); + monitoring.flushBuffer(); + control.endOfStream(); + control.readyToQuit(QuitRequest::Me); + return; + case TFReaderState::READ_FIRST_TABLE: + case TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE: + case TFReaderState::READ_NEXT_TABLE: + case TFReaderState::TRY_NEXT_FILE: + throw std::logic_error("Invalid timeframe read state"); } int64_t stopSize = totalSizeCompressed; int64_t bytesDelta = stopSize - startSize; diff --git a/Framework/AnalysisSupport/src/DataInputDirector.cxx b/Framework/AnalysisSupport/src/DataInputDirector.cxx index cfd578862fabd..4431a95f8e0fb 100644 --- a/Framework/AnalysisSupport/src/DataInputDirector.cxx +++ b/Framework/AnalysisSupport/src/DataInputDirector.cxx @@ -34,6 +34,7 @@ #include #include #include +#include #include #if __has_include() @@ -536,18 +537,25 @@ bool DataInputDescriptor::readTree(DataAllocator& outputs, header::DataHeader dh if (!format) { t.deactivate(); LOGP(debug, "Could not find tree {}. Trying in parent file.", fullpath.path()); - auto parentFile = getParentFile(counter, numTF, treename, wantedLevel, wantedOrigin); - if (parentFile != nullptr) { - int parentNumTF = parentFile->findDFNumber(0, folder.path()); - if (parentNumTF == -1) { - auto parentRootFS = std::dynamic_pointer_cast(parentFile->mCurrentFilesystem); - throw std::runtime_error(fmt::format(R"(DF {} listed in parent file map but not found in the corresponding file "{}")", folder.path(), parentRootFS->GetFile()->GetName())); - } - // first argument is 0 as the parent file object contains only 1 file - return parentFile->readTree(outputs, dh, 0, parentNumTF, treename, totalSizeCompressed, totalSizeUncompressed); + std::shared_ptr parentFile; + try { + parentFile = getParentFile(counter, numTF, treename, wantedLevel, wantedOrigin); + } catch (std::exception const& e) { + throw InvalidAODReadError(fmt::format("Unable to resolve parent file for tree {}: {}", treename, e.what())); + } catch (...) { + throw InvalidAODReadError(fmt::format("Unable to resolve parent file for tree {}", treename)); + } + if (parentFile == nullptr) { + auto rootFS = std::dynamic_pointer_cast(mCurrentFilesystem); + 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())); } - auto rootFS = std::dynamic_pointer_cast(mCurrentFilesystem); - 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())); + int parentNumTF = parentFile->findDFNumber(0, folder.path()); + if (parentNumTF == -1) { + auto parentRootFS = std::dynamic_pointer_cast(parentFile->mCurrentFilesystem); + throw InvalidAODReadError(fmt::format(R"(DF {} listed in parent file map but not found in the corresponding file "{}")", folder.path(), parentRootFS->GetFile()->GetName())); + } + // first argument is 0 as the parent file object contains only 1 file + return parentFile->readTree(outputs, dh, 0, parentNumTF, treename, totalSizeCompressed, totalSizeUncompressed); } auto schemaOpt = format->Inspect(fullpath); @@ -573,7 +581,23 @@ bool DataInputDescriptor::readTree(DataAllocator& outputs, header::DataHeader dh //// add branches to read //// fill the table f2b->setLabel(treename.c_str()); - f2b->fill(datasetSchema, format); + try { + f2b->fill(datasetSchema, format); + } catch (std::exception const& e) { + f2b.discard(); + throw InvalidAODReadError(fmt::format("Unable to read tree {}: {}", treename, e.what())); + } catch (...) { + f2b.discard(); + throw InvalidAODReadError(fmt::format("Unable to read tree {}", treename)); + } + + try { + f2b.release(); + } catch (std::exception const& e) { + throw InvalidAODReadError(fmt::format("Unable to finalize tree {}: {}", treename, e.what())); + } catch (...) { + throw InvalidAODReadError(fmt::format("Unable to finalize tree {}", treename)); + } return true; } diff --git a/Framework/AnalysisSupport/src/DataInputDirector.h b/Framework/AnalysisSupport/src/DataInputDirector.h index 17535f2935ba3..a810619530e10 100644 --- a/Framework/AnalysisSupport/src/DataInputDirector.h +++ b/Framework/AnalysisSupport/src/DataInputDirector.h @@ -21,6 +21,7 @@ #include #include +#include #include #include "rapidjson/fwd.h" @@ -32,6 +33,12 @@ class Monitoring; namespace o2::framework { +class InvalidAODReadError : public std::runtime_error +{ + public: + using std::runtime_error::runtime_error; +}; + struct FileNameHolder { std::string fileName; int numberOfTimeFrames = 0; diff --git a/Framework/Core/include/Framework/DataProcessingStats.h b/Framework/Core/include/Framework/DataProcessingStats.h index edb04c4c5f752..e164e11cb2134 100644 --- a/Framework/Core/include/Framework/DataProcessingStats.h +++ b/Framework/Core/include/Framework/DataProcessingStats.h @@ -74,6 +74,7 @@ enum struct ProcessingStatsId : short { CCDB_CACHE_FAILURE, CCDB_CACHE_FETCHED_BYTES, CCDB_CACHE_REQUESTED_BYTES, + AOD_INVALID_READ_SKIPPED_TIMEFRAMES, AVAILABLE_MANAGED_SHM_BASE = 512, }; diff --git a/Framework/Core/src/CommonServices.cxx b/Framework/Core/src/CommonServices.cxx index 2ac9dab40d20a..83cdf31833cce 100644 --- a/Framework/Core/src/CommonServices.cxx +++ b/Framework/Core/src/CommonServices.cxx @@ -1103,6 +1103,13 @@ o2::framework::ServiceSpec CommonServices::dataProcessingStats() MetricSpec{.name = "dropped_computations", .metricId = static_cast(ProcessingStatsId::DROPPED_COMPUTATIONS), .kind = Kind::UInt64, .minPublishInterval = quickUpdateInterval}, MetricSpec{.name = "dropped_incoming_messages", .metricId = static_cast(ProcessingStatsId::DROPPED_INCOMING_MESSAGES), .kind = Kind::UInt64, .minPublishInterval = quickUpdateInterval}, MetricSpec{.name = "relayed_messages", .metricId = static_cast(ProcessingStatsId::RELAYED_MESSAGES), .kind = Kind::UInt64, .minPublishInterval = quickUpdateInterval}, + MetricSpec{.name = "aod-invalid-read-skipped-timeframes", + .metricId = static_cast(ProcessingStatsId::AOD_INVALID_READ_SKIPPED_TIMEFRAMES), + .kind = Kind::UInt64, + .scope = Scope::DPL, + .minPublishInterval = 0, + .maxRefreshLatency = 10000, + .sendInitialValue = true}, MetricSpec{.name = "arrow-bytes-destroyed", .enabled = arrowAndResourceLimitingMetrics, .metricId = static_cast(ProcessingStatsId::ARROW_BYTES_DESTROYED),