Skip to content
Prev Previous commit
Next Next commit
Add internal slicer devices
Add the ArrowTableSlicer algorithm plugin, which builds the requested slice-info tables from the slice sources

The slice inputs are grouped based on the provider of the sources, so
that multiple slicers can be injected right after the providers to avoid
loops in topology
  • Loading branch information
aalkin committed Sep 29, 2026
commit e01d97d5cd3c1f7563fbcf76aa1c708af290a295
1 change: 1 addition & 0 deletions Framework/AnalysisSupport/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ endif()
o2_add_library(FrameworkOnDemandTablesSupport
SOURCES src/OnDemandPlugin.cxx
src/AODReaderHelpers.cxx
src/AODSliceHelpers.cxx
PRIVATE_INCLUDE_DIRECTORIES ${CMAKE_CURRENT_LIST_DIR}/src
PUBLIC_LINK_LIBRARIES O2::Framework ${EXTRA_TARGETS})

Expand Down
66 changes: 66 additions & 0 deletions Framework/AnalysisSupport/src/AODSliceHelpers.cxx
Original file line number Diff line number Diff line change
@@ -0,0 +1,66 @@
// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
// All rights not expressly granted are reserved.
//
// This software is distributed under the terms of the GNU General Public
// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
//
// In applying this license CERN does not waive the privileges and immunities
// granted to it by virtue of its status as an Intergovernmental Organization
// or submit itself to any jurisdiction.

#include "AODSliceHelpers.h"

#include "Framework/ArrowTableSlicingCache.h"
#include "Framework/ConfigParamRegistry.h"
#include "Framework/DanglingEdgesContext.h"
#include "Framework/DataAllocator.h"
#include "Framework/DataSpecUtils.h"
#include "Framework/InputRecord.h"
#include "Framework/TableConsumer.h"

namespace o2::framework::helpers
{
namespace
{
Entry sourceEntry(InputSpec const& spec)
{
auto source = DataSpecUtils::fromMetadataString(std::ranges::find_if(spec.metadata, [](ConfigParamSpec const& cps) { return cps.name.starts_with("slice-source"); })->defaultValue.get<std::string>());
return {source.binding, DataSpecUtils::asConcreteDataMatcher(source), std::ranges::find_if(spec.metadata, [](ConfigParamSpec const& cps) { return cps.name.starts_with("slice-key"); })->defaultValue.get<std::string>()};
}

struct Sliceable {
Entry entry;
ConcreteDataMatcher output;
bool sorted;

explicit Sliceable(InputSpec const& spec)
: entry{sourceEntry(spec)},
output{DataSpecUtils::asConcreteDataMatcher(spec)},
sorted{std::ranges::find_if(spec.metadata, [](ConfigParamSpec const& cps) { return cps.name.starts_with("sorted"); })->defaultValue.get<bool>()}
{
}

std::shared_ptr<arrow::Table> materialize(ProcessingContext& pc) const
{
auto source = pc.inputs().get<TableConsumer>(entry.matcher)->asArrowTable();
return sorted ? SliceInfo::makeSorted(entry, source) : SliceInfo::makeUnsorted(entry, source);
}
};
} // namespace

AlgorithmSpec AODSliceHelpers::arrowTablesSlicerCallback(ConfigContext const& /*ctx*/)
{
return AlgorithmSpec::InitCallback{[](InitContext& ic) {
// each slicer handles the group of slice infos for the tables from a single provider
auto const& requested = ic.services().get<DanglingEdgesContext>().slicerGroups[ic.options().get<int>("slicer-group")];
std::vector<Sliceable> sliceables;
sliceables.reserve(requested.size());
std::ranges::transform(requested, std::back_inserter(sliceables), [](auto const& i) { return Sliceable{i}; });
return [sliceables](ProcessingContext& pc) {
auto outputs = pc.outputs();
std::ranges::for_each(sliceables, [&pc, &outputs](auto const& sliceable) { outputs.adopt(Output{sliceable.output.origin, sliceable.output.description, sliceable.output.subSpec}, sliceable.materialize(pc)); });
};
}};
}
} // namespace o2::framework::helpers
25 changes: 25 additions & 0 deletions Framework/AnalysisSupport/src/AODSliceHelpers.h
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
// Copyright 2019-2020 CERN and copyright holders of ALICE O2.
// See https://alice-o2.web.cern.ch/copyright for details of the copyright holders.
// All rights not expressly granted are reserved.
//
// This software is distributed under the terms of the GNU General Public
// License v3 (GPL Version 3), copied verbatim in the file "COPYING".
//
// In applying this license CERN does not waive the privileges and immunities
// granted to it by virtue of its status as an Intergovernmental Organization
// or submit itself to any jurisdiction.

