79 auto& dec = ic.
services().get<DanglingEdgesContext>();
84 std::unordered_map<std::string, std::string> ccdbUrls;
85 std::unordered_map<std::string, std::string> runDependent;
86 for (
auto& input : dec.analysisCCDBInputs) {
87 for (
auto&
m : input.metadata) {
88 if (
m.name.starts_with(
"ccdb-run-dependent:")) {
89 runDependent.emplace(
m.name,
m.defaultValue.asString());
92 if (!
m.name.starts_with(
"ccdb:") || ccdbUrls.count(
m.name)) {
95 std::string
url =
m.defaultValue.asString();
97 url = options.
get<std::string>(
m.name.c_str());
99 LOGP(info,
"CCDB path resolved for {}: {}",
m.name,
url);
100 ccdbUrls.emplace(
m.name, std::move(
url));
103 std::vector<std::shared_ptr<arrow::Schema>> schemas;
104 for (
auto& input : dec.analysisCCDBInputs) {
105 auto schemaMetadata = std::make_shared<arrow::KeyValueMetadata>();
106 std::vector<std::shared_ptr<arrow::Field>>
fields;
108 schemaMetadata->Append(
"outputBinding", input.binding);
109 for (
auto&
m : input.metadata) {
110 if (
m.name.starts_with(
"input:")) {
111 auto name =
m.name.substr(6);
112 schemaMetadata->Append(
"sourceTable",
name);
116 if (
m.name ==
"timestamp-column" ||
m.name ==
"uniformity-column") {
117 schemaMetadata->Append(
m.name,
m.defaultValue.asString());
120 if (!
m.name.starts_with(
"ccdb:")) {
123 auto fieldMetadata = std::make_shared<arrow::KeyValueMetadata>();
124 auto it = ccdbUrls.find(
m.name);
125 fieldMetadata->Append(
"url", it != ccdbUrls.end() ? it->second :
m.defaultValue.asString());
126 auto runDep = runDependent.find(
"ccdb-run-dependent:" +
m.name.substr(strlen(
"ccdb:")));
127 fieldMetadata->Append(
"runDependent", runDep != runDependent.end() ? runDep->second :
"0");
128 auto columnName =
m.name.substr(strlen(
"ccdb:"));
131 schemas.emplace_back(std::make_shared<arrow::Schema>(
fields, schemaMetadata));
135 std::vector<std::vector<PathTable>> pathTables;
136 for (
auto const&
schema : schemas) {
137 auto& tables = pathTables.emplace_back();
138 for (
auto const& field :
schema->fields()) {
143 std::vector<std::pair<uint32_t, std::shared_ptr<arrow::FixedSizeListBuilder>>> allbuilders;
144 allbuilders.resize([&schemas]() {
size_t size = 0;
for (
auto&
schema : schemas) {
size +=
schema->num_fields(); };
return size; }());
145 auto* pool = arrow::default_memory_pool();
149 for (
auto const&
schema : schemas) {
150 for (
auto const& _ :
schema->fields()) {
151 auto value_builder = std::make_shared<arrow::Int64Builder>();
152 allbuilders[idx] = std::make_pair(sidx, std::make_shared<arrow::FixedSizeListBuilder>(pool, std::move(value_builder), 3));
158 std::shared_ptr<CCDBFetcherHelper> helper = std::make_shared<CCDBFetcherHelper>();
160 std::unordered_map<std::string, int> bindings;
161 fillValidRoutes(*helper, spec.
outputs, bindings);
165 O2_SIGNPOST_START(ccdb, sid,
"fetchFromAnalysisCCDB",
"Fetching CCDB objects for analysis%" PRIu64, (uint64_t)timingInfo.
timeslice);
166 std::ranges::for_each(allbuilders, [](
auto& builder) { builder.second->Reset(); });
167 for (
auto i = 0U;
i < schemas.size(); ++
i) {
169 std::vector<CCDBFetcherHelper::FetchOp> ops;
170 auto inputBinding = *
schema->metadata()->Get(
"sourceTable");
171 auto outRouteDesc = *
schema->metadata()->Get(
"outputRoute");
172 std::string outBinding = *
schema->metadata()->Get(
"outputBinding");
173 auto timestampColumnName =
schema->metadata()->Contains(
"timestamp-column") ? *
schema->metadata()->Get(
"timestamp-column") : std::string{
"fTimestamp"};
174 auto uniformityColumnName =
schema->metadata()->Contains(
"uniformity-column") ? *
schema->metadata()->Get(
"uniformity-column") : timestampColumnName;
176 "Fetching CCDB objects for %{public}s's columns with timestamps from %{public}s and putting them in route %{public}s",
177 outBinding.c_str(), inputBinding.c_str(), outRouteDesc.c_str());
181 std::shared_ptr<arrow::ChunkedArray> timestampColumn;
182 std::shared_ptr<arrow::ChunkedArray> uniformityColumn;
183 auto const& schemaKeys =
schema->metadata()->keys();
184 auto const& schemaValues =
schema->metadata()->values();
185 for (
size_t mi = 0; mi < schemaKeys.size(); ++mi) {
186 if (schemaKeys[mi] !=
"sourceMatcher") {
190 if (
auto column = sourceTable->GetColumnByName(timestampColumnName); column && !timestampColumn) {
191 timestampColumn = column;
193 if (
auto column = sourceTable->GetColumnByName(uniformityColumnName); column && !uniformityColumn) {
194 uniformityColumn = column;
197 if (!timestampColumn) {
198 LOGP(fatal,
"No source table of {} provides the timestamp column \"{}\"", outBinding, timestampColumnName);
200 if (!uniformityColumn) {
201 LOGP(fatal,
"No source table of {} provides the uniformity column \"{}\"", outBinding, uniformityColumnName);
205 if (uniformityColumn->length() != timestampColumn->length()) {
206 LOGP(fatal,
"Uniformity column \"{}\" has {} rows but timestamp column \"{}\" has {}; the two sources of {} are not row-aligned",
207 uniformityColumnName, uniformityColumn->length(), timestampColumnName, timestampColumn->length(), outBinding);
209 auto reserveSize = timestampColumn->length();
211 "There are %zu bindings available", bindings.size());
212 for (
auto const&
binding : bindings) {
217 int outputRouteIndex = bindings.at(outRouteDesc);
218 auto& spec = helper->routes[outputRouteIndex].matcher;
221 auto builders = allbuilders | std::views::filter([&
i](
auto const& builder) {
return builder.first ==
i; });
222 unsigned int numBuilders = std::ranges::count_if(allbuilders, [&
i](
auto const& builder) {
return builder.first ==
i; });
223 arrow::Status status;
224 std::ranges::for_each(builders, [&status, &reserveSize](
auto& builder) {
225 if (reserveSize > builder.second->capacity()) {
226 status &= builder.second->Reserve(reserveSize - builder.second->capacity());
239 std::vector<int64_t> uniformity;
240 bool const shortCircuit = uniformityColumn.get() != timestampColumn.get();
242 uniformity.reserve(reserveSize);
243 for (
auto uci = 0; uci < uniformityColumn->num_chunks(); ++uci) {
244 auto uchunk = uniformityColumn->chunk(uci);
245 auto const length = uchunk->data()->length;
246 switch (uchunk->type_id()) {
247 case arrow::Type::INT32:
248 for (int64_t ui = 0; ui <
length; ++ui) {
249 uniformity.push_back(uchunk->data()->GetValuesSafe<int32_t>(1)[ui]);
252 case arrow::Type::INT64:
253 case arrow::Type::UINT64:
254 for (int64_t ui = 0; ui <
length; ++ui) {
255 uniformity.push_back(uchunk->data()->GetValuesSafe<int64_t>(1)[ui]);
259 LOGP(fatal,
"Uniformity column \"{}\" of {} has unsupported arrow type {}",
260 uniformityColumnName, outBinding, uchunk->type()->ToString());
265 int64_t previousUniformity = 0;
266 bool haveResponses =
false;
267 std::vector<CCDBFetcherHelper::Response> responses;
269 for (
auto ci = 0; ci < timestampColumn->num_chunks(); ++ci) {
270 std::shared_ptr<arrow::Array> chunk = timestampColumn->chunk(ci);
271 auto const* timestamps = chunk->data()->GetValuesSafe<
size_t>(1);
273 for (int64_t ri = 0; ri < chunk->data()->length; ri++) {
275 bool const sameAsPrevious = shortCircuit && haveResponses && uniformity[
row] == previousUniformity;
277 previousUniformity = uniformity[
row];
280 int64_t timestamp = timestamps[ri];
283 int64_t
const uniformityKey = shortCircuit ? uniformity[
row] : timestamp;
285 for (
auto& field :
schema->fields()) {
286 auto const&
url = pathTables[
i][fi++].resolve(uniformityKey, field->name());
291 int const fieldRunDependent = field->metadata()->Contains(
"runDependent")
292 ? std::stoi(*field->metadata()->Get(
"runDependent"))
294 if (fieldRunDependent != 0 && uniformityColumnName !=
"fRunNumber") {
295 LOGP(fatal, R
"(Column "{}" of {} is declared run-dependent, but its table is uniform in "{}" rather than fRunNumber, so no run number is available to query with. Declare the table with DECLARE_SOA_UNIFORM_TABLE(..., aod::BCs, o2::aod::bc::RunNumber, ...).)",
296 field->name(), outBinding, uniformityColumnName);
301 .timestamp = timestamp,
302 .runNumber = fieldRunDependent != 0 ? static_cast<int>(uniformityKey) : 1,
303 .runDependent = fieldRunDependent,
307 if (!sameAsPrevious) {
309 haveResponses =
true;
312 "Got %zu responses from server.",
314 if (numBuilders != responses.size()) {
315 LOGP(fatal,
"Not enough responses (expected {}, found {})", numBuilders, responses.size());
320 for (
auto& builder : builders) {
321 auto& response = responses[bi];
322 auto& lastId = lastIds[bi];
323 if (response.id.value != lastId.value) {
324 lastId.value = response.id.value;
327 result &= builder.second->Append();
328 auto* value_builder =
dynamic_cast<arrow::Int64Builder*
>(builder.second->value_builder());
329 result &= value_builder->Append(response.id.handle);
330 result &= value_builder->Append(response.id.segment);
331 result &= value_builder->Append(response.size);
335 LOGP(fatal,
"Error adding results from CCDB");
337 O2_SIGNPOST_END(ccdb, sid,
"handlingResponses",
"Done processing responses");
340 arrow::ArrayVector
arrays;
341 std::ranges::for_each(builders, [&
arrays](
auto& builder) {
arrays.push_back(*builder.second->Finish()); });
348 O2_SIGNPOST_END(ccdb, sid,
"fetchFromAnalysisCCDB",
"Fetching CCDB objects");