Skip to content
Merged
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
1 change: 1 addition & 0 deletions Framework/Core/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,7 @@ o2_add_library(Framework
src/Array2D.cxx
src/Variant.cxx
src/VariantJSONHelpers.cxx
src/ExpressionJSONHelpers.cxx
src/VariantPropertyTreeHelpers.cxx
src/WorkflowCustomizationHelpers.cxx
src/WorkflowHelpers.cxx
Expand Down
5 changes: 1 addition & 4 deletions Framework/Core/include/Framework/AODReaderHelpers.h
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,7 @@
#ifndef O2_FRAMEWORK_AODREADERHELPERS_H_
#define O2_FRAMEWORK_AODREADERHELPERS_H_

#include "Framework/TableBuilder.h"
#include "Framework/AlgorithmSpec.h"
#include "Framework/Logger.h"
#include "Framework/RootMessageContext.h"
#include <uv.h>

namespace o2::framework::readers
Expand All @@ -24,7 +21,7 @@ namespace o2::framework::readers

struct AODReaderHelpers {
static AlgorithmSpec rootFileReaderCallback();
static AlgorithmSpec aodSpawnerCallback(std::vector<InputSpec>& requested);
static AlgorithmSpec aodSpawnerCallback(ConfigContext const& ctx);
static AlgorithmSpec indexBuilderCallback(std::vector<InputSpec>& requested);
};

Expand Down
10 changes: 9 additions & 1 deletion Framework/Core/include/Framework/ASoA.h
Original file line number Diff line number Diff line change
Expand Up @@ -1270,6 +1270,7 @@ struct TableIterator : IP, C... {

struct ArrowHelpers {
static std::shared_ptr<arrow::Table> joinTables(std::vector<std::shared_ptr<arrow::Table>>&& tables, std::span<const char* const> labels);
static std::shared_ptr<arrow::Table> joinTables(std::vector<std::shared_ptr<arrow::Table>>&& tables, std::span<const std::string> labels);
static std::shared_ptr<arrow::Table> concatTables(std::vector<std::shared_ptr<arrow::Table>>&& tables);
};

Expand All @@ -1293,7 +1294,14 @@ concept with_ccdb_urls = requires {
};

template <typename T>
concept with_base_table = not_void<typename aod::MetadataTrait<o2::aod::Hash<T::ref.desc_hash>>::metadata::base_table_t>;
concept with_base_table = requires {
typename aod::MetadataTrait<o2::aod::Hash<T::ref.desc_hash>>::metadata::base_table_t;
};

template <typename T>
concept with_expression_pack = requires {
typename T::expression_pack_t{};
};

template <size_t N1, std::array<TableRef, N1> os1, size_t N2, std::array<TableRef, N2> os2>
consteval bool is_compatible()
Expand Down
34 changes: 34 additions & 0 deletions Framework/Core/include/Framework/AnalysisHelpers.h
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,12 @@
#include "Framework/Traits.h"

#include <string>
namespace o2::framework
{
std::string serializeProjectors(std::vector<framework::expressions::Projector>& projectors);
std::string serializeSchema(std::shared_ptr<arrow::Schema>& schema);
} // namespace o2::framework

namespace o2::soa
{
template <TableRef R>
Expand Down Expand Up @@ -97,6 +103,32 @@ constexpr auto getCCDBMetadata() -> std::vector<framework::ConfigParamSpec>
{
return {};
}

template <soa::with_expression_pack T>
constexpr auto getExpressionMetadata() -> std::vector<framework::ConfigParamSpec>
{
using expression_pack_t = T::expression_pack_t;

auto projectors = []<typename... C>(framework::pack<C...>) -> std::vector<framework::expressions::Projector> {
std::vector<framework::expressions::Projector> result;
(result.emplace_back(std::move(C::Projector())), ...);
return result;
}(expression_pack_t{});

auto schema = std::make_shared<arrow::Schema>(o2::soa::createFieldsFromColumns(expression_pack_t{}));

auto json = framework::serializeProjectors(projectors);
return {framework::ConfigParamSpec{"projectors", framework::VariantType::String, json, {"\"\""}},
framework::ConfigParamSpec{"schema", framework::VariantType::String, framework::serializeSchema(schema), {"\"\""}}};
}

template <typename T>
requires(!soa::with_expression_pack<T>)
constexpr auto getExpressionMetadata() -> std::vector<framework::ConfigParamSpec>
{
return {};
}

} // namespace

template <TableRef R>
Expand All @@ -107,6 +139,8 @@ constexpr auto tableRef2InputSpec()
metadata.insert(metadata.end(), m.begin(), m.end());
auto ccdbMetadata = getCCDBMetadata<typename o2::aod::MetadataTrait<o2::aod::Hash<R.desc_hash>>::metadata>();
metadata.insert(metadata.end(), ccdbMetadata.begin(), ccdbMetadata.end());
auto p = getExpressionMetadata<typename o2::aod::MetadataTrait<o2::aod::Hash<R.desc_hash>>::metadata>();
metadata.insert(metadata.end(), p.begin(), p.end());

return framework::InputSpec{
o2::aod::label<R>(),
Expand Down
16 changes: 14 additions & 2 deletions Framework/Core/include/Framework/Expressions.h
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,8 @@ std::string upcastTo(atype::type f);

/// An expression tree node corresponding to a literal value
struct LiteralNode {
using var_t = LiteralValue::stored_type;

LiteralNode()
: value{-1},
type{atype::INT32}
Expand All @@ -120,7 +122,12 @@ struct LiteralNode {
{
}

using var_t = LiteralValue::stored_type;
LiteralNode(var_t v, atype::type t)
: value{v},
type{t}
{
}

var_t value;
atype::type type = atype::NA;
};
Expand Down Expand Up @@ -617,14 +624,19 @@ inline Node ncfg(T defaultValue, std::string path)
struct Filter {
Filter() = default;

Filter(std::unique_ptr<Node>&& ptr)
{
node = std::move(ptr);
(void)designateSubtrees(node.get());
}

Filter(Node&& node_) : node{std::make_unique<Node>(std::forward<Node>(node_))}
{
(void)designateSubtrees(node.get());
}

Filter(Filter&& other) : node{std::forward<std::unique_ptr<Node>>(other.node)}
{
(void)designateSubtrees(node.get());
}

Filter(std::string const& input_) : input{input_} {}
Expand Down
5 changes: 4 additions & 1 deletion Framework/Core/include/Framework/TableBuilder.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@
#include "arrow/type_traits.h"

// Apparently needs to be on top of the arrow includes.
#include <sstream>

#include <arrow/chunked_array.h>
#include <arrow/status.h>
Expand Down Expand Up @@ -796,6 +795,10 @@ auto makeEmptyTable(const char* name, framework::pack<Cs...> p)
std::shared_ptr<arrow::Table> spawnerHelper(std::shared_ptr<arrow::Table> const& fullTable, std::shared_ptr<arrow::Schema> newSchema, size_t nColumns,
expressions::Projector* projectors, const char* name, std::shared_ptr<gandiva::Projector>& projector);

std::shared_ptr<arrow::Table> spawnerHelper(std::shared_ptr<arrow::Table> const& fullTable, std::shared_ptr<arrow::Schema> newSchema,
const char* name, size_t nColumns,
const std::shared_ptr<gandiva::Projector>& projector);

/// Expression-based column generator to materialize columns
template <aod::is_aod_hash D>
requires(soa::has_configurable_extension<typename o2::aod::MetadataTrait<D>::metadata>)
Expand Down
197 changes: 131 additions & 66 deletions Framework/Core/src/AODReaderHelpers.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,16 @@
#include "Framework/AODReaderHelpers.h"
#include "Framework/AnalysisHelpers.h"
#include "Framework/AnalysisDataModelHelpers.h"
#include "Framework/DataProcessingHelpers.h"
#include "Framework/ExpressionHelpers.h"
#include "Framework/DataProcessingHelpers.h"
#include "Framework/AlgorithmSpec.h"
#include "Framework/ControlService.h"
#include "Framework/CallbackService.h"
#include "Framework/EndOfStreamContext.h"
#include "Framework/DataSpecUtils.h"
#include "ExpressionJSONHelpers.h"
#include "Framework/ConfigContext.h"
#include "Framework/AnalysisContext.h"

#include <Monitoring/Monitoring.h>

Expand All @@ -44,28 +47,6 @@ auto setEOSCallback(InitContext& ic)
});
}

template <typename... Ts>
static inline auto doExtractOriginal(framework::pack<Ts...>, ProcessingContext& pc)
{
if constexpr (sizeof...(Ts) == 1) {
return pc.inputs().get<TableConsumer>(aod::MetadataTrait<framework::pack_element_t<0, framework::pack<Ts...>>>::metadata::tableLabel())->asArrowTable();
} else {
return std::vector{pc.inputs().get<TableConsumer>(aod::MetadataTrait<Ts>::metadata::tableLabel())->asArrowTable()...};
}
}

template <typename... Os>
static inline auto extractOriginalsTuple(framework::pack<Os...>, ProcessingContext& pc)
{
return std::make_tuple(extractTypedOriginal<Os>(pc)...);
}

template <typename... Os>
static inline auto extractOriginalsVector(framework::pack<Os...>, ProcessingContext& pc)
{
return std::vector{extractOriginal<Os>(pc)...};
}

template <size_t N, std::array<soa::TableRef, N> refs>
static inline auto extractOriginals(ProcessingContext& pc)
{
Expand Down Expand Up @@ -156,53 +137,137 @@ auto make_spawn(InputSpec const& input, ProcessingContext& pc)
(typename metadata_t::expression_pack_t{});
return o2::framework::spawner<D>(extractOriginals<sources.size(), sources>(pc), input.binding.c_str(), projectors.data(), projector, schema);
}

struct Maker {
std::string binding;
std::vector<std::string> labels;
std::vector<std::shared_ptr<gandiva::Expression>> expressions;
std::shared_ptr<gandiva::Projector> projector = nullptr;
std::shared_ptr<arrow::Schema> schema;

header::DataOrigin origin;
header::DataDescription description;
header::DataHeader::SubSpecificationType version;

std::shared_ptr<arrow::Table> make(ProcessingContext& pc)
{
std::vector<std::shared_ptr<arrow::Table>> originals;
for (auto const& label : labels) {
originals.push_back(pc.inputs().get<TableConsumer>(label)->asArrowTable());
}
auto fullTable = soa::ArrowHelpers::joinTables(std::move(originals), std::span{labels.begin(), labels.size()});
if (projector == nullptr) {
auto s = gandiva::Projector::Make(
fullTable->schema(),
expressions,
&projector);
if (!s.ok()) {
throw o2::framework::runtime_error_f("Failed to create projector: %s", s.ToString().c_str());
}
}

return spawnerHelper(fullTable, schema, binding.c_str(), schema->num_fields(), projector);
}
};

struct Spawnable {
std::string binding;
std::vector<std::string> labels;
std::vector<expressions::Projector> projectors;
std::vector<std::shared_ptr<gandiva::Expression>> expressions;
std::shared_ptr<arrow::Schema> outputSchema;
std::shared_ptr<arrow::Schema> inputSchema;

header::DataOrigin origin;
header::DataDescription description;
header::DataHeader::SubSpecificationType version;

Spawnable(InputSpec const& spec)
: binding{spec.binding}
{
auto&& [origin_, description_, version_] = DataSpecUtils::asConcreteDataMatcher(spec);
origin = origin_;
description = description_;
version = version_;
auto loc = std::find_if(spec.metadata.begin(), spec.metadata.end(), [](ConfigParamSpec const& cps) { return cps.name.compare("projectors") == 0; });
std::stringstream iws(loc->defaultValue.get<std::string>());
projectors = ExpressionJSONHelpers::read(iws);

loc = std::find_if(spec.metadata.begin(), spec.metadata.end(), [](ConfigParamSpec const& cps) { return cps.name.compare("schema") == 0; });
iws.clear();
iws.str(loc->defaultValue.get<std::string>());
outputSchema = ArrowJSONHelpers::read(iws);

for (auto& i : spec.metadata) {
if (i.name.starts_with("input:")) {
labels.emplace_back(i.name.substr(6));
}
}

std::vector<std::shared_ptr<arrow::Field>> fields;
for (auto& p : projectors) {
expressions::walk(p.node.get(),
[&fields](expressions::Node* n) mutable {
if (n->self.index() == 1) {
auto& b = std::get<expressions::BindingNode>(n->self);
if (std::find_if(fields.begin(), fields.end(), [&b](std::shared_ptr<arrow::Field> const& field) { return field->name() == b.name; }) == fields.end()) {
fields.emplace_back(std::make_shared<arrow::Field>(b.name, expressions::concreteArrowType(b.type)));
}
}
});
}
inputSchema = std::make_shared<arrow::Schema>(fields);

int i = 0;
for (auto& p : projectors) {
expressions.push_back(
expressions::makeExpression(
expressions::createExpressionTree(
expressions::createOperations(p),
inputSchema),
outputSchema->field(i)));
++i;
}
}

std::shared_ptr<gandiva::Projector> makeProjector()
{
return expressions::createProjectorHelper(projectors.size(), projectors.data(), inputSchema, outputSchema->fields());
}

Maker createMaker()
{
return {
binding,
labels,
expressions,
nullptr,
outputSchema,
origin,
description,
version};
}
};

} // namespace

AlgorithmSpec AODReaderHelpers::aodSpawnerCallback(std::vector<InputSpec>& requested)
AlgorithmSpec AODReaderHelpers::aodSpawnerCallback(/*std::vector<InputSpec>& requested*/ ConfigContext const& ctx)
{
return AlgorithmSpec::InitCallback{[requested](InitContext& /*ic*/) {
return [requested](ProcessingContext& pc) {
auto& ac = ctx.services().get<AnalysisContext>();
return AlgorithmSpec::InitCallback{[requested = ac.spawnerInputs](InitContext& /*ic*/) {
std::vector<Spawnable> spawnables;
for (auto& i : requested) {
spawnables.emplace_back(i);
}
std::vector<Maker> makers;
for (auto& s : spawnables) {
makers.push_back(s.createMaker());
}

return [makers](ProcessingContext& pc) mutable {
auto outputs = pc.outputs();
// spawn tables
for (auto& input : requested) {
auto&& [origin, description, version] = DataSpecUtils::asConcreteDataMatcher(input);
if (description == header::DataDescription{"EXTRACK"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACK/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXTRACK_IU"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACK_IU/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXTRACKCOV"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACKCOV/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXTRACKCOV_IU"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACKCOV_IU/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXTRACKEXTRA"}) {
if (version == 0U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACKEXTRA/0"_h>>(input, pc));
} else if (version == 1U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACKEXTRA/1"_h>>(input, pc));
} else if (version == 2U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXTRACKEXTRA/2"_h>>(input, pc));
}
} else if (description == header::DataDescription{"EXMFTTRACK"}) {
if (version == 0U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXMFTTRACK/0"_h>>(input, pc));
} else if (version == 1U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXMFTTRACK/1"_h>>(input, pc));
}
} else if (description == header::DataDescription{"EXMFTTRACKCOV"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXMFTTRACKCOV/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXFWDTRACK"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXFWDTRACK/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXFWDTRACKCOV"}) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXFWDTRACKCOV/0"_h>>(input, pc));
} else if (description == header::DataDescription{"EXMCPARTICLE"}) {
if (version == 0U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXMCPARTICLE/0"_h>>(input, pc));
} else if (version == 1U) {
outputs.adopt(Output{origin, description, version}, make_spawn<o2::aod::Hash<"EXMCPARTICLE/1"_h>>(input, pc));
}
} else {
throw runtime_error("Not an extended table");
}
for (auto& maker : makers) {
outputs.adopt(Output{maker.origin, maker.description, maker.version}, maker.make(pc));
}
};
}};
Expand Down
Loading