#ifndef AODSLICEHELPERS_H
#define AODSLICEHELPERS_H

#include "Framework/AlgorithmSpec.h"
namespace o2::framework::helpers
{

struct AODSliceHelpers {
static AlgorithmSpec arrowTablesSlicerCallback(ConfigContext const& /*ctx*/);
};

} // namespace o2::framework::helpers

#endif // AODSLICEHELPERS_H
9 changes: 9 additions & 0 deletions Framework/AnalysisSupport/src/OnDemandPlugin.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@
#include "Framework/Plugins.h"
#include "Framework/AlgorithmSpec.h"
#include "AODReaderHelpers.h"
#include "AODSliceHelpers.h"

struct ExtendedTableSpawner : o2::framework::AlgorithmPlugin {
o2::framework::AlgorithmSpec create(o2::framework::ConfigContext const& config) override
Expand All @@ -26,7 +27,15 @@ struct IndexTableBuilder : o2::framework::AlgorithmPlugin {
}
};

struct ArrowTableSlicer : o2::framework::AlgorithmPlugin {
o2::framework::AlgorithmSpec create(o2::framework::ConfigContext const& config) override
{
return o2::framework::helpers::AODSliceHelpers::arrowTablesSlicerCallback(config);
}
};

DEFINE_DPL_PLUGINS_BEGIN
DEFINE_DPL_PLUGIN_INSTANCE(ExtendedTableSpawner, CustomAlgorithm);
DEFINE_DPL_PLUGIN_INSTANCE(IndexTableBuilder, CustomAlgorithm);
DEFINE_DPL_PLUGIN_INSTANCE(ArrowTableSlicer, CustomAlgorithm);
DEFINE_DPL_PLUGINS_END
9 changes: 9 additions & 0 deletions Framework/Core/include/Framework/AnalysisSupportHelpers.h
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,15 @@ struct AnalysisSupportHelpers {
std::vector<InputSpec>& requestedAODs,
std::vector<InputSpec>& requestedDYNs,
DataProcessorSpec& publisher);
static void addMissingOutputsToSlicer(std::vector<InputSpec> const& requestedSLCs,
DataProcessorSpec& publisher);
/// Split the requested slice infos into groups by the device providing the sliced table
/// (the AOD reader, if none of the providers has it) and create a slicer device for each
/// group, so that each slicer depends on a single device and does not create loops.
/// Each slicer is returned together with the name of its provider.
static std::vector<std::pair<std::string, DataProcessorSpec>> makeSlicers(std::vector<InputSpec> const& requestedSLCs,
std::vector<DataProcessorSpec const*> const& providers,
std::vector<std::vector<InputSpec>>& slicerGroups);

