Skip to content

Commit a419cc7

Browse files
committed
DPL Analysis: use shm metadata in CCDB tables instead of binary view
1 parent 335e4b5 commit a419cc7

8 files changed

Lines changed: 354 additions & 117 deletions

File tree

Framework/CCDBSupport/src/AnalysisCCDBHelpers.cxx

Lines changed: 60 additions & 13 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,49 @@ 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+
#if (FAIRMQ_VERSION_DEC >= 111000)
124+
std::vector<std::pair<uint32_t, std::shared_ptr<arrow::FixedSizeListBuilder>>> allbuilders;
125+
#else
126+
std::vector<std::pair<uint32_t, std::shared_ptr<arrow::BinaryViewBuilder>>> allbuilders;
127+
#endif
128+
allbuilders.resize([&schemas]() { size_t size = 0; for (auto& schema : schemas) { size += schema->num_fields(); }; return size; }());
129+
auto* pool = arrow::default_memory_pool();
130+
131+
int idx = 0;
132+
int sidx = 0;
133+
for (auto const& schema : schemas) {
134+
for (auto const& _ : schema->fields()) {
135+
#if (FAIRMQ_VERSION_DEC >= 111000)
136+
auto value_builder = std::make_shared<arrow::Int64Builder>();
137+
allbuilders[idx] = std::make_pair(sidx, std::make_shared<arrow::FixedSizeListBuilder>(pool, std::move(value_builder), 3));
138+
#else
139+
allbuilders[idx] = std::make_pair(sidx, std::make_shared<arrow::BinaryViewBuilder>());
140+
#endif
141+
++idx;
142+
}
143+
++sidx;
144+
}
145+
117146
std::shared_ptr<CCDBFetcherHelper> helper = std::make_shared<CCDBFetcherHelper>();
118147
CCDBFetcherHelper::initialiseHelper(*helper, options);
119148
std::unordered_map<std::string, int> bindings;
120149
fillValidRoutes(*helper, spec.outputs, bindings);
121150

