Skip to content

Commit b554115

Browse files
committed
DPL Analysis: use shm metadata in CCDB tables instead of binary view
1 parent f2a59e7 commit b554115

8 files changed

Lines changed: 242 additions & 19 deletions

File tree

Framework/CCDBSupport/src/AnalysisCCDBHelpers.cxx

Lines changed: 55 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111

1212
#include "AnalysisCCDBHelpers.h"
1313
#include "CCDBFetcherHelper.h"
14+
#include "Framework/ArrowTypes.h"
1415
#include "Framework/DataProcessingStats.h"
1516
#include "Framework/DeviceSpec.h"
1617
#include "Framework/TimingInfo.h"
@@ -22,6 +23,7 @@
2223
#include "Framework/DanglingEdgesContext.h"
2324
#include "Framework/ConfigContext.h"
2425
#include "Framework/ConfigParamsHelper.h"
26+
#include <fairmq/Version.h>
2527
#include <arrow/array/builder_binary.h>
2628
#include <arrow/type.h>
2729
#include <arrow/type_fwd.h>
@@ -109,20 +111,45 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
109111
auto it = ccdbUrls.find(m.name);
110112
fieldMetadata->Append("url", it != ccdbUrls.end() ? it->second : m.defaultValue.asString());
111113
auto columnName = m.name.substr(strlen("ccdb:"));
114+
#if (FAIRMQ_VERSION_DEC >= 111000)
115+
fields.emplace_back(std::make_shared<arrow::Field>(columnName, soa::asArrowDataType<int64_t[3]>(), false, fieldMetadata));
116+
#else
112117
fields.emplace_back(std::make_shared<arrow::Field>(columnName, arrow::binary_view(), false, fieldMetadata));
118+
#endif
113119
}
114120
schemas.emplace_back(std::make_shared<arrow::Schema>(fields, schemaMetadata));
115121
}
116122

123+
std::vector<std::pair<uint32_t, std::shared_ptr<arrow::FixedSizeListBuilder>>> allbuilders;
124+
allbuilders.resize([&schemas]() { size_t size = 0; for (auto& schema : schemas) { size += schema->num_fields(); }; return size; }());
125+
auto* pool = arrow::default_memory_pool();
126+
127+
int idx = 0;
128+
int sidx = 0;
129+
for (auto const& schema : schemas) {
130+
for (auto const& _ : schema->fields()) {
131+
#if (FAIRMQ_VERSION_DEC >= 111000)
132+
auto value_builder = std::make_shared<arrow::Int64Builder>();
133+
allbuilders[idx] = std::make_pair(sidx, std::make_shared<arrow::FixedSizeListBuilder>(pool, std::move(value_builder), 3));
134+
#else
135+
allbuilders[idx] = std::make_pair(sidx, std::make_shared<arrow::BinaryViewBuilder>());
136+
#endif
137+
++idx;
138+
}
139+
++sidx;
140+
}
141+
117142
std::shared_ptr<CCDBFetcherHelper> helper = std::make_shared<CCDBFetcherHelper>();
118143
CCDBFetcherHelper::initialiseHelper(*helper, options);
119144
std::unordered_map<std::string, int> bindings;
120145
fillValidRoutes(*helper, spec.outputs, bindings);
121146

