1use std::sync::Arc;
16use std::time::Duration;
17
18use bytes::{Bytes, BytesMut};
19use fail::fail_point;
20use futures::{StreamExt, stream};
21use opendal::layers::{RetryLayer, TimeoutLayer};
22use opendal::raw::BoxedStaticFuture;
23use opendal::services::Memory;
24use opendal::{Execute, Executor, Operator, Writer};
25use risingwave_common::config::ObjectStoreConfig;
26use risingwave_common::range::RangeBoundsExt;
27use thiserror_ext::AsReport;
28
29use crate::object::object_metrics::ObjectStoreMetrics;
30use crate::object::{
31 MonitoredObjectStore, ObjectDataStream, ObjectError, ObjectMetadata, ObjectMetadataIter,
32 ObjectRangeBounds, ObjectResult, ObjectStore, OperationType, StreamingUploader, prefix,
33};
34
35#[derive(Clone)]
37pub struct OpendalObjectStore {
38 pub(crate) op: Operator,
39 pub(crate) media_type: MediaType,
40 pub(crate) config: Arc<ObjectStoreConfig>,
41 pub(crate) metrics: Arc<ObjectStoreMetrics>,
42}
43
44#[derive(Clone)]
45pub enum MediaType {
46 Memory,
47 Hdfs,
48 Gcs,
49 Minio,
50 S3,
51 Obs,
52 Oss,
53 Webhdfs,
54 Azblob,
55 Fs,
56}
57
58impl MediaType {
59 pub fn as_str(&self) -> &'static str {
60 match self {
61 MediaType::Memory => "Memory",
62 MediaType::Hdfs => "Hdfs",
63 MediaType::Gcs => "Gcs",
64 MediaType::Minio => "Minio",
65 MediaType::S3 => "S3",
66 MediaType::Obs => "Obs",
67 MediaType::Oss => "Oss",
68 MediaType::Webhdfs => "Webhdfs",
69 MediaType::Azblob => "Azblob",
70 MediaType::Fs => "Fs",
71 }
72 }
73}
74
75impl OpendalObjectStore {
76 pub fn test_new_memory_engine() -> ObjectResult<Self> {
78 let builder = Memory::default();
80 let op: Operator = Operator::new(builder)?;
81
82 Ok(Self {
83 op,
84 media_type: MediaType::Memory,
85 config: Arc::new(ObjectStoreConfig::default()),
86 metrics: Arc::new(ObjectStoreMetrics::unused()),
87 })
88 }
89}
90
91#[inline]
92fn timestamp_to_secs(t: opendal::raw::Timestamp) -> f64 {
93 t.into_inner().as_second() as f64
94}
95
96#[async_trait::async_trait]
97impl ObjectStore for OpendalObjectStore {
98 type StreamingUploader = OpendalStreamingUploader;
99
100 fn get_object_prefix(&self, obj_id: u64, use_new_object_prefix_strategy: bool) -> String {
101 match self.media_type {
102 MediaType::S3 => prefix::s3::get_object_prefix(obj_id),
103 MediaType::Minio => prefix::s3::get_object_prefix(obj_id),
104 MediaType::Memory => String::default(),
105 MediaType::Hdfs
106 | MediaType::Gcs
107 | MediaType::Obs
108 | MediaType::Oss
109 | MediaType::Webhdfs
110 | MediaType::Azblob
111 | MediaType::Fs => {
112 prefix::opendal_engine::get_object_prefix(obj_id, use_new_object_prefix_strategy)
113 }
114 }
115 }
116
117 async fn upload(&self, path: &str, obj: Bytes) -> ObjectResult<()> {
118 if obj.is_empty() {
119 Err(ObjectError::internal("upload empty object"))
120 } else {
121 self.op.write(path, obj).await?;
122 Ok(())
123 }
124 }
125
126 async fn streaming_upload(&self, path: &str) -> ObjectResult<Self::StreamingUploader> {
127 if self.op.info().capability().write_can_multi {
128 Ok(OpendalStreamingUploader::Native(Box::new(
129 OpendalNativeStreamingUploader::new(
130 self.op.clone(),
131 path.to_owned(),
132 self.config.clone(),
133 self.metrics.clone(),
134 self.store_media_type(),
135 )
136 .await?,
137 )))
138 } else {
139 Ok(OpendalStreamingUploader::Buffered(
140 OpendalBufferedStreamingUploader::new(
141 self.clone()
142 .monitored(self.metrics.clone(), self.config.clone()),
143 path.to_owned(),
144 ),
145 ))
146 }
147 }
148
149 async fn read(&self, path: &str, range: impl ObjectRangeBounds) -> ObjectResult<Bytes> {
150 let data = if range.is_full() {
151 self.op.read(path).await?
152 } else {
153 self.op
154 .read_with(path)
155 .range(range.map(|v| *v as u64))
156 .await?
157 };
158
159 if let Some(len) = range.len()
160 && len != data.len()
161 {
162 return Err(ObjectError::internal(format!(
163 "mismatched size: expected {}, found {} when reading {} at {:?}",
164 len,
165 data.len(),
166 path,
167 range,
168 )));
169 }
170
171 Ok(data.to_bytes())
172 }
173
174 async fn streaming_read(
178 &self,
179 path: &str,
180 range: impl ObjectRangeBounds,
181 ) -> ObjectResult<ObjectDataStream> {
182 fail_point!("opendal_streaming_read_err", |_| Err(
183 ObjectError::internal("opendal streaming read error")
184 ));
185 let range = range.map(|v| *v as u64);
186
187 let reader = self
192 .op
193 .clone()
194 .layer(TimeoutLayer::new().with_io_timeout(Duration::from_millis(
195 self.config.retry.streaming_read_attempt_timeout_ms,
196 )))
197 .layer(
198 RetryLayer::new()
199 .with_min_delay(Duration::from_millis(
200 self.config.retry.req_backoff_interval_ms,
201 ))
202 .with_max_delay(Duration::from_millis(
203 self.config.retry.req_backoff_max_delay_ms,
204 ))
205 .with_max_times(self.config.retry.streaming_read_retry_attempts)
206 .with_factor(self.config.retry.req_backoff_factor as f32)
207 .with_jitter(),
208 )
209 .reader_with(path)
210 .await?;
211 let stream = reader.into_bytes_stream(range).await?.map(|item| {
212 item.map(|b| Bytes::copy_from_slice(b.as_ref()))
213 .map_err(|e| {
214 ObjectError::internal(format!("reader into_stream fail {}", e.as_report()))
215 })
216 });
217
218 Ok(Box::pin(stream))
219 }
220
221 async fn metadata(&self, path: &str) -> ObjectResult<ObjectMetadata> {
222 let opendal_metadata = self.op.stat(path).await?;
223 let key = path.to_owned();
224 let last_modified = match opendal_metadata.last_modified() {
225 Some(t) => timestamp_to_secs(t),
226 None => 0_f64,
227 };
228
229 let total_size = opendal_metadata.content_length() as usize;
230 let metadata = ObjectMetadata {
231 key,
232 last_modified,
233 total_size,
234 };
235 Ok(metadata)
236 }
237
238 async fn delete(&self, path: &str) -> ObjectResult<()> {
239 self.op.delete(path).await?;
240 Ok(())
241 }
242
243 async fn delete_objects(&self, paths: &[String]) -> ObjectResult<()> {
246 self.op.delete_iter(paths.to_vec()).await?;
247 Ok(())
248 }
249
250 async fn list(
251 &self,
252 prefix: &str,
253 start_after: Option<String>,
254 limit: Option<usize>,
255 ) -> ObjectResult<ObjectMetadataIter> {
256 let mut object_lister = self.op.lister_with(prefix).recursive(true);
257 if let Some(start_after) = start_after {
258 object_lister = object_lister.start_after(&start_after);
259 }
260 let object_lister = object_lister.await?;
261
262 let op = self.op.clone();
263 let stream = stream::unfold(object_lister, move |mut object_lister| {
264 let op = op.clone();
265
266 async move {
267 match object_lister.next().await {
268 Some(Ok(object)) => {
269 let key = object.path().to_owned();
270
271 let meta = object.metadata();
276 let mut last_modified = meta.last_modified().map(timestamp_to_secs);
277 let mut total_size = meta.content_length() as usize;
278 if last_modified.is_none() || total_size == 0 {
279 let stat_meta = op.stat(&key).await.ok()?;
280 last_modified = stat_meta.last_modified().map(timestamp_to_secs);
281 total_size = stat_meta.content_length() as usize;
282 }
283
284 let metadata = ObjectMetadata {
285 key,
286 last_modified: last_modified.unwrap_or(0_f64),
287 total_size,
288 };
289 Some((Ok(metadata), object_lister))
290 }
291 Some(Err(err)) => Some((Err(err.into()), object_lister)),
292 None => None,
293 }
294 }
295 });
296
297 Ok(stream.take(limit.unwrap_or(usize::MAX)).boxed())
298 }
299
300 fn store_media_type(&self) -> &'static str {
301 self.media_type.as_str()
302 }
303}
304
305impl OpendalObjectStore {
306 pub async fn copy(&self, from_path: &str, to_path: &str) -> ObjectResult<()> {
307 self.op.copy(from_path, to_path).await?;
308 Ok(())
309 }
310}
311
312struct OpendalStreamingUploaderExecute {
313 metrics: Arc<ObjectStoreMetrics>,
315 media_type: &'static str,
316}
317
318impl OpendalStreamingUploaderExecute {
319 const STREAMING_UPLOAD_TYPE: OperationType = OperationType::StreamingUpload;
320
321 fn new(metrics: Arc<ObjectStoreMetrics>, media_type: &'static str) -> Self {
322 Self {
323 metrics,
324 media_type,
325 }
326 }
327}
328
329impl Execute for OpendalStreamingUploaderExecute {
330 fn execute(&self, f: BoxedStaticFuture<()>) {
331 let operation_type_str = Self::STREAMING_UPLOAD_TYPE.as_str();
332 let media_type = self.media_type;
333
334 let metrics = self.metrics.clone();
335 let _handle = tokio::spawn(async move {
336 let _timer = metrics
337 .operation_latency
338 .with_label_values(&[media_type, operation_type_str])
339 .start_timer();
340
341 f.await
342 });
343 }
344}
345
346pub enum OpendalStreamingUploader {
347 Native(Box<OpendalNativeStreamingUploader>),
348 Buffered(OpendalBufferedStreamingUploader),
349}
350
351pub struct OpendalNativeStreamingUploader {
353 writer: Writer,
354 buf: Vec<Bytes>,
357 not_uploaded_len: usize,
359 is_valid: bool,
361
362 abort_on_err: bool,
363
364 upload_part_size: usize,
365}
366
367impl OpendalNativeStreamingUploader {
368 pub async fn new(
369 op: Operator,
370 path: String,
371 config: Arc<ObjectStoreConfig>,
372 metrics: Arc<ObjectStoreMetrics>,
373 media_type: &'static str,
374 ) -> ObjectResult<Self> {
375 let monitored_execute = OpendalStreamingUploaderExecute::new(metrics, media_type);
376 let executor = Executor::with(monitored_execute);
377 let ctx = op.base_context().with_executor(executor);
381 let op = op.with_context(ctx);
382
383 let writer = op
384 .clone()
385 .layer(TimeoutLayer::new().with_io_timeout(Duration::from_millis(
386 config.retry.streaming_upload_attempt_timeout_ms,
387 )))
388 .layer(
389 RetryLayer::new()
390 .with_min_delay(Duration::from_millis(config.retry.req_backoff_interval_ms))
391 .with_max_delay(Duration::from_millis(config.retry.req_backoff_max_delay_ms))
392 .with_max_times(config.retry.streaming_upload_retry_attempts)
393 .with_factor(config.retry.req_backoff_factor as f32)
394 .with_jitter(),
395 )
396 .writer_with(&path)
397 .concurrent(config.opendal_upload_concurrency)
398 .await?;
399 Ok(Self {
400 writer,
401 buf: vec![],
402 not_uploaded_len: 0,
403 is_valid: true,
404 abort_on_err: config.opendal_writer_abort_on_err,
405 upload_part_size: config.upload_part_size,
406 })
407 }
408
409 async fn flush(&mut self) -> ObjectResult<()> {
410 let data: Vec<Bytes> = self.buf.drain(..).collect();
411 debug_assert_eq!(
412 data.iter().map(|b| b.len()).sum::<usize>(),
413 self.not_uploaded_len
414 );
415 if let Err(err) = self.writer.write(data).await {
416 self.is_valid = false;
417 if self.abort_on_err {
418 self.writer.abort().await?;
419 }
420 return Err(err.into());
421 }
422 self.not_uploaded_len = 0;
423 Ok(())
424 }
425}
426
427impl StreamingUploader for OpendalNativeStreamingUploader {
428 async fn write_bytes(&mut self, data: Bytes) -> ObjectResult<()> {
429 assert!(self.is_valid);
430 self.not_uploaded_len += data.len();
431 self.buf.push(data);
432 if self.not_uploaded_len >= self.upload_part_size {
433 self.flush().await?;
434 }
435 Ok(())
436 }
437
438 async fn finish(mut self) -> ObjectResult<()> {
439 assert!(self.is_valid);
440 if self.not_uploaded_len > 0 {
441 self.flush().await?;
442 }
443
444 assert!(self.buf.is_empty());
445 assert_eq!(self.not_uploaded_len, 0);
446
447 self.is_valid = false;
448 match self.writer.close().await {
449 Ok(_) => (),
450 Err(err) => {
451 if self.abort_on_err {
452 self.writer.abort().await?;
453 }
454 return Err(err.into());
455 }
456 };
457
458 Ok(())
459 }
460
461 fn get_memory_usage(&self) -> u64 {
463 self.not_uploaded_len as u64
464 }
465}
466
467pub struct OpendalBufferedStreamingUploader {
468 store: MonitoredObjectStore<OpendalObjectStore>,
469 path: String,
470 buf: BytesMut,
471}
472
473impl OpendalBufferedStreamingUploader {
474 fn new(store: MonitoredObjectStore<OpendalObjectStore>, path: String) -> Self {
475 Self {
476 store,
477 path,
478 buf: BytesMut::new(),
479 }
480 }
481}
482
483impl StreamingUploader for OpendalBufferedStreamingUploader {
484 async fn write_bytes(&mut self, data: Bytes) -> ObjectResult<()> {
485 self.buf.extend_from_slice(&data);
486 Ok(())
487 }
488
489 async fn finish(self) -> ObjectResult<()> {
490 self.store.upload(&self.path, self.buf.freeze()).await
491 }
492
493 fn get_memory_usage(&self) -> u64 {
494 self.buf.len() as u64
495 }
496}
497
498impl StreamingUploader for OpendalStreamingUploader {
499 async fn write_bytes(&mut self, data: Bytes) -> ObjectResult<()> {
500 match self {
501 Self::Native(uploader) => uploader.write_bytes(data).await,
502 Self::Buffered(uploader) => uploader.write_bytes(data).await,
503 }
504 }
505
506 async fn finish(self) -> ObjectResult<()> {
507 match self {
508 Self::Native(uploader) => (*uploader).finish().await,
509 Self::Buffered(uploader) => uploader.finish().await,
510 }
511 }
512
513 fn get_memory_usage(&self) -> u64 {
514 match self {
515 Self::Native(uploader) => uploader.get_memory_usage(),
516 Self::Buffered(uploader) => uploader.get_memory_usage(),
517 }
518 }
519}
520
521#[cfg(test)]
522mod tests {
523 use stream::TryStreamExt;
524
525 use super::*;
526
527 async fn list_all(prefix: &str, store: &OpendalObjectStore) -> Vec<ObjectMetadata> {
528 store
529 .list(prefix, None, None)
530 .await
531 .unwrap()
532 .try_collect::<Vec<_>>()
533 .await
534 .unwrap()
535 }
536
537 #[tokio::test]
538 async fn test_memory_upload() {
539 let block = Bytes::from("123456");
540 let store = OpendalObjectStore::test_new_memory_engine().unwrap();
541 store.upload("/abc", block).await.unwrap();
542
543 store.read("/ab", 0..3).await.unwrap_err();
545
546 let bytes = store.read("/abc", 4..6).await.unwrap();
547 assert_eq!(String::from_utf8(bytes.to_vec()).unwrap(), "56".to_owned());
548
549 store.delete("/abc").await.unwrap();
550
551 store.read("/abc", 0..3).await.unwrap_err();
553 }
554
555 #[tokio::test]
556 async fn test_memory_read_out_of_range_returns_error() {
557 let block = Bytes::from("123456");
558 let store = OpendalObjectStore::test_new_memory_engine().unwrap();
559 store.upload("/abc", block).await.unwrap();
560
561 store.read("/abc", 4..44).await.unwrap_err();
563 }
564
565 #[tokio::test]
566 async fn test_memory_metadata() {
567 let block = Bytes::from("123456");
568 let path = "/abc".to_owned();
569 let obj_store = OpendalObjectStore::test_new_memory_engine().unwrap();
570 obj_store.upload("/abc", block).await.unwrap();
571
572 let err = obj_store.metadata("/not_exist").await.unwrap_err();
573 assert!(err.is_object_not_found_error());
574 let metadata = obj_store.metadata("/abc").await.unwrap();
575 assert_eq!(metadata.total_size, 6);
576 obj_store.delete(&path).await.unwrap();
577 }
578
579 #[tokio::test]
580 async fn test_buffered_streaming_upload() {
581 let store = OpendalObjectStore::test_new_memory_engine().unwrap();
582 let mut uploader = OpendalBufferedStreamingUploader::new(
583 store
584 .clone()
585 .monitored(store.metrics.clone(), store.config.clone()),
586 "/abc".to_owned(),
587 );
588
589 uploader.write_bytes(Bytes::from("123")).await.unwrap();
590 assert_eq!(uploader.get_memory_usage(), 3);
591 uploader.write_bytes(Bytes::from("456")).await.unwrap();
592 assert_eq!(uploader.get_memory_usage(), 6);
593 uploader.finish().await.unwrap();
594
595 let read_obj = store.read("/abc", ..).await.unwrap();
596 assert_eq!(read_obj, Bytes::from("123456"));
597 }
598
599 #[tokio::test]
600 async fn test_empty_buffered_streaming_upload() {
601 let store = OpendalObjectStore::test_new_memory_engine().unwrap();
602 let uploader = OpendalBufferedStreamingUploader::new(
603 store
604 .clone()
605 .monitored(store.metrics.clone(), store.config.clone()),
606 "/abc".to_owned(),
607 );
608
609 uploader.finish().await.unwrap_err();
610 }
611
612 #[tokio::test]
613 async fn test_memory_delete_objects_and_list_object() {
614 let block1 = Bytes::from("123456");
615 let block2 = Bytes::from("987654");
616 let store = OpendalObjectStore::test_new_memory_engine().unwrap();
617 store.upload("abc", Bytes::from("123456")).await.unwrap();
618 store.upload("prefix/abc", block1).await.unwrap();
619
620 store.upload("prefix/xyz", block2).await.unwrap();
621
622 assert_eq!(list_all("", &store).await.len(), 3);
623 assert_eq!(list_all("prefix/", &store).await.len(), 2);
624 let str_list = [String::from("prefix/abc"), String::from("prefix/xyz")];
625
626 store.delete_objects(&str_list).await.unwrap();
627
628 assert!(store.read("prefix/abc/", ..).await.is_err());
629 assert!(store.read("prefix/xyz/", ..).await.is_err());
630 assert_eq!(list_all("", &store).await.len(), 1);
631 assert_eq!(list_all("prefix/", &store).await.len(), 0);
632 }
633}