Read PK table with Deduplicate engine: level-0 splits go through KeyValueFileReader for sort-merge dedup, compacted splits use DataFileReader.
(
&self,
data_splits: &[DataSplit],
core_options: &CoreOptions,
)
| 98 | /// Read PK table with Deduplicate engine: level-0 splits go through |
| 99 | /// KeyValueFileReader for sort-merge dedup, compacted splits use DataFileReader. |
| 100 | fn read_pk( |
| 101 | &self, |
| 102 | data_splits: &[DataSplit], |
| 103 | core_options: &CoreOptions, |
| 104 | ) -> crate::Result<ArrowRecordBatchStream> { |
| 105 | if core_options.merge_engine()? == MergeEngine::PartialUpdate { |
| 106 | return self.read_kv(data_splits, core_options); |
| 107 | } |
| 108 | |
| 109 | let mut kv_splits = Vec::new(); |
| 110 | let mut raw_splits = Vec::new(); |
| 111 | for split in data_splits { |
| 112 | if split.data_files().iter().any(|f| f.level == 0) { |
| 113 | kv_splits.push(split.clone()); |
| 114 | } else { |
| 115 | raw_splits.push(split.clone()); |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | if raw_splits.is_empty() { |
| 120 | return self.read_kv(&kv_splits, core_options); |
| 121 | } |
| 122 | if kv_splits.is_empty() { |
| 123 | return self.read_raw(&raw_splits); |
| 124 | } |
| 125 | |
| 126 | let kv_stream = self.read_kv(&kv_splits, core_options)?; |
| 127 | let raw_stream = self.read_raw(&raw_splits)?; |
| 128 | Ok(Box::pin(futures::stream::select_all([ |
| 129 | kv_stream, raw_stream, |
| 130 | ]))) |
| 131 | } |
| 132 | |
| 133 | /// Read splits via KeyValueFileReader (sort-merge dedup). |
| 134 | fn read_kv( |
no test coverage detected