risingwave_storage/hummock/compactor/
block_stream.rs1use std::ops::Range;
16
17use await_tree::{InstrumentAwait, SpanExt};
18use bytes::Bytes;
19use fail::fail_point;
20use risingwave_hummock_sdk::sstable_info::SstableInfo;
21
22use crate::hummock::block_stream::BlockDataStream;
23use crate::hummock::sstable_store::SstableStoreRef;
24use crate::hummock::{HummockResult, TableHolder};
25
26pub(super) struct SstableBlockStream {
29 pub(super) sstable: TableHolder,
30 pub(super) sstable_info: SstableInfo,
31 sstable_store: SstableStoreRef,
32 block_stream: Option<BlockDataStream>,
33 remaining_blocks: Range<usize>,
35 io_retry_times: usize,
36 max_io_retry_times: usize,
37}
38
39impl SstableBlockStream {
40 pub(super) fn new(
41 sstable: TableHolder,
42 block_metas_range: Range<usize>,
43 sstable_info: SstableInfo,
44 sstable_store: SstableStoreRef,
45 max_io_retry_times: usize,
46 ) -> Self {
47 Self {
48 sstable,
49 sstable_info,
50 sstable_store,
51 block_stream: None,
52 remaining_blocks: block_metas_range,
53 io_retry_times: 0,
54 max_io_retry_times,
55 }
56 }
57
58 pub(super) fn next_block_index(&self) -> usize {
59 self.remaining_blocks.start
60 }
61
62 pub(super) fn has_next_block(&self) -> bool {
63 !self.remaining_blocks.is_empty()
64 }
65
66 pub(super) async fn next_block(&mut self) -> HummockResult<Option<(Bytes, usize)>> {
67 while self.has_next_block() {
68 if self.block_stream.is_none() {
69 self.block_stream = Some(
72 self.sstable_store
73 .get_stream_for_blocks(
74 self.sstable_info.object_id,
75 &self.sstable.meta.block_metas[self.remaining_blocks.clone()],
76 )
77 .instrument_await("stream_iter_get_stream".verbose())
78 .await?,
79 );
80 }
81
82 match self.block_stream.as_mut().unwrap().next_block().await {
83 Ok(Some(block)) => {
84 self.remaining_blocks.start += 1;
85 return Ok(Some(block));
86 }
87 Ok(None) => {
88 self.remaining_blocks.start = self.remaining_blocks.end;
89 }
90 Err(e) => {
91 if !e.is_object_error() || self.io_retry_times >= self.max_io_retry_times {
92 return Err(e);
93 }
94 self.block_stream = None;
97 self.io_retry_times += 1;
98 fail_point!("create_stream_err");
99 tracing::warn!(
100 object_id = %self.sstable_info.object_id,
101 sst_id = %self.sstable_info.sst_id,
102 meta_offset = self.sstable_info.meta_offset,
103 table_ids = ?self.sstable_info.table_ids,
104 block_index = self.next_block_index(),
105 io_retry_times = self.io_retry_times,
106 "retry create compactor block stream"
107 );
108 }
109 }
110 }
111 Ok(None)
112 }
113}
114
115#[cfg(test)]
116mod tests {
117 use std::sync::Arc;
118
119 use bytes::Bytes;
120 use risingwave_object_store::object::{MonitoredStreamingReader, ObjectError, ObjectResult};
121
122 use super::SstableBlockStream;
123 use crate::hummock::HummockValue;
124 use crate::hummock::block_stream::BlockDataStream;
125 use crate::hummock::iterator::test_utils::mock_sstable_store;
126 use crate::hummock::test_utils::{
127 default_builder_opt_for_test, default_writer_opt_for_test, gen_test_sstable_data, put_sst,
128 test_key_of, test_value_of,
129 };
130 use crate::monitor::{ObjectStoreMetrics, StoreLocalStatistic};
131
132 async fn test_stream(max_io_retry_times: usize) -> (SstableBlockStream, Bytes) {
133 let store = mock_sstable_store().await;
134 let mut options = default_builder_opt_for_test();
135 options.block_capacity = 128;
136 let (data, meta) = gen_test_sstable_data(
137 options,
138 (0..100).map(|i| (test_key_of(i), HummockValue::put(test_value_of(i)))),
139 )
140 .await;
141 assert!(meta.block_metas.len() > 5);
142 let info = put_sst(
143 0,
144 data.clone(),
145 meta,
146 store.clone(),
147 default_writer_opt_for_test(),
148 vec![0],
149 )
150 .await
151 .unwrap();
152 let table = store
153 .sstable(&info, &mut StoreLocalStatistic::default())
154 .await
155 .unwrap();
156 (
158 SstableBlockStream::new(table, 1..5, info, store, max_io_retry_times),
159 data,
160 )
161 }
162
163 fn inject_stream(stream: &mut SstableBlockStream, packets: Vec<ObjectResult<Bytes>>) {
164 let reader = MonitoredStreamingReader::new(
165 "test",
166 Box::pin(futures::stream::iter(packets)),
167 Arc::new(ObjectStoreMetrics::unused()),
168 None,
169 );
170 stream.block_stream = Some(BlockDataStream::new(
171 reader,
172 &stream.sstable.meta.block_metas[stream.remaining_blocks.clone()],
173 ));
174 }
175
176 async fn assert_next_block(stream: &mut SstableBlockStream, original: &Bytes, index: usize) {
177 let meta = stream.sstable.meta.block_metas[index].clone();
178 let (data, uncompressed_size) = stream.next_block().await.unwrap().unwrap();
179 assert_eq!(
180 data,
181 original.slice(meta.offset as usize..(meta.offset + meta.len) as usize)
182 );
183 assert_eq!(uncompressed_size, meta.uncompressed_size as usize);
184 assert_eq!(stream.next_block_index(), index + 1);
185 }
186
187 #[tokio::test]
188 async fn test_compactor_block_stream_partial_read() {
189 for unexpected_eof in [false, true] {
190 let (mut stream, original) = test_stream(1).await;
191 let first = &stream.sstable.meta.block_metas[1];
192 let second = &stream.sstable.meta.block_metas[2];
193 let cutoff = second.offset as usize + second.len as usize / 2;
194 let mut packets = vec![Ok(original.slice(first.offset as usize..cutoff))];
195 if !unexpected_eof {
196 packets.push(Err(ObjectError::internal("injected after partial block")));
197 }
198 inject_stream(&mut stream, packets);
199 assert_next_block(&mut stream, &original, 1).await;
200 assert_eq!(stream.io_retry_times, 0);
201 for index in 2..5 {
204 assert_next_block(&mut stream, &original, index).await;
205 }
206 assert_eq!(stream.io_retry_times, 1);
207 assert!(stream.next_block().await.unwrap().is_none());
208 assert!(stream.next_block().await.unwrap().is_none());
209 assert_eq!(stream.next_block_index(), 5);
210 }
211 }
212
213 #[tokio::test]
214 async fn test_compactor_block_stream_retry_budget() {
215 for budget in [0, 1, 2] {
216 let (mut stream, original) = test_stream(budget).await;
217 assert_next_block(&mut stream, &original, 1).await;
218 for attempt in 0..budget {
219 inject_stream(
220 &mut stream,
221 vec![Err(ObjectError::internal("injected read failure"))],
222 );
223 assert_next_block(&mut stream, &original, 2 + attempt).await;
224 assert_eq!(stream.io_retry_times, attempt + 1);
225 }
226 let failed_index = stream.next_block_index();
227 inject_stream(
228 &mut stream,
229 vec![Err(ObjectError::internal("retry budget exhausted"))],
230 );
231 assert!(stream.next_block().await.unwrap_err().is_object_error());
232 assert_eq!(stream.next_block_index(), failed_index);
233 assert_eq!(stream.io_retry_times, budget);
234 }
235 }
236}