/// Match all inputs of kind ATSK and write them to a ROOT file,
/// one root file per originating task.
Expand Down
4 changes: 4 additions & 0 deletions Framework/Core/include/Framework/DanglingEdgesContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,10 @@ struct DanglingEdgesContext {
// ccdb tables
std::vector<OutputSpec> providedTIMs;
std::vector<InputSpec> requestedTIMs;
// slice infos
std::vector<InputSpec> requestedSLCs;
// slice infos grouped by the device providing the sliced tables, one slicer device per group
std::vector<std::vector<InputSpec>> slicerGroups;
// output objects
std::vector<OutputSpec> providedOutputObjHist;
// inputs for the extension spawner
Expand Down
57 changes: 57 additions & 0 deletions Framework/Core/src/AnalysisSupportHelpers.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -188,6 +188,63 @@ void AnalysisSupportHelpers::addMissingOutputsToBuilder(std::vector<InputSpec> c
sinks::update_input_list{requestedDYNs}; // update requestedDYNs
}

void AnalysisSupportHelpers::addMissingOutputsToSlicer(std::vector<InputSpec> const& requestedSLCs,
DataProcessorSpec& publisher)
{
requestedSLCs |
views::input_to_output_specs() |
sinks::append_to{publisher.outputs};

for (auto const& input : requestedSLCs) {
input.metadata |
views::filter_string_params_starts_with("slice-source:") |
views::params_to_input_specs() |
sinks::update_input_list{publisher.inputs};
}
}

std::vector<std::pair<std::string, DataProcessorSpec>> AnalysisSupportHelpers::makeSlicers(std::vector<InputSpec> const& requestedSLCs,
std::vector<DataProcessorSpec const*> const& providers,
std::vector<std::vector<InputSpec>>& slicerGroups)
{
// find the device providing the sliced table, if there is none the table is read from file
auto providerFor = [&providers](InputSpec const& request) -> std::string {
auto sources = request.metadata |
views::filter_string_params_starts_with("slice-source:") |
views::params_to_input_specs();
auto matcher = DataSpecUtils::asConcreteDataMatcher(*sources.begin());
auto provider = std::ranges::find_if(providers, [&matcher](DataProcessorSpec const* spec) {
return std::ranges::any_of(spec->outputs, [&matcher](OutputSpec const& output) { return DataSpecUtils::match(output, matcher); });
});
return provider != providers.end() ? (*provider)->name : "internal-dpl-aod-reader";
};

slicerGroups.clear();
std::vector<std::string> groupProviders;
for (auto const& request : requestedSLCs) {
auto provider = providerFor(request);
auto locate = std::ranges::find(groupProviders, provider);
if (locate == groupProviders.end()) {
groupProviders.push_back(provider);
slicerGroups.push_back({request});
} else {
slicerGroups[std::distance(groupProviders.begin(), locate)].push_back(request);
}
}

std::vector<std::pair<std::string, DataProcessorSpec>> slicers;
for (auto i = 0u; i < slicerGroups.size(); ++i) {
DataProcessorSpec slicer{.name = "internal-dpl-aod-slicer-" + std::to_string(i),
.inputs = {},
.outputs = {},
.algorithm = AlgorithmSpec::dummyAlgorithm(), // real algorithm will be set in adjustTopology
.options = {ConfigParamSpec{"slicer-group", VariantType::Int, static_cast<int>(i), {"index of the slice info group handled by this slicer"}}}};
addMissingOutputsToSlicer(slicerGroups[i], slicer);
slicers.emplace_back(groupProviders[i], std::move(slicer));
}
return slicers;
}

// =============================================================================
DataProcessorSpec AnalysisSupportHelpers::getOutputObjHistSink(ConfigContext const& ctx)
{
Expand Down
20 changes: 20 additions & 0 deletions Framework/Core/src/ArrowSupport.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -698,6 +698,26 @@ o2::framework::ServiceSpec ArrowSupport::arrowBackendSpec()
}
}

// slicers are recreated from scratch, grouped by the devices providing the sliced tables
std::erase_if(workflow, [](DataProcessorSpec const& spec) { return spec.name.starts_with("internal-dpl-aod-slicer"); });
dec.requestedSLCs.clear();
for (auto& d : workflow) {
d.inputs |
views::filter_with_params_by_name_starting("slice-source:") |
sinks::update_input_list{dec.requestedSLCs};
}
std::ranges::sort(dec.requestedSLCs, inputSpecLessThan);
std::vector<DataProcessorSpec const*> slicedTablesProviders;
std::ranges::transform(workflow, std::back_inserter(slicedTablesProviders), [](DataProcessorSpec const& spec) { return &spec; });
auto slicers = AnalysisSupportHelpers::makeSlicers(dec.requestedSLCs, slicedTablesProviders, dec.slicerGroups);
// the slicers are placed after their providers in a pre-sorted workflow
for (auto& [providerName, slicer] : slicers) {
// load real AlgorithmSpec before deployment
slicer.algorithm = PluginManager::loadAlgorithmFromPlugin("O2FrameworkOnDemandTablesSupport", "ArrowTableSlicer", ctx);
auto provider = std::ranges::find(workflow, providerName, &DataProcessorSpec::name);
workflow.insert(provider == workflow.end() ? workflow.begin() : std::next(provider), std::move(slicer));
}

