Skip to main content

risingwave_storage/hummock/compactor/
block_stream.rs

1// Copyright 2026 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
26/// Streams a caller-selected range of physical SST blocks, resuming at block boundaries on I/O
27/// errors. Decoding and deciding which blocks may be copied belong to the consuming iterator.
28pub(super) struct SstableBlockStream {
29    pub(super) sstable: TableHolder,
30    pub(super) sstable_info: SstableInfo,
31    sstable_store: SstableStoreRef,
32    block_stream: Option<BlockDataStream>,
33    /// Absolute SST indices. Advance the start only after a complete block has been read.
34    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                // Opening a stream already uses the object store's initialization retry policy.
70                // An exhausted initialization failure propagates without consuming this budget.
71                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                    // Discard any partial block and reopen at its original offset. The budget
95                    // is cumulative for this SST iterator, not reset after each successful block.
96                    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        // Both ends exclude physical SST blocks, so recovery must preserve the selected range.
157        (
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            // Block 2 begins in the buffered packet. Its remaining bytes fail to arrive, so it
202            // must be reread in full from the backing object, with no duplicate or skipped block.
203            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}