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");
153 int parentAccessLevel = 0;
154 if (ctx.options().isSet(
"aod-parent-access-level")) {
155 parentAccessLevel = ctx.options().get<
int>(
"aod-parent-access-level");
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);
167 std::string
key(pair.substr(0, colonPos));
168 std::string_view valueStr = pair.substr(colonPos + 1);
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);
174 LOGP(fatal,
"Unable to parse level in aod-origin-level-mapping entry: \"{}\"", pair);
191 if (!options.
isSet(
"aod-file-private")) {
192 LOGP(fatal,
"No input file defined!");
193 throw std::runtime_error(
"Processing is stopped!");
196 auto filename = options.
get<std::string>(
"aod-file-private");
198 auto maxRate = options.
get<
float>(
"aod-max-io-rate");
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!");
214 bool reportTFN =
false;
215 bool reportTFFileName =
false;
218 std::vector<OutputRoute> requestedTables;
220 for (
auto route :
routes) {
223 TFNumberHeader =
header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
227 TFFileNameHeader =
header::DataHeader(concrete.description, concrete.origin, concrete.subSpec);
228 reportTFFileName =
true;
230 requestedTables.emplace_back(route);
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();
245 didir, reportTFN, reportTFFileName,
level](
Monitoring& monitoring,
DataAllocator& outputs,
ControlService& control,
DeviceSpec const& device,
DataProcessingStats& dpstats,
ArrowContext& arrowContext,
MessageContext& messageContext,
StringContext& stringContext) {
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));
259 static size_t totalSizeUncompressed = 0;
260 static size_t totalSizeCompressed = 0;
261 static uint64_t totalDFSent = 0;
262 static uint64_t totalInvalidReadSkipped = 0;
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();
275 int64_t startTime = uv_hrtime();
276 int64_t startSize = totalSizeCompressed;
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();
285 stringContext.
clear();
290 enum class TFReaderState {
292 READ_FIRST_TABLE_FROM_NEXT_FILE,
298 auto readState = TFReaderState::READ_FIRST_TABLE;
299 [[maybe_unused]]
auto stateName = [](TFReaderState
state) ->
char const* {
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";
317 auto transitionTo = [&](TFReaderState nextState) {
319 "%{public}s -> %{public}s (fileCounter %d, timeFrame %d)",
320 stateName(readState), stateName(nextState), fcnt, ntf);
321 readState = nextState;
323 size_t routeIndex = 0;
324 auto reportTimeframe = [&didir, &fcnt, &ntf, &outputs, &TFNumberHeader, &TFFileNameHeader, reportTFN, reportTFFileName](
header::DataHeader const& dh) {
327 auto timeFrameNumber = didir->getTimeFrameNumber(dh, fcnt, ntf);
328 auto o =
Output(TFNumberHeader);
329 outputs.make<uint64_t>(o) = timeFrameNumber;
332 if (reportTFFileName) {
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] !=
'/') {
341 static std::string pwd = gSystem->pwd() + std::string(
"/");
342 currentFilename = pwd + std::string(
f->GetName());
344 outputs.make<std::string>(
o2) = currentFilename;
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) {
352 if (routeIndex == requestedTables.size()) {
353 return TFReaderState::TIMEFRAME_READ;
356 auto& route = requestedTables[routeIndex];
359 bool wasAOD = std::ranges::any_of(route.matcher.metadata, [](
ConfigParamSpec const& p) { return p.name.starts_with(
"aod-origin-replaced"); });
361 if (currentState == TFReaderState::READ_FIRST_TABLE || currentState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE) {
362 didir->beginTimeFrame();
365 if (!didir->readTree(outputs, dh, fcnt, ntf, totalSizeCompressed, totalSizeUncompressed, wasAOD)) {
366 return TFReaderState::TRY_NEXT_FILE;
369 if (!skipInvalidReads) {
372 skipInvalidRead(concrete, e);
373 return TFReaderState::INVALID_TIMEFRAME;
376 if (currentState == TFReaderState::READ_FIRST_TABLE || currentState == TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE) {
380 return TFReaderState::READ_NEXT_TABLE;
382 while (readState != TFReaderState::TIMEFRAME_READ) {
384 case TFReaderState::READ_FIRST_TABLE:
385 transitionTo(tryReadTable(readState));
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) {
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!");
397 case TFReaderState::TRY_NEXT_FILE:
398 didir->finishTimeFrame();
400 if (didir->atEnd(fcnt)) {
401 LOGP(info,
"No input files left to read for reader {}!", device.
inputTimesliceId);
402 didir->closeInputFiles();
403 monitoring.flushBuffer();
410 transitionTo(TFReaderState::READ_FIRST_TABLE_FROM_NEXT_FILE);
412 case TFReaderState::INVALID_TIMEFRAME:
414 case TFReaderState::TIMEFRAME_READ:
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;
423 if (ceil(maxRate) > 0.) {
424 float extraTime = (bytesDelta / 1000000 - currentDelta * maxRate) / maxRate;
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);
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));
448 auto firstRoute = std::ranges::find_if(requestedTables, [&didir,
level](
auto const& route) {
450 return didir->getLevelForOrigin(concrete.
origin) ==
level;
454 auto fileAndFolder = didir->getFileFolder(dh, fcnt, ntf);
457 if (fileAndFolder.filesystem() ==
nullptr) {
460 if (didir->atEnd(fcnt)) {
461 LOGP(info,
"No input files left to read for reader {}!", device.
inputTimesliceId);
462 didir->closeInputFiles();
463 monitoring.flushBuffer();