75 auto& dec = ic.
services().get<DanglingEdgesContext>();
80 std::unordered_map<std::string, std::string> ccdbUrls;
81 for (
auto& input : dec.analysisCCDBInputs) {
82 for (
auto&
m : input.metadata) {
83 if (!
m.name.starts_with(
"ccdb:") || ccdbUrls.count(
m.name)) {
86 std::string
url =
m.defaultValue.asString();
88 url = options.
get<std::string>(
m.name.c_str());
90 LOGP(info,
"CCDB path resolved for {}: {}",
m.name,
url);
91 ccdbUrls.emplace(
m.name, std::move(
url));
94 std::vector<std::shared_ptr<arrow::Schema>> schemas;
95 for (
auto& input : dec.analysisCCDBInputs) {
96 auto schemaMetadata = std::make_shared<arrow::KeyValueMetadata>();
97 std::vector<std::shared_ptr<arrow::Field>>
fields;
99 schemaMetadata->Append(
"outputBinding", input.binding);
100 for (
auto&
m : input.metadata) {
101 if (
m.name.starts_with(
"input:")) {
102 auto name =
m.name.substr(6);
103 schemaMetadata->Append(
"sourceTable",
name);
107 if (!
m.name.starts_with(
"ccdb:")) {
110 auto fieldMetadata = std::make_shared<arrow::KeyValueMetadata>();
111 auto it = ccdbUrls.find(
m.name);
112 fieldMetadata->Append(
"url", it != ccdbUrls.end() ? it->second :
m.defaultValue.asString());
113 auto columnName =
m.name.substr(strlen(
"ccdb:"));
114#if (FAIRMQ_VERSION_DEC >= 111000)
117 fields.emplace_back(std::make_shared<arrow::Field>(columnName, arrow::binary_view(),
false, fieldMetadata));
120 schemas.emplace_back(std::make_shared<arrow::Schema>(
fields, schemaMetadata));
123#if (FAIRMQ_VERSION_DEC >= 111000)
124 std::vector<std::pair<uint32_t, std::shared_ptr<arrow::FixedSizeListBuilder>>> allbuilders;
126 std::vector<std::pair<uint32_t, std::shared_ptr<arrow::BinaryViewBuilder>>> allbuilders;
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();
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));
139 allbuilders[idx] = std::make_pair(sidx, std::make_shared<arrow::BinaryViewBuilder>());
146 std::shared_ptr<CCDBFetcherHelper> helper = std::make_shared<CCDBFetcherHelper>();
148 std::unordered_map<std::string, int> bindings;
149 fillValidRoutes(*helper, spec.
outputs, bindings);
153 O2_SIGNPOST_START(ccdb, sid,
"fetchFromAnalysisCCDB",
"Fetching CCDB objects for analysis%" PRIu64, (uint64_t)timingInfo.
timeslice);
154 std::ranges::for_each(allbuilders, [](
auto& builder) { builder.second->Reset(); });
155 for (
auto i = 0U;
i < schemas.size(); ++
i) {
157 std::vector<CCDBFetcherHelper::FetchOp> ops;
158 auto inputBinding = *
schema->metadata()->Get(
"sourceTable");
160 auto outRouteDesc = *
schema->metadata()->Get(
"outputRoute");
161 std::string outBinding = *
schema->metadata()->Get(
"outputBinding");
163 "Fetching CCDB objects for %{public}s's columns with timestamps from %{public}s and putting them in route %{public}s",
164 outBinding.c_str(), inputBinding.c_str(), outRouteDesc.c_str());
167 auto timestampColumn = table->GetColumnByName(
"fTimestamp");
168 auto reserveSize = timestampColumn->length();
170 "There are %zu bindings available", bindings.size());
171 for (
auto const&
binding : bindings) {
176 int outputRouteIndex = bindings.at(outRouteDesc);
177 auto& spec = helper->routes[outputRouteIndex].matcher;
180 auto builders = allbuilders | std::views::filter([&
i](
auto const& builder) {
return builder.first ==
i; });
181 unsigned int numBuilders = std::ranges::count_if(allbuilders, [&
i](
auto const& builder) {
return builder.first ==
i; });
182 arrow::Status status;
183 std::ranges::for_each(builders, [&status, &reserveSize](
auto& builder) {
184 if (reserveSize > builder.second->capacity()) {
185 status &= builder.second->Reserve(reserveSize - builder.second->capacity());
194 for (
auto ci = 0; ci < timestampColumn->num_chunks(); ++ci) {
195 std::shared_ptr<arrow::Array> chunk = timestampColumn->chunk(ci);
196 auto const* timestamps = chunk->data()->GetValuesSafe<
size_t>(1);
198 for (int64_t ri = 0; ri < chunk->data()->length; ri++) {
200 int64_t timestamp = timestamps[ri];
201 for (
auto& field :
schema->fields()) {
202 auto url = *field->metadata()->Get(
"url");
207 .timestamp = timestamp,
215 "Got %zu responses from server.",
217 if (numBuilders != responses.size()) {
218 LOGP(fatal,
"Not enough responses (expected {}, found {})", numBuilders, responses.size());
223 for (
auto& builder : builders) {
224 auto& response = responses[bi];
225 auto& lastId = lastIds[bi];
226 if (response.id.value != lastId.value) {
227 lastId.value = response.id.value;
230#if (FAIRMQ_VERSION_DEC >= 111000)
231 result &= builder.second->Append();
232 auto* value_builder =
dynamic_cast<arrow::Int64Builder*
>(builder.second->value_builder());
233 result &= value_builder->Append(response.id.handle);
234 result &= value_builder->Append(response.id.segment);
235 result &= value_builder->Append(response.size);
237 char const*
address =
reinterpret_cast<char const*
>(response.id.value);
238 result &= builder.second->Append(std::string_view(
address, response.size));
243 LOGP(fatal,
"Error adding results from CCDB");
245 O2_SIGNPOST_END(ccdb, sid,
"handlingResponses",
"Done processing responses");
248 arrow::ArrayVector
arrays;
249 std::ranges::for_each(builders, [&
arrays](
auto& builder) {
arrays.push_back(*builder.second->Finish()); });
256 O2_SIGNPOST_END(ccdb, sid,
"fetchFromAnalysisCCDB",
"Fetching CCDB objects");
#define O2_SIGNPOST_EVENT_EMIT_INFO(log, id, name, format,...)
#define O2_SIGNPOST_END(log, id, name, format,...)
#define O2_SIGNPOST_ID_GENERATE(name, log)
#define O2_SIGNPOST_START(log, id, name, format,...)