auto writer = std::ranges::find_if(workflow, [](DataProcessorSpec const& spec) { return spec.name.starts_with("internal-dpl-aod-writer"); });
if (writer != workflow.end()) {
workflow.erase(writer);
Expand Down
25 changes: 23 additions & 2 deletions Framework/Core/src/WorkflowHelpers.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -282,9 +282,11 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
bool hasProjectors = false;
bool hasIndexRecords = false;
bool hasCCDBURLs = false;
bool hasSliceSource = false;
bool wasAOD = false;
// all three options are exclusive
// all options are exclusive
for (auto const& p : input.metadata) {
// wasAOD can be true or false for all of the options
if (p.name.starts_with("aod-origin-replaced")) {
wasAOD = true;
}
Expand All @@ -300,6 +302,10 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
hasCCDBURLs = true;
break;
}
if (p.name.starts_with("slice-source")) {
hasSliceSource = true;
break;
}
}
switch (input.lifetime) {
case Lifetime::Timer: {
Expand Down Expand Up @@ -346,6 +352,8 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
DataSpecUtils::updateInputList(dec.requestedIDXs, InputSpec{input});
} else if (hasCCDBURLs) {
DataSpecUtils::updateInputList(dec.requestedTIMs, InputSpec{input});
} else if (hasSliceSource) {
DataSpecUtils::updateInputList(dec.requestedSLCs, InputSpec{input});
} else if (DataSpecUtils::partialMatch(input, AODOrigins) || wasAOD) {
DataSpecUtils::updateInputList(dec.requestedAODs, InputSpec{input});
}
Expand All @@ -358,8 +366,11 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
bool hasIndexRecords = false;
bool hasCCDBURLs = false;
bool wasAOD = false;
// all three options are exclusive
// all options are exclusive
// provided slice outputs are ignored, they can only come from the slicer device
// that will be re-added in adjust topology
for (auto const& p : output.metadata) {
// wasAOD can be true or false for all of the options
if (p.name.starts_with("aod-origin-replaced")) {
wasAOD = true;
}
Expand Down Expand Up @@ -439,6 +450,13 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
std::ranges::sort(providedCCDBs, outputSpecLessThan);
AnalysisSupportHelpers::addMissingOutputsToReader(providedCCDBs, requestedCCDBs, ccdbBackend);

// slicers are grouped by the devices providing the sliced tables
std::ranges::sort(dec.requestedSLCs, inputSpecLessThan);
std::vector<DataProcessorSpec const*> slicedTablesProviders;
std::ranges::transform(workflow, std::back_inserter(slicedTablesProviders), [](DataProcessorSpec const& spec) { return &spec; });
slicedTablesProviders.insert(slicedTablesProviders.end(), {&aodSpawner, &indexBuilder, &aodReader});
auto aodSlicers = AnalysisSupportHelpers::makeSlicers(dec.requestedSLCs, slicedTablesProviders, dec.slicerGroups);

std::vector<DataProcessorSpec> extraSpecs;

if (transientStore.outputs.empty() == false) {
Expand All @@ -456,6 +474,9 @@ void WorkflowHelpers::injectServiceDevices(WorkflowSpec& workflow, ConfigContext
extraSpecs.push_back(indexBuilder);
}

// here the slicers are just added, unlike in adjustTopology
std::ranges::transform(aodSlicers, std::back_inserter(extraSpecs), [](auto&& pair){ return pair.second; });

// add the Analysys CCDB backend which reads CCDB objects using a provided table
DeploymentMode deploymentMode = DefaultsHelpers::deploymentMode();
if (deploymentMode != DeploymentMode::OnlineDDS && deploymentMode != DeploymentMode::OnlineECS) {
Expand Down