MCPcopy Create free account
hub / github.com/apache/arrow / MaterializeFilterInput

Method MaterializeFilterInput

cpp/src/arrow/acero/swiss_join.cc:2180–2231  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

2178}
2179
2180Result<ExecBatch> JoinResidualFilter::MaterializeFilterInput(
2181 const ExecBatch& keypayload_batch, int num_batch_rows, const uint16_t* batch_row_ids,
2182 const uint32_t* key_ids_maybe_null, const uint32_t* payload_ids_maybe_null) const {
2183 ExecBatch out;
2184 out.length = num_batch_rows;
2185 out.values.resize(probe_filter_to_key_and_payload_.size() + num_build_keys_referred_ +
2186 num_build_payloads_referred_);
2187
2188 if (probe_filter_to_key_and_payload_.size() > 0) {
2189 ExecBatchBuilder probe_batch_builder;
2190 RETURN_NOT_OK(probe_batch_builder.AppendSelected(
2191 pool_, keypayload_batch, num_batch_rows, batch_row_ids,
2192 static_cast<int>(probe_filter_to_key_and_payload_.size()),
2193 probe_filter_to_key_and_payload_.data()));
2194 ExecBatch probe_batch = probe_batch_builder.Flush();
2195 ARROW_DCHECK(probe_batch.values.size() == probe_filter_to_key_and_payload_.size());
2196 for (size_t i = 0; i < probe_batch.values.size(); ++i) {
2197 out.values[i] = std::move(probe_batch.values[i]);
2198 }
2199 }
2200
2201 if (num_build_keys_referred_ > 0 || num_build_payloads_referred_ > 0) {
2202 ARROW_DCHECK(num_build_keys_referred_ == 0 || key_ids_maybe_null);
2203 ARROW_DCHECK(num_build_payloads_referred_ == 0 || payload_ids_maybe_null);
2204
2205 int num_build_cols = build_schemas_->num_cols(HashJoinProjection::FILTER);
2206 auto to_key =
2207 build_schemas_->map(HashJoinProjection::FILTER, HashJoinProjection::KEY);
2208 auto to_payload =
2209 build_schemas_->map(HashJoinProjection::FILTER, HashJoinProjection::PAYLOAD);
2210 for (int i = 0; i < num_build_cols; ++i) {
2211 ResizableArrayData column_data;
2212 // Allocate at least 8 rows for the convenience of SIMD decoding.
2213 int log_num_rows_min = std::max(3, bit_util::Log2(num_batch_rows));
2214 RETURN_NOT_OK(
2215 column_data.Init(build_schemas_->data_type(HashJoinProjection::FILTER, i),
2216 pool_, log_num_rows_min));
2217 if (auto idx = to_key.get(i); idx != SchemaProjectionMap::kMissingField) {
2218 RETURN_NOT_OK(build_keys_->DecodeSelected(&column_data, idx, num_batch_rows,
2219 key_ids_maybe_null, pool_));
2220 } else if (idx = to_payload.get(i); idx != SchemaProjectionMap::kMissingField) {
2221 RETURN_NOT_OK(build_payloads_->DecodeSelected(&column_data, idx, num_batch_rows,
2222 payload_ids_maybe_null, pool_));
2223 } else {
2224 ARROW_DCHECK(false);
2225 }
2226 out.values[probe_filter_to_key_and_payload_.size() + i] = column_data.array_data();
2227 }
2228 }
2229
2230 return out;
2231}
2232
2233void JoinProbeProcessor::Init(int num_key_columns, JoinType join_type,
2234 SwissTableForJoin* hash_table,

Callers

nothing calls this directly

Calls 13

resizeMethod · 0.80
AppendSelectedMethod · 0.80
num_colsMethod · 0.80
mapMethod · 0.80
DecodeSelectedMethod · 0.80
array_dataMethod · 0.80
Log2Function · 0.50
sizeMethod · 0.45
dataMethod · 0.45
FlushMethod · 0.45
InitMethod · 0.45
data_typeMethod · 0.45

Tested by

no test coverage detected