122-
return adaptStateless([schemas, bindings, helper](InputRecord& inputs, DataTakingContext& dtc, DataAllocator& allocator, TimingInfo& timingInfo, DataProcessingStats& stats) {
151+
return adaptStateless([schemas, bindings, helper, allbuilders](InputRecord& inputs, DataTakingContext& dtc, DataAllocator& allocator, TimingInfo& timingInfo, DataProcessingStats& stats) {
123152
O2_SIGNPOST_ID_GENERATE(sid, ccdb);
124153
O2_SIGNPOST_START(ccdb, sid, "fetchFromAnalysisCCDB", "Fetching CCDB objects for analysis%" PRIu64, (uint64_t)timingInfo.timeslice);
125-
for (auto& schema : schemas) {
154+
std::ranges::for_each(allbuilders, [](auto& builder) { builder.second->Reset(); });
155+
for (auto i = 0U; i < schemas.size(); ++i) {
156+
auto& schema = schemas[i];
126157
std::vector<CCDBFetcherHelper::FetchOp> ops;
127158
auto inputBinding = *schema->metadata()->Get("sourceTable");
128159
auto inputMatcher = DataSpecUtils::fromString(*schema->metadata()->Get("sourceMatcher"));
@@ -134,6 +165,7 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
134165
auto table = inputs.get<TableConsumer>(inputMatcher)->asArrowTable();
135166
// FIXME: make the fTimestamp column configurable.
136167
auto timestampColumn = table->GetColumnByName("fTimestamp");
168+
auto reserveSize = timestampColumn->length();
137169
O2_SIGNPOST_EVENT_EMIT_INFO(ccdb, sid, "fetchFromAnalysisCCDB",
138170
"There are %zu bindings available", bindings.size());
139171
for (auto& binding : bindings) {
@@ -143,9 +175,16 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
143175
}
144176
int outputRouteIndex = bindings.at(outRouteDesc);
145177
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>());
178+
auto builders = allbuilders | std::views::filter([&i](auto const& builder) { return builder.first == i; });
179+
unsigned int numBuilders = std::ranges::count_if(allbuilders, [&i](auto const& builder) { return builder.first == i; });
180+
arrow::Status status;
181+
std::ranges::for_each(builders, [&status, &reserveSize](auto& builder) {
182+
if (reserveSize > builder.second->capacity()) {
183+
status &= builder.second->Reserve(reserveSize - builder.second->capacity());
184+
}
185+
});
186+
if (!status.ok()) {
187+
throw framework::runtime_error_f("Failed to reserve arrays: ", status.ToString().c_str());
149188
}
150189

151190
for (auto ci = 0; ci < timestampColumn->num_chunks(); ++ci) {
@@ -171,15 +210,25 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
171210
O2_SIGNPOST_START(ccdb, sid, "handlingResponses",
172211
"Got %zu responses from server.",
173212
responses.size());
174-
if (builders.size() != responses.size()) {
175-
LOGP(fatal, "Not enough responses (expected {}, found {})", builders.size(), responses.size());
213+
if (numBuilders != responses.size()) {
214+
LOGP(fatal, "Not enough responses (expected {}, found {})", numBuilders, responses.size());
176215
}
177216
arrow::Status result;
178-
for (size_t bi = 0; bi < responses.size(); bi++) {
179-
auto& builder = builders[bi];
217+
218+
int bi = 0;
219+
for (auto& builder : builders) {
180220
auto& response = responses[bi];
221+
#if (FAIRMQ_VERSION_DEC >= 111000)
222+
result &= builder.second->Append();
223+
auto* value_builder = dynamic_cast<arrow::Int64Builder*>(builder.second->value_builder());
224+
result &= value_builder->Append(response.id.handle);
225+
result &= value_builder->Append(response.id.segment);
226+
result &= value_builder->Append(response.size);
227+
#else
181228
char const* address = reinterpret_cast<char const*>(response.id.value);
182-
result &= builder->Append(std::string_view(address, response.size));
229+
result &= builder.second->Append(std::string_view(address, response.size));
230+
#endif
231+
++bi;
183232
}
184233
if (!result.ok()) {
185234
LOGP(fatal, "Error adding results from CCDB");
@@ -188,9 +237,7 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
188237
}
189238
}
190239
arrow::ArrayVector arrays;
191-
for (auto& builder : builders) {
192-
arrays.push_back(*builder->Finish());
193-
}
240+
std::ranges::for_each(builders, [&arrays](auto& builder) { arrays.push_back(*builder.second->Finish()); });
194241
auto outTable = arrow::Table::Make(schema, arrays);
195242
auto concrete = DataSpecUtils::asConcreteDataMatcher(spec);
196243
allocator.adopt(Output{concrete.origin, concrete.description, concrete.subSpec}, outTable);

Framework/Core/include/Framework/ASoA.h

Lines changed: 91 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,17 @@
4041
#include <cstring>
4142
#include <gsl/span> // IWYU pragma: export
4243

44+
namespace fair::mq::shmem
45+
{
46+
struct MetaHeader;
47+
}
48+
4349
namespace o2::framework
4450
{
4551
using ListVector = std::vector<std::vector<int64_t>>;
52+
#if (FAIRMQ_VERSION_DEC >= 111000)
53+
using PointerReconstructor = std::function<std::byte*(fair::mq::shmem::MetaHeader&&)>;
54+
#endif
4655

4756
std::string cutString(std::string&& str);
4857
std::string strToUpper(std::string&& str);
@@ -1014,6 +1023,9 @@ struct ColumnDataHolder {
10141023
arrow::ChunkedArray* second;
10151024
};
10161025

1026+
template <typename C>
1027+
concept needs_ptr_rec = C::needs_ptr_rec;
1028+
10171029
template <typename D, typename O, typename IP, typename... C>
10181030
struct TableIterator : IP, C... {
10191031
public:
@@ -1182,6 +1194,21 @@ struct TableIterator : IP, C... {
11821194
{
11831195
doSetCurrentInternal(internal_index_columns_t{}, table);
11841196
}
1197+
#if (FAIRMQ_VERSION_DEC >= 111000)
1198+
void setPointerReconstructor(framework::PointerReconstructor const& pointerReconstructor)
1199+
{
1200+
[&pointerReconstructor, this]<typename... Cs>(framework::pack<Cs...>) {
1201+
([&pointerReconstructor, this]<typename CC>() {
1202+
if constexpr (needs_ptr_rec<CC>) {
1203+
if (pointerReconstructor) {
1204+
CC::ptrRec = &pointerReconstructor;
1205+
}
1206+
}
1207+
}.template operator()<Cs>(),
1208+
...);
1209+
}(all_columns{});
1210+
}
1211+
#endif
11851212

11861213
private:
11871214
/// Helper to move at the end of columns which actually have an iterator.
@@ -2146,7 +2173,12 @@ class Table
21462173
{
21472174
return self_t{mTable->Slice(0, 0), 0};
21482175
}
2149-
2176+
#if (FAIRMQ_VERSION_DEC >= 111000)
2177+
void setPointerReconstructor(framework::PointerReconstructor const& pointerReconstructor)
2178+
{
2179+
mBegin.setPointerReconstructor(pointerReconstructor);
2180+
}
2181+
#endif
21502182
private:
21512183
template <typename T>
21522184
arrow::ChunkedArray* lookupColumn()
@@ -2325,6 +2357,57 @@ consteval static std::string_view namespace_prefix()
23252357
}; \
23262358
[[maybe_unused]] static constexpr o2::framework::expressions::BindingNode _Getter_ { _Label_, _Name_::hash, o2::framework::expressions::selectArrowType<_Type_>() }
23272359

2360+
#if (FAIRMQ_VERSION_DEC >= 111000)
2361+
#define DECLARE_SOA_CCDB_COLUMN_FULL(_Name_, _Label_, _Getter_, _ConcreteType_, _CCDBQuery_) \
2362+
struct _Name_ : o2::soa::Column<int64_t[3], _Name_> { \
2363+
static constexpr const char* mLabel = _Label_; \
2364+
static constexpr const char* query = _CCDBQuery_; \
2365+
static constexpr const uint32_t hash = crc32(namespace_prefix<_Name_>(), std::string_view{#_Getter_}); \
2366+
static constexpr bool needs_ptr_rec = true; \
2367+
std::function<std::byte*(fair::mq::shmem::MetaHeader&&)> const* ptrRec = nullptr; \
2368+
using base = o2::soa::Column<int64_t[3], _Name_>; \
2369+
using type = int64_t[3]; \
2370+
using column_t = _Name_; \
2371+
_Name_(arrow::ChunkedArray const* column) \
2372+
: o2::soa::Column<int64_t[3], _Name_>(o2::soa::ColumnIterator<int64_t[3]>(column)) \
2373+
{ \
2374+
} \
2375+
\
2376+
_Name_() = default; \
2377+
_Name_(_Name_ const& other) = default; \
2378+
_Name_& operator=(_Name_ const& other) = default; \
2379+
\
2380+
decltype(auto) _Getter_() const \
2381+
{ \
2382+
auto& [handle, segment, size] = *mColumnIterator; \
2383+
auto span = std::span<std::byte>{(*ptrRec)(fair::mq::shmem::MetaHeader{ \
2384+
static_cast<size_t>(size), \
2385+
0, handle, 0, 0, \
2386+
static_cast<uint16_t>(segment), true}), \
2387+
static_cast<size_t>(size)}; \
2388+
if constexpr (std::same_as<_ConcreteType_, std::span<std::byte>>) { \
2389+
return span; \
2390+
} else { \
2391+
static std::byte* payload = nullptr; \
2392+
static _ConcreteType_* deserialised = nullptr; \
2393+
static TClass* c = TClass::GetClass(#_ConcreteType_); \
2394+
if (payload != (std::byte*)span.data()) { \
2395+
payload = (std::byte*)span.data(); \
2396+
delete deserialised; \
2397+
TBufferFile f(TBufferFile::EMode::kRead, span.size(), (char*)span.data(), kFALSE); \
2398+
deserialised = (_ConcreteType_*)soa::extractCCDBPayload((char*)payload, span.size(), c, "ccdb_object"); \
2399+
} \
2400+
return *deserialised; \
2401+
} \
2402+
} \
2403+
\
2404+
decltype(auto) \
2405+
get() const \
2406+
{ \
2407+
return _Getter_(); \
2408+
} \
2409+
};
2410+
#else
23282411
#define DECLARE_SOA_CCDB_COLUMN_FULL(_Name_, _Label_, _Getter_, _ConcreteType_, _CCDBQuery_) \
23292412
struct _Name_ : o2::soa::Column<std::span<std::byte>, _Name_> { \
23302413
static constexpr const char* mLabel = _Label_; \
@@ -2367,6 +2450,7 @@ consteval static std::string_view namespace_prefix()
23672450
return _Getter_(); \
23682451
} \
23692452
};
2453+
#endif
23702454

23712455
#define DECLARE_SOA_CCDB_COLUMN(_Name_, _Getter_, _ConcreteType_, _CCDBQuery_) \
23722456
DECLARE_SOA_CCDB_COLUMN_FULL(_Name_, "f" #_Name_, _Getter_, _ConcreteType_, _CCDBQuery_)
@@ -3715,7 +3799,12 @@ class FilteredBase : public T
37153799
{
37163800
return mCached;
37173801
}
3718-
3802+
#if (FAIRMQ_VERSION_DEC >= 111000)
3803+
void setPointerReconstructor(framework::PointerReconstructor const& pointerReconstructor)
3804+
{
3805+
mFilteredBegin.setPointerReconstructor(pointerReconstructor);
3806+
}
3807+
#endif
37193808
private:
37203809
void resetRanges()
37213810
{

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)