| 2178 | } |
| 2179 | |
| 2180 | Result<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 | |
| 2233 | void JoinProbeProcessor::Init(int num_key_columns, JoinType join_type, |
| 2234 | SwissTableForJoin* hash_table, |
nothing calls this directly
no test coverage detected