Project
Loading...
Searching...
No Matches
AODJAlienReaderHelpers.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
13#include <algorithm>
14#include <charconv>
15#include <cctype>
16#include <cstdlib>
17#include <exception>
18#include <memory>
19#include <ranges>
20#include <string>
21#include <string_view>
22#include <vector>
38#include "Framework/Signpost.h"
41#include "DataInputDirector.h"
44#include "Framework/Logger.h"
45
46#if __has_include(<TJAlienFile.h>)
47#include <TJAlienFile.h>
48#endif
49#include <TGrid.h>
50#include <TFile.h>
51#include <TTreeCache.h>
52#include <TSystem.h>
53
54#include <arrow/ipc/reader.h>
55#include <arrow/ipc/writer.h>
56#include <arrow/io/interfaces.h>
57#include <arrow/table.h>
58#include <arrow/util/key_value_metadata.h>
59#include <arrow/dataset/dataset.h>
60#include <arrow/dataset/file_base.h>
61
62using namespace o2;
63using namespace o2::aod;
64
66
69 uint64_t startTime;
70 uint64_t lastTime;
71 double runTime;
72 uint64_t runTimeLimit;
73
74 RuntimeWatchdog(Long64_t limit)
75 {
77 startTime = uv_hrtime();
79 runTime = 0.;
80 runTimeLimit = limit;
81 }
82
83 bool update()
84 {
86 if (runTimeLimit <= 0) {
87 return true;
88 }
89
90 auto nowTime = uv_hrtime();
91
92 // time spent to process the time frame
93 double time_spent = numberTimeFrames < 1 ? (double)(nowTime - lastTime) / 1.E9 : 0.;
94 runTime += time_spent;
95 lastTime = nowTime;
96
97 return ((double)(lastTime - startTime) / 1.E9 + runTime / (numberTimeFrames + 1)) < runTimeLimit;
98 }
99
100 void printOut()
101 {
102 LOGP(info, "RuntimeWatchdog");
103 LOGP(info, " run time limit: {}", runTimeLimit);
104 LOGP(info, " number of time frames: {}", numberTimeFrames);
105 LOGP(info, " estimated run time per time frame: {}", (numberTimeFrames >= 0) ? runTime / (numberTimeFrames + 1) : 0.);
106 LOGP(info, " estimated total run time: {}", (double)(lastTime - startTime) / 1.E9 + ((numberTimeFrames >= 0) ? runTime / (numberTimeFrames + 1) : 0.));
107 }
108};
109
110using o2::monitoring::Metric;
111using o2::monitoring::Monitoring;
112using o2::monitoring::tags::Key;
113using o2::monitoring::tags::Value;
114
116{
117static bool shouldSkipInvalidReads()
118{
119 auto const* envValue = getenv("DPL_AOD_READER_SKIP_INVALID");
120 if (envValue == nullptr) {
121 return false;
122 }
123
124 std::string value{envValue};
125 std::ranges::transform(value, value.begin(), [](unsigned char c) { return std::tolower(c); });
126 return !value.empty() && value != "0" && value != "false";
127}
128
129static std::string describeException(std::exception const& exception)
130{
131 std::string description{exception.what()};
132 try {
133 std::rethrow_if_nested(exception);
134 } catch (std::exception const& nested) {
135 description += fmt::format(": {}", describeException(nested));
136 } catch (RuntimeErrorRef const& ref) {
137 description += fmt::format(": {}", error_from_ref(ref).what);
138 } catch (...) {
139 description += ": unknown exception";
140 }
141 return description;
142}
143
145{
146 // aod-parent-base-path-replacement is now a workflow option, so it needs to be
147 // retrieved from the ConfigContext. This is because we do not allow workflow options
148 // to change over start-stop-start because they can affect the topology generation.
149 std::string parentFileReplacement;
150 if (ctx.options().isSet("aod-parent-base-path-replacement")) {
151 parentFileReplacement = ctx.options().get<std::string>("aod-parent-base-path-replacement");
152 }
153 int parentAccessLevel = 0;
154 if (ctx.options().isSet("aod-parent-access-level")) {
155 parentAccessLevel = ctx.options().get<int>("aod-parent-access-level");
156 }
157 std::vector<std::pair<std::string, int>> originLevelMapping;
158 if (ctx.options().isSet("aod-origin-level-mapping")) {
159 auto originLevelMappingStr = ctx.options().get<std::string>("aod-origin-level-mapping");
160 for (auto pairRange : originLevelMappingStr | std::views::split(',')) {
161 std::string_view pair{pairRange.begin(), pairRange.end()};
162 auto colonPos = pair.find(':');
163 if (colonPos == std::string_view::npos) {
164 LOGP(fatal, "Badly formatted aod-origin-level-mapping entry: \"{}\"", pair);
165 continue;
166 }
167 std::string key(pair.substr(0, colonPos));
168 std::string_view valueStr = pair.substr(colonPos + 1);
169 int value{};
170 auto [ptr, ec] = std::from_chars(valueStr.data(), valueStr.data() + valueStr.size(), value);
171 if (ec == std::errc{}) {
172 originLevelMapping.emplace_back(std::move(key), value);
173 } else {
174 LOGP(fatal, "Unable to parse level in aod-origin-level-mapping entry: \"{}\"", pair);
175 }
176 }
177 }
178 auto callback = AlgorithmSpec{adaptStateful([parentFileReplacement, parentAccessLevel, originLevelMapping](ConfigParamRegistry const& options,
179 DeviceSpec const& spec,
180 Monitoring& monitoring,
181 DataProcessingStats& stats) {
182 // FIXME: not actually needed, since data processing stats can specify that we should
183 // send the initial value.
184 stats.updateStats({static_cast<short>(ProcessingStatsId::ARROW_BYTES_CREATED), DataProcessingStats::Op::Set, 0});
185 stats.updateStats({static_cast<short>(ProcessingStatsId::ARROW_MESSAGES_CREATED), DataProcessingStats::Op::Set, 0});
186 stats.updateStats({static_cast<short>(ProcessingStatsId::ARROW_BYTES_DESTROYED), DataProcessingStats::Op::Set, 0});
187 stats.updateStats({static_cast<short>(ProcessingStatsId::ARROW_MESSAGES_DESTROYED), DataProcessingStats::Op::Set, 0});
188 stats.updateStats({static_cast<short>(ProcessingStatsId::ARROW_BYTES_EXPIRED), DataProcessingStats::Op::Set, 0});
189 stats.updateStats({static_cast<short>(ProcessingStatsId::CONSUMED_TIMEFRAMES), DataProcessingStats::Op::Set, 0});
190
191 if (!options.isSet("aod-file-private")) {
192 LOGP(fatal, "No input file defined!");
193 throw std::runtime_error("Processing is stopped!");
194 }
195
196 auto filename = options.get<std::string>("aod-file-private");
197
198 auto maxRate = options.get<float>("aod-max-io-rate");
199
200 // create a DataInputDirector
201 auto didir = std::make_shared<DataInputDirector>(std::vector<std::string>{filename}, DataInputDirectorContext{&monitoring, parentAccessLevel, parentFileReplacement, originLevelMapping});
202 if (options.isSet("aod-reader-json")) {
203 auto jsonFile = options.get<std::string>("aod-reader-json");
204 if (!didir->readJson(jsonFile)) {
205 LOGP(error, "Check the JSON document! Can not be properly parsed!");
206 }
207 }
208
209 // get the run time watchdog
210 auto* watchdog = new RuntimeWatchdog(options.get<int64_t>("time-limit"));
211
212 // selected the TFN input and
213 // create list of requested tables
214 bool reportTFN = false;
215 bool reportTFFileName = false;
216 header::DataHeader TFNumberHeader;
217 header::DataHeader TFFileNameHeader;
218 std::vector<OutputRoute> requestedTables;
219 std::vector<OutputRoute> routes(spec.outputs);
220 for (auto route : routes) {
221 if (DataSpecUtils::partialMatch(route.matcher, header::DataOrigin("TFN"))) {
222 auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher);
223 TFNumberHeader = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
224 reportTFN = true;
225 } else if (DataSpecUtils::partialMatch(route.matcher, header::DataOrigin("TFF"))) {
226 auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher);
227 TFFileNameHeader = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
228 reportTFFileName = true;
229 } else {
230 requestedTables.emplace_back(route);
231 }
232 }
233 int level = originLevelMapping.empty() ? -1 : 0;
234 auto fileCounter = std::make_shared<int>(0);
235 auto numTF = std::make_shared<int>(-1);
236 bool const skipInvalidReads = shouldSkipInvalidReads();
237 return adaptStateless([TFNumberHeader,
238 TFFileNameHeader,
239 requestedTables,
240 fileCounter,
241 numTF,
242 watchdog,
243 maxRate,
244 skipInvalidReads,
245 didir, reportTFN, reportTFFileName, level](Monitoring& monitoring, DataAllocator& outputs, ControlService& control, DeviceSpec const& device, DataProcessingStats& dpstats, ArrowContext& arrowContext, MessageContext& messageContext, StringContext& stringContext) {
246 // Each parallel reader device.inputTimesliceId reads the files fileCounter*device.maxInputTimeslices+device.inputTimesliceId
247 // the TF to read is numTF
248 assert(device.inputTimesliceId < device.maxInputTimeslices);
249 int fcnt = (*fileCounter * device.maxInputTimeslices) + device.inputTimesliceId;
250 int ntf = *numTF + 1;
251 static int currentFileCounter = -1;
252 static int filesProcessed = 0;
253 if (currentFileCounter != *fileCounter) {
254 currentFileCounter = *fileCounter;
255 monitoring.send(Metric{(uint64_t)++filesProcessed, "files-opened"}.addTag(Key::Subsystem, monitoring::tags::Value::DPL));
256 }
257
258 // loop over requested tables
259 static size_t totalSizeUncompressed = 0;
260 static size_t totalSizeCompressed = 0;
261 static uint64_t totalDFSent = 0;
262 static uint64_t totalInvalidReadSkipped = 0;
263
264 // check if RuntimeLimit is reached
265 if (!watchdog->update()) {
266 LOGP(info, "Run time exceeds run time limit of {} seconds. Exiting gracefully...", watchdog->runTimeLimit);
267 LOGP(info, "Stopping reader {} after time frame {}.", device.inputTimesliceId, watchdog->numberTimeFrames - 1);
268 didir->closeInputFiles();
269 monitoring.flushBuffer();
270 control.endOfStream();
272 return;
273 }
274
275 int64_t startTime = uv_hrtime();
276 int64_t startSize = totalSizeCompressed;
277 auto skipInvalidRead = [&](ConcreteDataMatcher const& concrete, InvalidAODReadError const& e) {
278 auto skippedTimeframes = ++totalInvalidReadSkipped;
279 LOGP(error, "Invalid AOD read for table {}: fileCounter {}, timeFrame {}. Skipping timeframe (skipped timeframes: {}). Reason: {}",
280 concrete.origin.as<std::string>(), fcnt, ntf, skippedTimeframes, describeException(e));
282 didir->finishTimeFrame(true);
283 arrowContext.clear();
284 messageContext.discard();
285 stringContext.clear();
287 *fileCounter = (fcnt - device.inputTimesliceId) / device.maxInputTimeslices;
288 *numTF = ntf;
289 };
290 enum class TFReaderState {
291 READ_FIRST_TABLE,
292 READ_FIRST_TABLE_FROM_NEXT_FILE,
293 READ_NEXT_TABLE,
294 TRY_NEXT_FILE,
295 TIMEFRAME_READ,
296 INVALID_TIMEFRAME,
297 };
298 auto readState = TFReaderState::READ_FIRST_TABLE;
299 [[maybe_unused]] auto stateName = [](TFReaderState state) -> char const* {
300 switch (state) {
301 case TFReaderState::READ_FIRST_TABLE:
302 return "READ_FIRST_TABLE";
303 case TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE:
304 return "READ_FIRST_TABLE_FROM_NEXT_FILE";
305 case TFReaderState::READ_NEXT_TABLE:
306 return "READ_NEXT_TABLE";
307 case TFReaderState::TRY_NEXT_FILE:
308 return "TRY_NEXT_FILE";
309 case TFReaderState::TIMEFRAME_READ:
310 return "TIMEFRAME_READ";
311 case TFReaderState::INVALID_TIMEFRAME:
312 return "INVALID_TIMEFRAME";
313 }
314 return "UNKNOWN";
315 };
316 O2_SIGNPOST_ID_FROM_POINTER(readerStateId, aod_reader, &readState);
317 auto transitionTo = [&](TFReaderState nextState) {
318 O2_SIGNPOST_EVENT_EMIT(aod_reader, readerStateId, "state transition",
319 "%{public}s -> %{public}s (fileCounter %d, timeFrame %d)",
320 stateName(readState), stateName(nextState), fcnt, ntf);
321 readState = nextState;
322 };
323 size_t routeIndex = 0;
324 auto reportTimeframe = [&didir, &fcnt, &ntf, &outputs, &TFNumberHeader, &TFFileNameHeader, reportTFN, reportTFFileName](header::DataHeader const& dh) {
325 if (reportTFN) {
326 // TF number
327 auto timeFrameNumber = didir->getTimeFrameNumber(dh, fcnt, ntf);
328 auto o = Output(TFNumberHeader);
329 outputs.make<uint64_t>(o) = timeFrameNumber;
330 }
331
332 if (reportTFFileName) {
333 // Origin file name for derived output map
334 auto o2 = Output(TFFileNameHeader);
335 auto fileAndFolder = didir->getFileFolder(dh, fcnt, ntf);
336 auto rootFS = std::dynamic_pointer_cast<TFileFileSystem>(fileAndFolder.filesystem());
337 auto* f = dynamic_cast<TFile*>(rootFS->GetFile());
338 std::string currentFilename(f->GetFile()->GetName());
339 if (strcmp(f->GetEndpointUrl()->GetProtocol(), "file") == 0 && f->GetEndpointUrl()->GetFile()[0] != '/') {
340 // This is not an absolute local path. Make it absolute.
341 static std::string pwd = gSystem->pwd() + std::string("/");
342 currentFilename = pwd + std::string(f->GetName());
343 }
344 outputs.make<std::string>(o2) = currentFilename;
345 }
346 };
347 auto tryReadTable = [&device, &didir, &fcnt, &ntf, &outputs, &reportTimeframe, &requestedTables, &routeIndex, &skipInvalidRead, skipInvalidReads](TFReaderState currentState) -> TFReaderState {
348 while (routeIndex < requestedTables.size() &&
349 (device.inputTimesliceId % requestedTables[routeIndex].maxTimeslices) != requestedTables[routeIndex].timeslice) {
350 ++routeIndex;
351 }
352 if (routeIndex == requestedTables.size()) {
353 return TFReaderState::TIMEFRAME_READ;
354 }
355
356 auto& route = requestedTables[routeIndex];
357 auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher);
358 auto dh = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
359 bool wasAOD = std::ranges::any_of(route.matcher.metadata, [](ConfigParamSpec const& p) { return p.name.starts_with("aod-origin-replaced"); });
360
361 if (currentState == TFReaderState::READ_FIRST_TABLE || currentState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE) {
362 didir->beginTimeFrame();
363 }
364 try {
365 if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) {
366 return TFReaderState::TRY_NEXT_FILE;
367 }
368 } catch (InvalidAODReadError const& e) {
369 if (!skipInvalidReads) {
370 throw;
371 }
372 skipInvalidRead(concrete, e);
373 return TFReaderState::INVALID_TIMEFRAME;
374 }
375
376 if (currentState == TFReaderState::READ_FIRST_TABLE || currentState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE) {
377 reportTimeframe(dh);
378 }
379 ++routeIndex;
380 return TFReaderState::READ_NEXT_TABLE;
381 };
382 while (readState != TFReaderState::TIMEFRAME_READ) {
383 switch (readState) {
384 case TFReaderState::READ_FIRST_TABLE:
385 transitionTo(tryReadTable(readState));
386 break;
387 case TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE:
388 case TFReaderState::READ_NEXT_TABLE:
389 transitionTo(tryReadTable(readState));
390 if (readState == TFReaderState::TRY_NEXT_FILE) {
391 // Once a file has been selected, every requested table must exist.
392 auto concrete = DataSpecUtils::asConcreteDataMatcher(requestedTables[routeIndex].matcher);
393 LOGP(fatal, "Can not retrieve tree for table {}: fileCounter {}, timeFrame {}", concrete.origin.as<std::string>(), fcnt, ntf);
394 throw std::runtime_error("Processing is stopped!");
395 }
396 break;
397 case TFReaderState::TRY_NEXT_FILE:
398 didir->finishTimeFrame();
399 fcnt += device.maxInputTimeslices;
400 if (didir->atEnd(fcnt)) {
401 LOGP(info, "No input files left to read for reader {}!", device.inputTimesliceId);
402 didir->closeInputFiles();
403 monitoring.flushBuffer();
404 control.endOfStream();
406 return;
407 }
408 ntf = 0;
409 routeIndex = 0;
410 transitionTo(TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE);
411 break;
412 case TFReaderState::INVALID_TIMEFRAME:
413 return;
414 case TFReaderState::TIMEFRAME_READ:
415 break;
416 }
417 }
418 didir->finishTimeFrame();
419 int64_t stopSize = totalSizeCompressed;
420 int64_t bytesDelta = stopSize - startSize;
421 int64_t stopTime = uv_hrtime();
422 float currentDelta = float(stopTime - startTime) / 1000000000; // in s
423 if (ceil(maxRate) > 0.) {
424 float extraTime = (bytesDelta / 1000000 - currentDelta * maxRate) / maxRate;
425 // We only sleep if we read faster than the max-read-rate.
426 if (extraTime > 0.) {
427 LOGP(info, "Read {} MB in {} s. Sleeping for {} seconds to stay within {} MB/s limit.", bytesDelta / 1000000, currentDelta, extraTime, maxRate);
428 uv_sleep(extraTime * 1000); // in milliseconds
429 }
430 }
431 totalDFSent++;
432
433 // Use the new API for sending TIMESLICE_NUMBER_STARTED
435 dpstats.processCommandQueue();
436 monitoring.send(Metric{(uint64_t)totalDFSent, "df-sent"}.addTag(Key::Subsystem, monitoring::tags::Value::DPL));
437 monitoring.send(Metric{(uint64_t)totalSizeUncompressed / 1000, "aod-bytes-read-uncompressed"}.addTag(Key::Subsystem, monitoring::tags::Value::DPL));
438 monitoring.send(Metric{(uint64_t)totalSizeCompressed / 1000, "aod-bytes-read-compressed"}.addTag(Key::Subsystem, monitoring::tags::Value::DPL));
439
440 // save file number and time frame
441 *fileCounter = (fcnt - device.inputTimesliceId) / device.maxInputTimeslices;
442 *numTF = ntf;
443
444 // Check if the next timeframe is available or
445 // if there are more files to be processed. If not, simply exit.
446 ntf = *numTF + 1;
447 // first route with level 0 or -1 if no mapping requested
448 auto firstRoute = std::ranges::find_if(requestedTables, [&didir, level](auto const& route) {
449 auto concrete = DataSpecUtils::asConcreteDataMatcher(route.matcher);
450 return didir->getLevelForOrigin(concrete.origin) == level;
451 });
452 auto concrete = DataSpecUtils::asConcreteDataMatcher(firstRoute->matcher);
453 auto dh = header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
454 auto fileAndFolder = didir->getFileFolder(dh, fcnt, ntf);
455
456 // In case the filesource is empty, move to the next one.
457 if (fileAndFolder.filesystem() == nullptr) {
458 fcnt += 1;
459 ntf = 0;
460 if (didir->atEnd(fcnt)) {
461 LOGP(info, "No input files left to read for reader {}!", device.inputTimesliceId);
462 didir->closeInputFiles();
463 monitoring.flushBuffer();
464 control.endOfStream();
466 return;
467 }
468 }
469 });
470 })};
471
472 return callback;
473}
474
475} // namespace o2::framework::readers
header::DataDescription description
std::vector< OutputRoute > routes
o2::monitoring::Metric Metric
SurfaceTrackState state
uint32_t c
Definition RawData.h:2
#define O2_DECLARE_DYNAMIC_LOG(name)
Definition Signpost.h:490
#define O2_SIGNPOST_ID_FROM_POINTER(name, log, pointer)
Definition Signpost.h:506
#define O2_SIGNPOST_EVENT_EMIT(log, id, name, format,...)
Definition Signpost.h:523
TBranch * ptr
o2::monitoring::Monitoring Monitoring
StringRef key
void readyToQuit(bool all)
Compatibility with old API.
void endOfStream()
Signal that we are done with the current stream.
GLdouble f
Definition glcorearb.h:310
GLsizei const GLfloat * value
Definition glcorearb.h:819
GLint level
Definition glcorearb.h:275
GLint ref
Definition glcorearb.h:291
constexpr framework::ConcreteDataMatcher matcher()
Definition ASoA.h:388
void clean_all_runtime_errors()
RuntimeError & error_from_ref(RuntimeErrorRef)
@ Me
Only quit this data processor.
AlgorithmSpec::ProcessCallback adaptStateless(LAMBDA l)
AlgorithmSpec::InitCallback adaptStateful(LAMBDA l)
a couple of static helper functions to create timestamp values for CCDB queries or override obsolete ...
std::string filename()
RuntimeWatchdog(Long64_t limit)
header::DataHeader::SubSpecificationType subSpec
Helper struct to hold statistics about the data processing happening.
@ Add
Update the rate of the metric given the amount since the last time.
static bool partialMatch(InputSpec const &spec, o2::header::DataOrigin const &origin)
static ConcreteDataMatcher asConcreteDataMatcher(InputSpec const &input)
size_t maxInputTimeslices
The maximum number of time pipelining for this device.
Definition DeviceSpec.h:70
size_t inputTimesliceId
The time pipelining id of this particular device.
Definition DeviceSpec.h:68
std::vector< OutputRoute > outputs
Definition DeviceSpec.h:63
static AlgorithmSpec rootFileReaderCallback(ConfigContext const &context)
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