122-
return adaptStateless([schemas, bindings, helper](InputRecord& inputs, DataTakingContext& dtc, DataAllocator& allocator, TimingInfo& timingInfo, DataProcessingStats& stats) {
147+
return adaptStateless([schemas, bindings, helper, allbuilders](InputRecord& inputs, DataTakingContext& dtc, DataAllocator& allocator, TimingInfo& timingInfo, DataProcessingStats& stats) {
123148
O2_SIGNPOST_ID_GENERATE(sid, ccdb);
124149
O2_SIGNPOST_START(ccdb, sid, "fetchFromAnalysisCCDB", "Fetching CCDB objects for analysis%" PRIu64, (uint64_t)timingInfo.timeslice);
125-
for (auto& schema : schemas) {
150+
std::ranges::for_each(allbuilders, [](auto& builder) { builder.second->Reset(); });
151+
for (auto i = 0U; i < schemas.size(); ++i) {
152+
auto& schema = schemas[i];
126153
std::vector<CCDBFetcherHelper::FetchOp> ops;
127154
auto inputBinding = *schema->metadata()->Get("sourceTable");
128155
auto inputMatcher = DataSpecUtils::fromString(*schema->metadata()->Get("sourceMatcher"));
@@ -134,6 +161,7 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
134161
auto table = inputs.get<TableConsumer>(inputMatcher)->asArrowTable();
135162
// FIXME: make the fTimestamp column configurable.
136163
auto timestampColumn = table->GetColumnByName("fTimestamp");
164+
auto reserveSize = timestampColumn->length();
137165
O2_SIGNPOST_EVENT_EMIT_INFO(ccdb, sid, "fetchFromAnalysisCCDB",
138166
"There are %zu bindings available", bindings.size());
139167
for (auto& binding : bindings) {
@@ -143,9 +171,16 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
143171
}
144172
int outputRouteIndex = bindings.at(outRouteDesc);
145173
auto& spec = helper->routes[outputRouteIndex].matcher;
146-
std::vector<std::shared_ptr<arrow::BinaryViewBuilder>> builders;
147-
for (auto const& _ : schema->fields()) {
148-
builders.emplace_back(std::make_shared<arrow::BinaryViewBuilder>());
174+
auto builders = allbuilders | std::views::filter([&i](auto const& builder) { return builder.first == i; });
175+
unsigned int numBuilders = std::ranges::count_if(allbuilders, [&i](auto const& builder) { return builder.first == i; });
176+
arrow::Status status;
177+
std::ranges::for_each(builders, [&status, &reserveSize](auto& builder) {
178+
if (reserveSize > builder.second->capacity()) {
179+
status &= builder.second->Reserve(reserveSize - builder.second->capacity());
180+
}
181+
});
182+
if (!status.ok()) {
183+
throw framework::runtime_error_f("Failed to reserve arrays: ", status.ToString().c_str());
149184
}
150185

151186
for (auto ci = 0; ci < timestampColumn->num_chunks(); ++ci) {
@@ -171,15 +206,25 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
171206
O2_SIGNPOST_START(ccdb, sid, "handlingResponses",
172207
"Got %zu responses from server.",
173208
responses.size());
174-
if (builders.size() != responses.size()) {
175-
LOGP(fatal, "Not enough responses (expected {}, found {})", builders.size(), responses.size());
209+
if (numBuilders != responses.size()) {
210+
LOGP(fatal, "Not enough responses (expected {}, found {})", numBuilders, responses.size());
176211
}
177212
arrow::Status result;
178-
for (size_t bi = 0; bi < responses.size(); bi++) {
179-
auto& builder = builders[bi];
213+
214+
int bi = 0;
215+
for (auto& builder : builders) {
180216
auto& response = responses[bi];
217+
#if (FAIRMQ_VERSION_DEC >= 111000)
218+
result &= builder.second->Append();
219+
auto* value_builder = dynamic_cast<arrow::Int64Builder*>(builder.second->value_builder());
220+
result &= value_builder->Append(response.id.handle);
221+
result &= value_builder->Append(response.id.segment);
222+
result &= value_builder->Append(response.size);
223+
#else
181224
char const* address = reinterpret_cast<char const*>(response.id.value);
182225
result &= builder->Append(std::string_view(address, response.size));
226+
#endif
227+
++bi;
183228
}
184229
if (!result.ok()) {
185230
LOGP(fatal, "Error adding results from CCDB");
@@ -188,9 +233,7 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
188233
}
189234
}
190235
arrow::ArrayVector arrays;
191-
for (auto& builder : builders) {
192-
arrays.push_back(*builder->Finish());
193-
}
236+
std::ranges::for_each(builders, [&arrays](auto& builder) { arrays.push_back(*builder.second->Finish()); });
194237
auto outTable = arrow::Table::Make(schema, arrays);
195238
auto concrete = DataSpecUtils::asConcreteDataMatcher(spec);
196239
allocator.adopt(Output{concrete.origin, concrete.description, concrete.subSpec}, outTable);

Framework/Core/include/Framework/ASoA.h

Lines changed: 89 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@
2828
#include "Framework/ArrowTableSlicingCache.h" // IWYU pragma: export
2929
#include "Framework/SliceCache.h" // IWYU pragma: export
3030
#include "Framework/VariantHelpers.h" // IWYU pragma: export
31+
#include <fairmq/Version.h>
3132
#include <arrow/array/array_binary.h>
3233
#include <arrow/table.h> // IWYU pragma: export
3334
#include <arrow/array.h> // IWYU pragma: export
@@ -40,9 +41,16 @@
4041
#include <cstring>
4142
#include <gsl/span> // IWYU pragma: export
4243

44+
namespace fair::mq::shmem {
45+
struct MetaHeader;
46+
}
47+
4348
namespace o2::framework
4449
{
4550
using ListVector = std::vector<std::vector<int64_t>>;
51+
#if (FAIRMQ_VERSION_DEC >= 111000)
52+
using PointerReconstructor = std::function<std::byte*(fair::mq::shmem::MetaHeader&&)>;
53+
#endif
4654

4755
std::string cutString(std::string&& str);
4856
std::string strToUpper(std::string&& str);
@@ -1085,6 +1093,9 @@ concept can_bind = requires(T&& t) {
10851093
template <typename... C>
10861094
concept has_index = (is_indexing_column<C> || ...);
10871095

1096+
template <typename C>
1097+
concept needs_ptr_rec = C::needs_ptr_rec;
1098+
10881099
template <typename D, typename O, typename IP, typename... C>
10891100
struct TableIterator : IP, C... {
10901101
public:
@@ -1252,6 +1263,21 @@ struct TableIterator : IP, C... {
12521263
{
12531264
doSetCurrentInternal(internal_index_columns_t{}, table);
12541265
}
1266+
#if (FAIRMQ_VERSION_DEC >= 111000)
1267+
void setPointerReconstructor(framework::PointerReconstructor const& pointerReconstructor)
1268+
{
1269+
[&pointerReconstructor, this]<typename... Cs>(framework::pack<Cs...>) {
1270+
([&pointerReconstructor, this]<typename CC>() {
1271+
if constexpr (needs_ptr_rec<CC>) {
1272+
if (pointerReconstructor) {
1273+
CC::ptrRec = &pointerReconstructor;
1274+
}
1275+
}
1276+
}.template operator()<Cs>(),
1277+
...);
1278+
}(all_columns{});
1279+
}
1280+
#endif
12551281

12561282
private:
12571283
/// Helper to move at the end of columns which actually have an iterator.
@@ -2289,7 +2315,12 @@ class Table
22892315
{
22902316
return self_t{mTable->Slice(0, 0), 0};
22912317
}
2292-
2318+
#if (FAIRMQ_VERSION_DEC >= 111000)
2319+
void setPointerReconstructor(framework::PointerReconstructor const& pointerReconstructor)
2320+
{
2321+
mBegin.setPointerReconstructor(pointerReconstructor);
2322+
}
2323+
#endif
22932324
private:
22942325
template <typename T>
22952326
arrow::ChunkedArray* lookupColumn()
@@ -2468,6 +2499,56 @@ consteval static std::string_view namespace_prefix()
24682499
}; \
24692500
[[maybe_unused]] static constexpr o2::framework::expressions::BindingNode _Getter_ { _Label_, _Name_::hash, o2::framework::expressions::selectArrowType<_Type_>() }
24702501

2502+
#if (FAIRMQ_VERSION_DEC >= 111000)
2503+
#define DECLARE_SOA_CCDB_COLUMN_FULL(_Name_, _Label_, _Getter_, _ConcreteType_, _CCDBQuery_) \
2504+
struct _Name_ : o2::soa::Column<int64_t[3], _Name_> { \
2505+
static constexpr const char* mLabel = _Label_; \
2506+
static constexpr const char* query = _CCDBQuery_; \
2507+
static constexpr const uint32_t hash = crc32(namespace_prefix<_Name_>(), std::string_view{#_Getter_}); \
2508+
static constexpr bool needs_ptr_rec = true; \
2509+
std::function<std::byte*(fair::mq::shmem::MetaHeader&&)> const* ptrRec = nullptr; \
2510+
using base = o2::soa::Column<int64_t[3], _Name_>; \
2511+
using type = int64_t[3]; \
2512+
using column_t = _Name_; \
2513+
_Name_(arrow::ChunkedArray const* column) \
2514+
: o2::soa::Column<int64_t[3], _Name_>(o2::soa::ColumnIterator<int64_t[3]>(column)) \
2515+
{ \
2516+
} \
2517+
\
2518+
_Name_() = default; \
2519+
_Name_(_Name_ const& other) = default; \
2520+
_Name_& operator=(_Name_ const& other) = default; \
2521+
\
2522+
decltype(auto) _Getter_() const \
2523+
{ \
2524+
auto& [handle, segment, size] = *mColumnIterator; \
2525+
auto span = std::span<std::byte>{(*ptrRec)(fair::mq::shmem::MetaHeader{ \
2526+
static_cast<size_t>(size), \
2527+
0, handle, 0, 0, \
2528+
static_cast<uint16_t>(segment), true}), static_cast<size_t>(size)}; \
2529+
if constexpr (std::same_as<_ConcreteType_, std::span<std::byte>>) { \
2530+
return span; \
2531+
} else { \
2532+
static std::byte* payload = nullptr; \
2533+
static _ConcreteType_* deserialised = nullptr; \
2534+
static TClass* c = TClass::GetClass(#_ConcreteType_); \
2535+
if (payload != (std::byte*)span.data()) { \
2536+
payload = (std::byte*)span.data(); \
2537+
delete deserialised; \
2538+
TBufferFile f(TBufferFile::EMode::kRead, span.size(), (char*)span.data(), kFALSE); \
2539+
deserialised = (_ConcreteType_*)soa::extractCCDBPayload((char*)payload, span.size(), c, "ccdb_object"); \
2540+
} \
2541+
return *deserialised; \
2542+
} \
2543+
} \
2544+
\
2545+
decltype(auto) \
2546+
get() const \
2547+
{ \
2548+
return _Getter_(); \
2549+
} \
2550+
};
2551+
#else
24712552
#define DECLARE_SOA_CCDB_COLUMN_FULL(_Name_, _Label_, _Getter_, _ConcreteType_, _CCDBQuery_) \
24722553
struct _Name_ : o2::soa::Column<std::span<std::byte>, _Name_> { \
24732554
static constexpr const char* mLabel = _Label_; \
@@ -2510,6 +2591,7 @@ consteval static std::string_view namespace_prefix()
25102591
return _Getter_(); \
25112592
} \
25122593
};
2594+
#endif
25132595

25142596
#define DECLARE_SOA_CCDB_COLUMN(_Name_, _Getter_, _ConcreteType_, _CCDBQuery_) \
25152597
DECLARE_SOA_CCDB_COLUMN_FULL(_Name_, "f" #_Name_, _Getter_, _ConcreteType_, _CCDBQuery_)
@@ -3853,7 +3935,12 @@ class FilteredBase : public T
38533935
{
38543936
return mCached;
38553937
}
3856-
3938+
#if (FAIRMQ_VERSION_DEC >= 111000)
3939+
void setPointerReconstructor(framework::PointerReconstructor const& pointerReconstructor)
3940+
{
3941+
mFilteredBegin.setPointerReconstructor(pointerReconstructor);
3942+
}
3943+
#endif
38573944
private:
38583945
void resetRanges()
38593946
{

Framework/Core/include/Framework/AnalysisDataModel.h

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,11 @@
2626
#include "SimulationDataFormat/MCGenProperties.h"
2727
#include "Framework/PID.h"
2828

29+
#include <fairmq/Version.h>
30+
#if (FAIRMQ_VERSION_DEC >= 111000)
31+
#include <fairmq/shmem/Common.h>
32+
#endif
33+
2934
namespace o2
3035
{
3136
namespace aod

0 commit comments

Comments
 (0)