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
85 changes: 80 additions & 5 deletions Framework/CCDBSupport/src/AnalysisCCDBHelpers.cxx
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,10 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
schemaMetadata->Append("sourceMatcher", DataSpecUtils::describe(std::get<ConcreteDataMatcher>(DataSpecUtils::fromMetadataString(m.defaultValue.get<std::string>()).matcher)));
continue;
}
if (m.name == "timestamp-column" || m.name == "uniformity-column") {
schemaMetadata->Append(m.name, m.defaultValue.asString());
continue;
}
if (!m.name.starts_with("ccdb:")) {
continue;
}
Expand Down Expand Up @@ -144,15 +148,44 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
auto& schema = schemas[i];
std::vector<CCDBFetcherHelper::FetchOp> ops;
auto inputBinding = *schema->metadata()->Get("sourceTable");
auto inputMatcher = DataSpecUtils::fromString(*schema->metadata()->Get("sourceMatcher"));
auto outRouteDesc = *schema->metadata()->Get("outputRoute");
std::string outBinding = *schema->metadata()->Get("outputBinding");
auto timestampColumnName = schema->metadata()->Contains("timestamp-column") ? *schema->metadata()->Get("timestamp-column") : std::string{"fTimestamp"};
auto uniformityColumnName = schema->metadata()->Contains("uniformity-column") ? *schema->metadata()->Get("uniformity-column") : timestampColumnName;
O2_SIGNPOST_EVENT_EMIT_INFO(ccdb, sid, "fetchFromAnalysisCCDB",
"Fetching CCDB objects for %{public}s's columns with timestamps from %{public}s and putting them in route %{public}s",
outBinding.c_str(), inputBinding.c_str(), outRouteDesc.c_str());
auto table = inputs.get<TableConsumer>(inputMatcher)->asArrowTable();
// FIXME: make the fTimestamp column configurable.
auto timestampColumn = table->GetColumnByName("fTimestamp");
// The timestamp and uniformity columns may live in different source tables (the
// run number is on aod::BCs, the timestamp on aod::Timestamps). Locate each by
// name across every declared source, and read them positionally.
std::shared_ptr<arrow::ChunkedArray> timestampColumn;
std::shared_ptr<arrow::ChunkedArray> uniformityColumn;
auto const& schemaKeys = schema->metadata()->keys();
auto const& schemaValues = schema->metadata()->values();
for (size_t mi = 0; mi < schemaKeys.size(); ++mi) {
if (schemaKeys[mi] != "sourceMatcher") {
continue;
}
auto sourceTable = inputs.get<TableConsumer>(DataSpecUtils::fromString(schemaValues[mi]))->asArrowTable();
if (auto column = sourceTable->GetColumnByName(timestampColumnName); column && !timestampColumn) {
timestampColumn = column;
}
if (auto column = sourceTable->GetColumnByName(uniformityColumnName); column && !uniformityColumn) {
uniformityColumn = column;
}
}
if (!timestampColumn) {
LOGP(fatal, "No source table of {} provides the timestamp column \"{}\"", outBinding, timestampColumnName);
}
if (!uniformityColumn) {
LOGP(fatal, "No source table of {} provides the uniformity column \"{}\"", outBinding, uniformityColumnName);
}
// Positional reading is only sound if the two sources are row-aligned; ASoA has
// no type-level way to state that, so it is checked here.
if (uniformityColumn->length() != timestampColumn->length()) {
LOGP(fatal, "Uniformity column \"{}\" has {} rows but timestamp column \"{}\" has {}; the two sources of {} are not row-aligned",
uniformityColumnName, uniformityColumn->length(), timestampColumnName, timestampColumn->length(), outBinding);
}
auto reserveSize = timestampColumn->length();
O2_SIGNPOST_EVENT_EMIT_INFO(ccdb, sid, "fetchFromAnalysisCCDB",
"There are %zu bindings available", bindings.size());
Expand All @@ -179,11 +212,50 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)

std::vector<DataAllocator::CacheId> lastIds(numBuilders, DataAllocator::CacheId{.value = -1, .handle = -1, .segment = -1});

// Rows sharing a uniformity value resolve to the same objects, so the query is
// issued once per distinct value and the resulting handles are repeated for the
// rest of the run. When uniformity is the timestamp itself (the default) this
// degenerates to the previous behaviour, one query per row.
std::vector<int64_t> uniformity;
bool const shortCircuit = uniformityColumn.get() != timestampColumn.get();
if (shortCircuit) {
uniformity.reserve(reserveSize);
for (auto uci = 0; uci < uniformityColumn->num_chunks(); ++uci) {
auto uchunk = uniformityColumn->chunk(uci);
auto const length = uchunk->data()->length;
switch (uchunk->type_id()) {
case arrow::Type::INT32:
for (int64_t ui = 0; ui < length; ++ui) {
uniformity.push_back(uchunk->data()->GetValuesSafe<int32_t>(1)[ui]);
}
break;
case arrow::Type::INT64:
case arrow::Type::UINT64:
for (int64_t ui = 0; ui < length; ++ui) {
uniformity.push_back(uchunk->data()->GetValuesSafe<int64_t>(1)[ui]);
}
break;
default:
LOGP(fatal, "Uniformity column \"{}\" of {} has unsupported arrow type {}",
uniformityColumnName, outBinding, uchunk->type()->ToString());
}
}
}
int64_t row = -1;
int64_t previousUniformity = 0;
bool haveResponses = false;
std::vector<CCDBFetcherHelper::Response> responses;

for (auto ci = 0; ci < timestampColumn->num_chunks(); ++ci) {
std::shared_ptr<arrow::Array> chunk = timestampColumn->chunk(ci);
auto const* timestamps = chunk->data()->GetValuesSafe<size_t>(1);

for (int64_t ri = 0; ri < chunk->data()->length; ri++) {
++row;
bool const sameAsPrevious = shortCircuit && haveResponses && uniformity[row] == previousUniformity;
if (shortCircuit) {
previousUniformity = uniformity[row];
}
ops.clear();
int64_t timestamp = timestamps[ri];
for (auto& field : schema->fields()) {
Expand All @@ -198,7 +270,10 @@ AlgorithmSpec AnalysisCCDBHelpers::fetchFromCCDB(ConfigContext const& /*ctx*/)
.queryRate = 0,
});
}
auto responses = CCDBFetcherHelper::populateCacheWith(helper, ops, timingInfo, dtc, allocator);
if (!sameAsPrevious) {
responses = CCDBFetcherHelper::populateCacheWith(helper, ops, timingInfo, dtc, allocator);
haveResponses = true;
}
O2_SIGNPOST_START(ccdb, sid, "handlingResponses",
"Got %zu responses from server.",
responses.size());
Expand Down
Loading
Loading