risingwave_batch_executors/executor/join/
distributed_lookup_join.rs1use std::marker::PhantomData;
16use std::mem::swap;
17
18use anyhow::anyhow;
19use futures::pin_mut;
20use itertools::Itertools;
21use risingwave_batch::task::ShutdownToken;
22use risingwave_common::bitmap::Bitmap;
23use risingwave_common::catalog::{ColumnDesc, ColumnId, Field, Schema};
24use risingwave_common::hash::{HashKey, HashKeyDispatcher, VnodeCountCompat};
25use risingwave_common::memory::MemoryContext;
26use risingwave_common::row::OwnedRow;
27use risingwave_common::types::{DataType, Datum};
28use risingwave_common::util::chunk_coalesce::DataChunkBuilder;
29use risingwave_common::util::iter_util::ZipEqFast;
30use risingwave_common::util::scan_range::ScanRange;
31use risingwave_expr::expr::{BoxedExpression, build_from_prost};
32use risingwave_pb::batch_plan::plan_node::NodeBody;
33use risingwave_pb::common::BatchQueryEpoch;
34use risingwave_storage::store::PrefetchOptions;
35use risingwave_storage::table::TableIter;
36use risingwave_storage::table::batch_table::BatchTable;
37use risingwave_storage::{StateStore, dispatch_state_store};
38
39use super::AsOfDesc;
40use crate::error::Result;
41use crate::executor::join::JoinType;
42use crate::executor::{
43 BoxedDataChunkStream, BoxedExecutor, BoxedExecutorBuilder, BufferChunkExecutor, Executor,
44 ExecutorBuilder, LookupExecutorBuilder, LookupJoinBase,
45};
46
47pub struct DistributedLookupJoinExecutor<K, S: StateStore> {
57 base: LookupJoinBase<K, InnerSideExecutorBuilder<S>>,
58}
59
60impl<K: HashKey, S: StateStore> Executor for DistributedLookupJoinExecutor<K, S> {
61 fn schema(&self) -> &Schema {
62 &self.base.schema
63 }
64
65 fn identity(&self) -> &str {
66 &self.base.identity
67 }
68
69 fn execute(self: Box<Self>) -> BoxedDataChunkStream {
70 Box::new(self.base).do_execute()
71 }
72}
73
74impl<K, S: StateStore> DistributedLookupJoinExecutor<K, S> {
75 fn new(base: LookupJoinBase<K, InnerSideExecutorBuilder<S>>) -> Self {
76 Self { base }
77 }
78}
79
80pub struct DistributedLookupJoinExecutorBuilder {}
81
82impl BoxedExecutorBuilder for DistributedLookupJoinExecutorBuilder {
83 async fn new_boxed_executor(
84 source: &ExecutorBuilder<'_>,
85 inputs: Vec<BoxedExecutor>,
86 ) -> Result<BoxedExecutor> {
87 let [outer_side_input]: [_; 1] = inputs.try_into().unwrap();
88
89 let distributed_lookup_join_node = try_match_expand!(
90 source.plan_node().get_node_body().unwrap(),
91 NodeBody::DistributedLookupJoin
92 )?;
93
94 let join_type = JoinType::from_prost(distributed_lookup_join_node.get_join_type()?);
95 let condition = match distributed_lookup_join_node.get_condition() {
96 Ok(cond_prost) => Some(build_from_prost(cond_prost)?),
97 Err(_) => None,
98 };
99
100 let output_indices: Vec<usize> = distributed_lookup_join_node
101 .get_output_indices()
102 .iter()
103 .map(|&x| x as usize)
104 .collect();
105
106 let outer_side_data_types = outer_side_input.schema().data_types();
107
108 let table_desc = distributed_lookup_join_node.get_inner_side_table_desc()?;
109 let inner_side_column_ids = distributed_lookup_join_node
110 .get_inner_side_column_ids()
111 .clone();
112
113 let inner_side_schema = Schema {
114 fields: inner_side_column_ids
115 .iter()
116 .map(|&id| {
117 let column = table_desc
118 .columns
119 .iter()
120 .find(|c| c.column_id == id)
121 .unwrap();
122 Field::from(&ColumnDesc::from(column))
123 })
124 .collect_vec(),
125 };
126
127 let fields = if join_type == JoinType::LeftSemi || join_type == JoinType::LeftAnti {
128 outer_side_input.schema().fields.clone()
129 } else {
130 [
131 outer_side_input.schema().fields.clone(),
132 inner_side_schema.fields.clone(),
133 ]
134 .concat()
135 };
136
137 let original_schema = Schema { fields };
138 let actual_schema = output_indices
139 .iter()
140 .map(|&idx| original_schema[idx].clone())
141 .collect();
142
143 let mut outer_side_key_idxs = vec![];
144 for outer_side_key in distributed_lookup_join_node.get_outer_side_key() {
145 outer_side_key_idxs.push(*outer_side_key as usize)
146 }
147
148 let outer_side_key_types: Vec<DataType> = outer_side_key_idxs
149 .iter()
150 .map(|&i| outer_side_data_types[i].clone())
151 .collect_vec();
152
153 let lookup_prefix_len: usize =
154 distributed_lookup_join_node.get_lookup_prefix_len() as usize;
155
156 let mut inner_side_key_idxs = vec![];
157 for inner_side_key in distributed_lookup_join_node.get_inner_side_key() {
158 inner_side_key_idxs.push(*inner_side_key as usize)
159 }
160
161 let inner_side_key_types = inner_side_key_idxs
162 .iter()
163 .map(|&i| inner_side_schema.fields[i].data_type.clone())
164 .collect_vec();
165
166 let null_safe = distributed_lookup_join_node.get_null_safe().clone();
167
168 let chunk_size = source.context().get_config().developer.chunk_size;
169
170 let asof_desc = distributed_lookup_join_node
171 .asof_desc
172 .map(|desc| AsOfDesc::from_protobuf(&desc))
173 .transpose()?;
174
175 let column_ids = inner_side_column_ids
176 .iter()
177 .copied()
178 .map(ColumnId::from)
179 .collect();
180
181 let vnodes = Some(Bitmap::ones(table_desc.vnode_count()).into());
187
188 dispatch_state_store!(source.context().state_store(), state_store, {
189 let table = BatchTable::new_partial(state_store, column_ids, vnodes, table_desc);
190 let inner_side_builder = InnerSideExecutorBuilder::new(
191 outer_side_key_types,
192 inner_side_key_types.clone(),
193 lookup_prefix_len,
194 distributed_lookup_join_node
195 .query_epoch
196 .ok_or_else(|| anyhow!("query_epoch not set in distributed lookup join"))?,
197 vec![],
198 table,
199 chunk_size,
200 );
201
202 let identity = source.plan_node().get_identity().clone();
203
204 Ok(DistributedLookupJoinExecutorArgs {
205 join_type,
206 condition,
207 outer_side_input,
208 outer_side_data_types,
209 outer_side_key_idxs,
210 inner_side_builder,
211 inner_side_key_types,
212 inner_side_key_idxs,
213 null_safe,
214 lookup_prefix_len,
215 chunk_builder: DataChunkBuilder::new(original_schema.data_types(), chunk_size),
216 schema: actual_schema,
217 output_indices,
218 chunk_size,
219 asof_desc,
220 identity: identity.clone(),
221 shutdown_rx: source.shutdown_rx().clone(),
222 mem_ctx: source.context().create_executor_mem_context(&identity),
223 }
224 .dispatch())
225 })
226 }
227}
228
229struct DistributedLookupJoinExecutorArgs<S: StateStore> {
230 join_type: JoinType,
231 condition: Option<BoxedExpression>,
232 outer_side_input: BoxedExecutor,
233 outer_side_data_types: Vec<DataType>,
234 outer_side_key_idxs: Vec<usize>,
235 inner_side_builder: InnerSideExecutorBuilder<S>,
236 inner_side_key_types: Vec<DataType>,
237 inner_side_key_idxs: Vec<usize>,
238 null_safe: Vec<bool>,
239 lookup_prefix_len: usize,
240 chunk_builder: DataChunkBuilder,
241 schema: Schema,
242 output_indices: Vec<usize>,
243 chunk_size: usize,
244 asof_desc: Option<AsOfDesc>,
245 identity: String,
246 shutdown_rx: ShutdownToken,
247 mem_ctx: MemoryContext,
248}
249
250impl<S: StateStore> HashKeyDispatcher for DistributedLookupJoinExecutorArgs<S> {
251 type Output = BoxedExecutor;
252
253 fn dispatch_impl<K: HashKey>(self) -> Self::Output {
254 Box::new(DistributedLookupJoinExecutor::<K, S>::new(LookupJoinBase {
255 join_type: self.join_type,
256 condition: self.condition,
257 outer_side_input: self.outer_side_input,
258 outer_side_data_types: self.outer_side_data_types,
259 outer_side_key_idxs: self.outer_side_key_idxs,
260 inner_side_builder: self.inner_side_builder,
261 inner_side_key_types: self.inner_side_key_types,
262 inner_side_key_idxs: self.inner_side_key_idxs,
263 null_safe: self.null_safe,
264 lookup_prefix_len: self.lookup_prefix_len,
265 chunk_builder: self.chunk_builder,
266 schema: self.schema,
267 output_indices: self.output_indices,
268 chunk_size: self.chunk_size,
269 asof_desc: self.asof_desc,
270 identity: self.identity,
271 shutdown_rx: self.shutdown_rx,
272 mem_ctx: self.mem_ctx,
273 _phantom: PhantomData,
274 }))
275 }
276
277 fn data_types(&self) -> &[DataType] {
278 &self.inner_side_key_types
279 }
280}
281
282struct InnerSideExecutorBuilder<S: StateStore> {
284 outer_side_key_types: Vec<DataType>,
285 inner_side_key_types: Vec<DataType>,
286 lookup_prefix_len: usize,
287 epoch: BatchQueryEpoch,
288 row_list: Vec<OwnedRow>,
289 table: BatchTable<S>,
290 chunk_size: usize,
291}
292
293impl<S: StateStore> InnerSideExecutorBuilder<S> {
294 fn new(
295 outer_side_key_types: Vec<DataType>,
296 inner_side_key_types: Vec<DataType>,
297 lookup_prefix_len: usize,
298 epoch: BatchQueryEpoch,
299 row_list: Vec<OwnedRow>,
300 table: BatchTable<S>,
301 chunk_size: usize,
302 ) -> Self {
303 Self {
304 outer_side_key_types,
305 inner_side_key_types,
306 lookup_prefix_len,
307 epoch,
308 row_list,
309 table,
310 chunk_size,
311 }
312 }
313}
314
315impl<S: StateStore> LookupExecutorBuilder for InnerSideExecutorBuilder<S> {
316 fn reset(&mut self) {
317 }
319
320 async fn add_scan_range(&mut self, key_datums: Vec<Datum>) -> Result<()> {
322 let mut scan_range = ScanRange::full_table_scan();
323
324 for ((datum, outer_type), inner_type) in key_datums
325 .into_iter()
326 .zip_eq_fast(
327 self.outer_side_key_types
328 .iter()
329 .take(self.lookup_prefix_len),
330 )
331 .zip_eq_fast(
332 self.inner_side_key_types
333 .iter()
334 .take(self.lookup_prefix_len),
335 )
336 {
337 let datum = if inner_type == outer_type {
338 datum
339 } else {
340 bail!("Join key types are not aligned: LHS: {outer_type:?}, RHS: {inner_type:?}");
341 };
342
343 scan_range.eq_conds.push(datum);
344 }
345
346 let pk_prefix = OwnedRow::new(scan_range.eq_conds);
347
348 if self.lookup_prefix_len == self.table.pk_indices().len() {
349 let row = self.table.get_row(&pk_prefix, self.epoch.into()).await?;
350
351 if let Some(row) = row {
352 self.row_list.push(row);
353 }
354 } else {
355 let iter = self
356 .table
357 .batch_iter_with_pk_bounds(
358 self.epoch.into(),
359 &pk_prefix,
360 ..,
361 false,
362 PrefetchOptions::default(),
363 )
364 .await?;
365
366 pin_mut!(iter);
367 while let Some(row) = iter.next_row().await? {
368 self.row_list.push(row);
369 }
370 }
371
372 Ok(())
373 }
374
375 async fn build_executor(&mut self) -> Result<BoxedExecutor> {
377 let mut data_chunk_builder =
378 DataChunkBuilder::new(self.table.schema().data_types(), self.chunk_size);
379 let mut chunk_list = Vec::new();
380
381 let mut new_row_list = vec![];
382 swap(&mut new_row_list, &mut self.row_list);
383
384 for row in new_row_list {
385 if let Some(chunk) = data_chunk_builder.append_one_row(row) {
386 chunk_list.push(chunk);
387 }
388 }
389 if let Some(chunk) = data_chunk_builder.consume_all() {
390 chunk_list.push(chunk);
391 }
392
393 Ok(Box::new(BufferChunkExecutor::new(
394 self.table.schema().clone(),
395 chunk_list,
396 )))
397 }
398}