1use std::sync::Arc;
18use std::time::{Duration, Instant};
19
20use common_base::readable_size::ReadableSize;
21use common_telemetry::{debug, info};
22use futures::AsyncWriteExt;
23use object_store::ObjectStore;
24use snafu::ResultExt;
25use store_api::storage::RegionId;
26use tokio::sync::mpsc::{UnboundedSender, unbounded_channel};
27
28use crate::access_layer::{
29 FilePathProvider, Metrics, OperationType, RegionFilePathFactory, SstInfoArray, SstWriteRequest,
30 TempFileCleaner, WriteCachePathProvider, WriteType, new_fs_cache_store,
31};
32use crate::cache::file_cache::{FileCache, FileCacheRef, FileType, IndexKey, IndexValue};
33use crate::cache::manifest_cache::ManifestCache;
34use crate::error::{self, Result};
35use crate::metrics::UPLOAD_BYTES_TOTAL;
36use crate::region::opener::RegionLoadCacheTask;
37use crate::sst::file::RegionFileId;
38use crate::sst::index::IndexerBuilderImpl;
39use crate::sst::index::intermediate::IntermediateManager;
40use crate::sst::index::puffin_manager::{PuffinManagerFactory, SstPuffinManager};
41use crate::sst::parquet::WriteOptions;
42use crate::sst::parquet::writer::ParquetWriter;
43use crate::sst::{DEFAULT_WRITE_BUFFER_SIZE, DEFAULT_WRITE_CONCURRENCY};
44
45pub trait WriteCacheUploadStoreWrapper: Send + Sync {
47 fn wrap(&self, store: ObjectStore, op_type: OperationType) -> ObjectStore;
54}
55
56pub type WriteCacheUploadStoreWrapperRef = Arc<dyn WriteCacheUploadStoreWrapper>;
57
58pub struct WriteCache {
62 file_cache: FileCacheRef,
64 puffin_manager_factory: PuffinManagerFactory,
66 intermediate_manager: IntermediateManager,
68 task_sender: UnboundedSender<RegionLoadCacheTask>,
70 manifest_cache: Option<ManifestCache>,
72 upload_store_wrapper: Option<WriteCacheUploadStoreWrapperRef>,
74}
75
76pub type WriteCacheRef = Arc<WriteCache>;
77
78impl WriteCache {
79 #[allow(clippy::too_many_arguments)]
82 pub async fn new(
83 local_store: ObjectStore,
84 cache_capacity: ReadableSize,
85 ttl: Option<Duration>,
86 index_cache_percent: Option<u8>,
87 enable_background_worker: bool,
88 puffin_manager_factory: PuffinManagerFactory,
89 intermediate_manager: IntermediateManager,
90 manifest_cache: Option<ManifestCache>,
91 ) -> Result<Self> {
92 let (task_sender, task_receiver) = unbounded_channel();
93
94 let file_cache = Arc::new(FileCache::new(
95 local_store,
96 cache_capacity,
97 ttl,
98 index_cache_percent,
99 enable_background_worker,
100 ));
101 file_cache.recover(false, Some(task_receiver)).await;
102
103 Ok(Self {
104 file_cache,
105 puffin_manager_factory,
106 intermediate_manager,
107 task_sender,
108 manifest_cache,
109 upload_store_wrapper: None,
110 })
111 }
112
113 #[allow(clippy::too_many_arguments)]
115 pub async fn new_fs(
116 cache_dir: &str,
117 cache_capacity: ReadableSize,
118 ttl: Option<Duration>,
119 index_cache_percent: Option<u8>,
120 enable_background_worker: bool,
121 puffin_manager_factory: PuffinManagerFactory,
122 intermediate_manager: IntermediateManager,
123 manifest_cache_capacity: ReadableSize,
124 ) -> Result<Self> {
125 info!("Init write cache on {cache_dir}, capacity: {cache_capacity}");
126
127 let local_store = new_fs_cache_store(cache_dir).await?;
128
129 let manifest_cache = if manifest_cache_capacity.as_bytes() > 0 {
131 Some(ManifestCache::new(local_store.clone(), manifest_cache_capacity, ttl, false).await)
132 } else {
133 None
134 };
135
136 Self::new(
137 local_store,
138 cache_capacity,
139 ttl,
140 index_cache_percent,
141 enable_background_worker,
142 puffin_manager_factory,
143 intermediate_manager,
144 manifest_cache,
145 )
146 .await
147 }
148
149 pub(crate) fn file_cache(&self) -> FileCacheRef {
151 self.file_cache.clone()
152 }
153
154 pub(crate) fn manifest_cache(&self) -> Option<ManifestCache> {
156 self.manifest_cache.clone()
157 }
158
159 pub(crate) fn with_upload_store_wrapper(
161 mut self,
162 upload_store_wrapper: Option<WriteCacheUploadStoreWrapperRef>,
163 ) -> Self {
164 self.upload_store_wrapper = upload_store_wrapper;
165 self
166 }
167
168 pub(crate) fn build_puffin_manager(&self) -> SstPuffinManager {
170 let store = self.file_cache.local_store();
171 let path_provider = WriteCachePathProvider::new(self.file_cache.clone());
172 self.puffin_manager_factory.build(store, path_provider)
173 }
174
175 pub(crate) async fn put_and_upload_sst(
177 &self,
178 data: &bytes::Bytes,
179 region_file_id: RegionFileId,
180 upload_request: SstUploadRequest,
181 write_buffer_size: ReadableSize,
182 ) -> Result<Metrics> {
183 let region_id = region_file_id.region_id();
184 let file_id = region_file_id.file_id();
185 let mut metrics = Metrics::new(WriteType::Flush);
186
187 let parquet_key = IndexKey::new(region_id, file_id, FileType::Parquet);
189
190 let cache_start = Instant::now();
192 let cache_path = self.file_cache.cache_file_path(parquet_key);
193 let store = self.file_cache.local_store();
194 let cleaner = TempFileCleaner::new(region_id, store.clone());
195 let write_res = store
196 .write_with(&cache_path, data.clone())
197 .chunk(write_buffer_size.as_bytes() as usize)
198 .await
199 .context(crate::error::OpenDalSnafu);
200 if let Err(e) = write_res {
201 cleaner.clean_by_file_id(file_id).await;
202 return Err(e);
203 }
204
205 metrics.write_batch = cache_start.elapsed();
206
207 let upload_start = Instant::now();
209 let remote_path = upload_request
210 .dest_path_provider
211 .build_sst_file_path(region_file_id);
212
213 if let Err(e) = self
214 .upload(
215 parquet_key,
216 &remote_path,
217 &upload_request.remote_store,
218 UploadOptions {
219 op_type: OperationType::Flush,
221 write_buffer_size,
222 },
223 )
224 .await
225 {
226 self.remove(parquet_key).await;
228 return Err(e);
229 }
230
231 metrics.upload_parquet = upload_start.elapsed();
232 Ok(metrics)
233 }
234
235 pub(crate) fn intermediate_manager(&self) -> &IntermediateManager {
237 &self.intermediate_manager
238 }
239
240 pub(crate) async fn write_and_upload_sst(
242 &self,
243 write_request: SstWriteRequest,
244 upload_request: SstUploadRequest,
245 write_opts: &WriteOptions,
246 metrics: &mut Metrics,
247 ) -> Result<SstInfoArray> {
248 let region_id = write_request.metadata.region_id;
249 let override_sequence = if write_request.preserve_row_sequence {
250 None
251 } else {
252 write_request.max_sequence
253 };
254
255 let store = self.file_cache.local_store();
256 let path_provider = WriteCachePathProvider::new(self.file_cache.clone());
257 let indexer = IndexerBuilderImpl {
258 build_type: write_request.op_type.into(),
259 metadata: write_request.metadata.clone(),
260 puffin_manager: self
261 .puffin_manager_factory
262 .build(store.clone(), path_provider.clone()),
263 write_cache_enabled: true,
264 intermediate_manager: self.intermediate_manager.clone(),
265 index_options: write_request.index_options,
266 inverted_index_config: write_request.inverted_index_config,
267 fulltext_index_config: write_request.fulltext_index_config,
268 bloom_filter_index_config: write_request.bloom_filter_index_config,
269 #[cfg(feature = "vector_index")]
270 vector_index_config: write_request.vector_index_config,
271 };
272
273 let cleaner = TempFileCleaner::new(region_id, store.clone());
274 let mut writer = ParquetWriter::new_with_object_store(
276 store.clone(),
277 write_request.metadata,
278 write_request.index_config,
279 indexer,
280 path_provider.clone(),
281 metrics,
282 )
283 .await
284 .with_file_cleaner(cleaner);
285
286 let sst_info = match write_request.sst_write_format {
287 crate::sst::FormatType::PrimaryKey => {
288 writer
289 .write_all_flat_as_primary_key(
290 write_request.source,
291 override_sequence,
292 write_opts,
293 )
294 .await?
295 }
296 crate::sst::FormatType::Flat => {
297 writer
298 .write_all_flat(write_request.source, override_sequence, write_opts)
299 .await?
300 }
301 };
302
303 if sst_info.is_empty() {
305 return Ok(sst_info);
306 }
307
308 let mut upload_tracker = UploadTracker::new(region_id);
309 let mut err = None;
310 let op_type = write_request.op_type;
311 let remote_store = &upload_request.remote_store;
312 for sst in &sst_info {
313 let parquet_key = IndexKey::new(region_id, sst.file_id, FileType::Parquet);
314 let parquet_path = upload_request
315 .dest_path_provider
316 .build_sst_file_path(RegionFileId::new(region_id, sst.file_id));
317 let start = Instant::now();
318 if let Err(e) = self
319 .upload(
320 parquet_key,
321 &parquet_path,
322 remote_store,
323 UploadOptions {
324 op_type,
325 write_buffer_size: write_opts.write_buffer_size,
326 },
327 )
328 .await
329 {
330 err = Some(e);
331 break;
332 }
333 metrics.upload_parquet += start.elapsed();
334 upload_tracker.push_uploaded_file(parquet_path);
335
336 if sst.index_metadata.file_size > 0 {
337 let puffin_key = IndexKey::new(region_id, sst.file_id, FileType::Puffin(0));
338 let puffin_path = upload_request
339 .dest_path_provider
340 .build_index_file_path(RegionFileId::new(region_id, sst.file_id));
341 let start = Instant::now();
342 if let Err(e) = self
343 .upload(
344 puffin_key,
345 &puffin_path,
346 remote_store,
347 UploadOptions {
348 op_type,
349 write_buffer_size: DEFAULT_WRITE_BUFFER_SIZE,
350 },
351 )
352 .await
353 {
354 err = Some(e);
355 break;
356 }
357 metrics.upload_puffin += start.elapsed();
358 upload_tracker.push_uploaded_file(puffin_path);
359 }
360 }
361
362 if let Some(err) = err {
363 upload_tracker
365 .clean(&sst_info, &self.file_cache, remote_store)
366 .await;
367 return Err(err);
368 }
369
370 Ok(sst_info)
371 }
372
373 pub(crate) async fn remove(&self, index_key: IndexKey) {
375 self.file_cache.remove(index_key).await
376 }
377
378 pub(crate) async fn download(
381 &self,
382 index_key: IndexKey,
383 remote_path: &str,
384 remote_store: &ObjectStore,
385 file_size: u64,
386 ) -> Result<()> {
387 self.file_cache
388 .download(index_key, remote_path, remote_store, file_size)
389 .await
390 }
391
392 pub(crate) async fn download_if_absent(
397 &self,
398 index_key: IndexKey,
399 remote_path: &str,
400 remote_store: &ObjectStore,
401 file_size: u64,
402 ) -> Result<bool> {
403 if self.file_cache.contains_key(&index_key) {
404 debug!(
405 "Skip downloading file already in write cache, region: {}, file: {}",
406 index_key.region_id, index_key.file_id
407 );
408 return Ok(false);
409 }
410
411 self.download(index_key, remote_path, remote_store, file_size)
412 .await?;
413 Ok(true)
414 }
415
416 pub(crate) async fn upload(
418 &self,
419 index_key: IndexKey,
420 upload_path: &str,
421 remote_store: &ObjectStore,
422 options: UploadOptions,
423 ) -> Result<()> {
424 let region_id = index_key.region_id;
425 let file_id = index_key.file_id;
426 let file_type = index_key.file_type;
427 let cache_path = self.file_cache.cache_file_path(index_key);
428
429 let start = Instant::now();
430 let cached_value = self
431 .file_cache
432 .local_store()
433 .stat(&cache_path)
434 .await
435 .context(error::OpenDalSnafu)?;
436 let reader = self
437 .file_cache
438 .local_store()
439 .reader(&cache_path)
440 .await
441 .context(error::OpenDalSnafu)?
442 .into_futures_async_read(0..cached_value.content_length())
443 .await
444 .context(error::OpenDalSnafu)?;
445
446 let upload_store = self.upload_store_wrapper.as_ref().map_or_else(
447 || remote_store.clone(),
448 |wrapper| wrapper.wrap(remote_store.clone(), options.op_type),
449 );
450 let mut writer = upload_store
451 .writer_with(upload_path)
452 .chunk(options.write_buffer_size.as_bytes() as usize)
453 .concurrent(DEFAULT_WRITE_CONCURRENCY)
454 .await
455 .context(error::OpenDalSnafu)?
456 .into_futures_async_write();
457
458 let bytes_written =
459 futures::io::copy(reader, &mut writer)
460 .await
461 .context(error::UploadSnafu {
462 region_id,
463 file_id,
464 file_type,
465 })?;
466
467 writer.close().await.context(error::UploadSnafu {
469 region_id,
470 file_id,
471 file_type,
472 })?;
473
474 UPLOAD_BYTES_TOTAL.inc_by(bytes_written);
475
476 debug!(
477 "Successfully upload file to remote, region: {}, file: {}, upload_path: {}, cost: {:?}",
478 region_id,
479 file_id,
480 upload_path,
481 start.elapsed(),
482 );
483
484 let index_value = IndexValue {
485 file_size: bytes_written as _,
486 };
487 self.file_cache.put(index_key, index_value).await;
489
490 Ok(())
491 }
492
493 pub(crate) fn load_region_cache(&self, task: RegionLoadCacheTask) {
497 let _ = self.task_sender.send(task);
498 }
499}
500
501pub struct SstUploadRequest {
503 pub dest_path_provider: RegionFilePathFactory,
505 pub remote_store: ObjectStore,
507}
508
509pub(crate) struct UploadOptions {
511 pub op_type: OperationType,
513 pub write_buffer_size: ReadableSize,
515}
516
517pub(crate) struct UploadTracker {
519 region_id: RegionId,
521 files_uploaded: Vec<String>,
523}
524
525impl UploadTracker {
526 pub(crate) fn new(region_id: RegionId) -> Self {
528 Self {
529 region_id,
530 files_uploaded: Vec::new(),
531 }
532 }
533
534 pub(crate) fn push_uploaded_file(&mut self, path: String) {
536 self.files_uploaded.push(path);
537 }
538
539 pub(crate) async fn clean(
541 &self,
542 sst_info: &SstInfoArray,
543 file_cache: &FileCacheRef,
544 remote_store: &ObjectStore,
545 ) {
546 common_telemetry::info!(
547 "Start cleaning files on upload failure, region: {}, num_ssts: {}",
548 self.region_id,
549 sst_info.len()
550 );
551
552 for sst in sst_info {
554 let parquet_key = IndexKey::new(self.region_id, sst.file_id, FileType::Parquet);
555 file_cache.remove(parquet_key).await;
556
557 if sst.index_metadata.file_size > 0 {
558 let puffin_key = IndexKey::new(
559 self.region_id,
560 sst.file_id,
561 FileType::Puffin(sst.index_metadata.version),
562 );
563 file_cache.remove(puffin_key).await;
564 }
565 }
566
567 for file_path in &self.files_uploaded {
569 if let Err(e) = remote_store.delete(file_path).await {
570 common_telemetry::error!(e; "Failed to delete file {}", file_path);
571 }
572 }
573 }
574}
575
576#[cfg(test)]
577mod tests {
578 use std::sync::Mutex;
579
580 use bytes::Bytes;
581 use common_test_util::temp_dir::create_temp_dir;
582 use object_store::services::Memory;
583 use object_store::{ATOMIC_WRITE_DIR, ObjectStore};
584 use parquet::file::metadata::PageIndexPolicy;
585 use store_api::region_request::PathType;
586 use store_api::storage::FileId;
587
588 use super::*;
589 use crate::access_layer::OperationType;
590 use crate::cache::file_cache::IndexValue;
591 use crate::cache::test_util::{assert_parquet_metadata_equal, new_fs_store};
592 use crate::cache::{CacheManager, CacheStrategy};
593 use crate::error::InvalidBatchSnafu;
594 use crate::read::FlatSource;
595 use crate::region::options::IndexOptions;
596 use crate::sst::parquet::reader::ParquetReaderBuilder;
597 use crate::test_util::TestEnv;
598 use crate::test_util::sst_util::{
599 WriteChunkRecorder, new_flat_source_from_record_batches, new_record_batch_by_range,
600 sst_file_handle_with_file_id, sst_region_metadata,
601 };
602
603 struct RedirectUploadStoreWrapper {
604 target_store: ObjectStore,
605 last_op_type: Mutex<Option<OperationType>>,
606 }
607
608 impl WriteCacheUploadStoreWrapper for RedirectUploadStoreWrapper {
609 fn wrap(&self, _store: ObjectStore, op_type: OperationType) -> ObjectStore {
610 *self.last_op_type.lock().unwrap() = Some(op_type);
611 self.target_store.clone()
612 }
613 }
614
615 #[tokio::test]
616 async fn test_upload_uses_wrapped_remote_store() {
617 let env = TestEnv::new().await;
618 let local_store = ObjectStore::new(Memory::default()).unwrap();
619 let original_store = ObjectStore::new(Memory::default()).unwrap();
620 let target_store = ObjectStore::new(Memory::default()).unwrap();
621 let wrapper = Arc::new(RedirectUploadStoreWrapper {
622 target_store: target_store.clone(),
623 last_op_type: Mutex::new(None),
624 });
625 let write_cache = WriteCache::new(
626 local_store.clone(),
627 ReadableSize::mb(10),
628 None,
629 None,
630 false,
631 env.get_puffin_manager(),
632 env.get_intermediate_manager(),
633 None,
634 )
635 .await
636 .unwrap()
637 .with_upload_store_wrapper(Some(wrapper.clone()));
638
639 let region_id = RegionId::new(1024, 1);
640 let file_id = FileId::random();
641 let key = IndexKey::new(region_id, file_id, FileType::Parquet);
642 let cache_path = write_cache.file_cache.cache_file_path(key);
643 let upload_path = "wrapped-upload.parquet";
644 let data = Bytes::from_static(b"wrapped upload data");
645 local_store.write(&cache_path, data.clone()).await.unwrap();
646
647 write_cache
648 .upload(
649 key,
650 upload_path,
651 &original_store,
652 UploadOptions {
653 op_type: OperationType::Compact,
654 write_buffer_size: DEFAULT_WRITE_BUFFER_SIZE,
655 },
656 )
657 .await
658 .unwrap();
659
660 assert_eq!(
661 *wrapper.last_op_type.lock().unwrap(),
662 Some(OperationType::Compact)
663 );
664
665 assert!(original_store.stat(upload_path).await.is_err());
666 assert_eq!(
667 target_store.read(upload_path).await.unwrap().to_vec(),
668 data.to_vec()
669 );
670 }
671
672 #[rstest::rstest]
673 #[tokio::test]
674 async fn test_write_and_upload_sst(
675 #[values(OperationType::Flush, OperationType::Compact)] op_type: OperationType,
676 ) {
677 let mut env = TestEnv::new().await;
680 let remote_chunks = WriteChunkRecorder::default();
681 let mock_store = env.init_object_store_manager().layer(remote_chunks.layer());
682 let path_provider = RegionFilePathFactory::new("test".to_string(), PathType::Bare);
683
684 let local_dir = create_temp_dir("");
685 let local_chunks = WriteChunkRecorder::default();
686 let local_store =
687 new_fs_store(local_dir.path().to_str().unwrap()).layer(local_chunks.layer());
688
689 let write_cache = env
690 .create_write_cache(local_store.clone(), ReadableSize::mb(10))
691 .await;
692
693 let metadata = Arc::new(sst_region_metadata());
695 let region_id = metadata.region_id;
696 let source = new_flat_source_from_record_batches(vec![
697 new_record_batch_by_range(&["a", "d"], 0, 60),
698 new_record_batch_by_range(&["b", "f"], 0, 40),
699 new_record_batch_by_range(&["b", "h"], 100, 200),
700 ]);
701
702 let write_request = SstWriteRequest {
703 op_type,
704 metadata,
705 source,
706 storage: None,
707 max_sequence: None,
708 sst_write_format: Default::default(),
709 cache_manager: Default::default(),
710 preserve_row_sequence: false,
711 index_options: IndexOptions::default(),
712 index_config: Default::default(),
713 inverted_index_config: Default::default(),
714 fulltext_index_config: Default::default(),
715 bloom_filter_index_config: Default::default(),
716 #[cfg(feature = "vector_index")]
717 vector_index_config: Default::default(),
718 };
719
720 let upload_request = SstUploadRequest {
721 dest_path_provider: path_provider.clone(),
722 remote_store: mock_store.clone(),
723 };
724
725 let write_opts = WriteOptions {
726 write_buffer_size: ReadableSize::kb(1),
727 row_group_size: 512,
728 ..Default::default()
729 };
730
731 let mut metrics = Metrics::new(WriteType::Flush);
733 let mut sst_infos = write_cache
734 .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
735 .await
736 .unwrap();
737 let sst_info = sst_infos.remove(0);
738
739 let file_id = sst_info.file_id;
740 let sst_upload_path =
741 path_provider.build_sst_file_path(RegionFileId::new(region_id, file_id));
742 let index_upload_path =
743 path_provider.build_index_file_path(RegionFileId::new(region_id, file_id));
744
745 let key = IndexKey::new(region_id, file_id, FileType::Parquet);
747 assert!(write_cache.file_cache.contains_key(&key));
748
749 let remote_data = mock_store.read(&sst_upload_path).await.unwrap();
751 let cache_data = local_store
752 .read(&write_cache.file_cache.cache_file_path(key))
753 .await
754 .unwrap();
755 assert_eq!(remote_data.to_vec(), cache_data.to_vec());
756 let chunk_size = write_opts.write_buffer_size.as_bytes() as usize;
757 assert!(remote_data.len() > chunk_size);
758 remote_chunks.assert_chunks(&sst_upload_path, chunk_size, remote_data.len());
759 local_chunks.assert_chunks(
760 &write_cache.file_cache.cache_file_path(key),
761 chunk_size,
762 cache_data.len(),
763 );
764
765 let index_key = IndexKey::new(region_id, file_id, FileType::Puffin(0));
767 assert!(write_cache.file_cache.contains_key(&index_key));
768
769 let remote_index_data = mock_store.read(&index_upload_path).await.unwrap();
770 let cache_index_data = local_store
771 .read(&write_cache.file_cache.cache_file_path(index_key))
772 .await
773 .unwrap();
774 assert_eq!(remote_index_data.to_vec(), cache_index_data.to_vec());
775 remote_chunks.assert_chunks(
776 &index_upload_path,
777 DEFAULT_WRITE_BUFFER_SIZE.as_bytes() as usize,
778 remote_index_data.len(),
779 );
780
781 let sst_index_key = IndexKey::new(region_id, file_id, FileType::Parquet);
783 write_cache.remove(sst_index_key).await;
784 assert!(!write_cache.file_cache.contains_key(&sst_index_key));
785 write_cache.remove(index_key).await;
786 assert!(!write_cache.file_cache.contains_key(&index_key));
787 }
788
789 #[tokio::test]
790 async fn test_read_metadata_from_write_cache() {
791 common_telemetry::init_default_ut_logging();
792 let mut env = TestEnv::new().await;
793 let data_home = env.data_home().display().to_string();
794 let mock_store = env.init_object_store_manager();
795
796 let local_dir = create_temp_dir("");
797 let local_path = local_dir.path().to_str().unwrap();
798 let local_store = new_fs_store(local_path);
799
800 let write_cache = env
802 .create_write_cache(local_store.clone(), ReadableSize::mb(10))
803 .await;
804 let cache_manager = Arc::new(
805 CacheManager::builder()
806 .write_cache(Some(write_cache.clone()))
807 .build(),
808 );
809 assert!(!cache_manager.sst_meta_cache_enabled());
810
811 let metadata = Arc::new(sst_region_metadata());
813
814 let source = new_flat_source_from_record_batches(vec![
815 new_record_batch_by_range(&["a", "d"], 0, 60),
816 new_record_batch_by_range(&["b", "f"], 0, 40),
817 new_record_batch_by_range(&["b", "h"], 100, 200),
818 ]);
819
820 let write_request = SstWriteRequest {
822 op_type: OperationType::Flush,
823 metadata,
824 source,
825 storage: None,
826 max_sequence: None,
827 sst_write_format: Default::default(),
828 cache_manager: cache_manager.clone(),
829 preserve_row_sequence: false,
830 index_options: IndexOptions::default(),
831 index_config: Default::default(),
832 inverted_index_config: Default::default(),
833 fulltext_index_config: Default::default(),
834 bloom_filter_index_config: Default::default(),
835 #[cfg(feature = "vector_index")]
836 vector_index_config: Default::default(),
837 };
838 let write_opts = WriteOptions {
839 row_group_size: 512,
840 ..Default::default()
841 };
842 let upload_request = SstUploadRequest {
843 dest_path_provider: RegionFilePathFactory::new(data_home.clone(), PathType::Bare),
844 remote_store: mock_store.clone(),
845 };
846
847 let mut metrics = Metrics::new(WriteType::Flush);
848 let mut sst_infos = write_cache
849 .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
850 .await
851 .unwrap();
852 let sst_info = sst_infos.remove(0);
853 let write_parquet_metadata = sst_info.file_metadata.unwrap();
854
855 let handle = sst_file_handle_with_file_id(sst_info.file_id, 0, 1000);
857 let builder = ParquetReaderBuilder::new(
858 data_home,
859 PathType::Bare,
860 handle.clone(),
861 mock_store.clone(),
862 )
863 .cache(CacheStrategy::EnableAll(cache_manager.clone()))
864 .page_index_policy(PageIndexPolicy::Optional);
865 let reader = builder.build().await.unwrap().unwrap();
866 let cached_write_parquet_metadata = crate::cache::CachedSstMeta::try_new(
867 "test.sst",
868 Arc::unwrap_or_clone(write_parquet_metadata),
869 )
870 .unwrap()
871 .parquet_metadata();
872
873 assert_parquet_metadata_equal(cached_write_parquet_metadata, reader.parquet_metadata());
875 }
876
877 #[tokio::test]
878 async fn test_write_cache_clean_tmp_files() {
879 common_telemetry::init_default_ut_logging();
880 let mut env = TestEnv::new().await;
881 let data_home = env.data_home().display().to_string();
882 let mock_store = env.init_object_store_manager();
883
884 let write_cache_dir = create_temp_dir("");
885 let write_cache_path = write_cache_dir.path().to_str().unwrap();
886 let write_cache = env
887 .create_write_cache_from_path(write_cache_path, ReadableSize::mb(10))
888 .await;
889
890 let cache_manager = Arc::new(
892 CacheManager::builder()
893 .write_cache(Some(write_cache.clone()))
894 .build(),
895 );
896
897 let metadata = Arc::new(sst_region_metadata());
899
900 let record_batch = new_record_batch_by_range(&["a", "d"], 0, 60);
902 let schema = record_batch.schema();
903 let iter = Box::new(
904 [
905 Ok(record_batch),
906 InvalidBatchSnafu {
907 reason: "Abort the writer",
908 }
909 .fail(),
910 ]
911 .into_iter(),
912 );
913 let source = FlatSource::new_iter(schema, iter);
914
915 let write_request = SstWriteRequest {
917 op_type: OperationType::Flush,
918 metadata,
919 source,
920 storage: None,
921 max_sequence: None,
922 sst_write_format: Default::default(),
923 cache_manager: cache_manager.clone(),
924 preserve_row_sequence: false,
925 index_options: IndexOptions::default(),
926 index_config: Default::default(),
927 inverted_index_config: Default::default(),
928 fulltext_index_config: Default::default(),
929 bloom_filter_index_config: Default::default(),
930 #[cfg(feature = "vector_index")]
931 vector_index_config: Default::default(),
932 };
933 let write_opts = WriteOptions {
934 row_group_size: 512,
935 ..Default::default()
936 };
937 let upload_request = SstUploadRequest {
938 dest_path_provider: RegionFilePathFactory::new(data_home.clone(), PathType::Bare),
939 remote_store: mock_store.clone(),
940 };
941
942 let mut metrics = Metrics::new(WriteType::Flush);
943 write_cache
944 .write_and_upload_sst(write_request, upload_request, &write_opts, &mut metrics)
945 .await
946 .unwrap_err();
947 let atomic_write_dir = write_cache_dir.path().join(ATOMIC_WRITE_DIR);
948 let mut entries = tokio::fs::read_dir(&atomic_write_dir).await.unwrap();
949 let mut has_files = false;
950 while let Some(entry) = entries.next_entry().await.unwrap() {
951 if entry.file_type().await.unwrap().is_dir() {
952 continue;
953 }
954 has_files = true;
955 common_telemetry::warn!(
956 "Found remaining temporary file in atomic dir: {}",
957 entry.path().display()
958 );
959 }
960
961 assert!(!has_files);
962 }
963
964 #[tokio::test]
965 async fn test_download_if_absent_skips_when_cached() {
966 let mut env = TestEnv::new().await;
967 let remote_store = env.init_object_store_manager();
968
969 let local_dir = create_temp_dir("");
970 let local_store = new_fs_store(local_dir.path().to_str().unwrap());
971 let write_cache = env
972 .create_write_cache(local_store.clone(), ReadableSize::mb(10))
973 .await;
974
975 let region_id = RegionId::new(1024, 1);
976 let file_id = FileId::random();
977 let key = IndexKey::new(region_id, file_id, FileType::Parquet);
978 write_cache
979 .file_cache()
980 .put(key, IndexValue { file_size: 1 })
981 .await;
982
983 let downloaded = write_cache
984 .download_if_absent(key, "missing/path.parquet", &remote_store, 1)
985 .await
986 .unwrap();
987
988 assert!(!downloaded);
989 }
990
991 #[tokio::test]
992 async fn test_download_if_absent_downloads_when_missing() {
993 let mut env = TestEnv::new().await;
994 let remote_store = env.init_object_store_manager();
995
996 let local_dir = create_temp_dir("");
997 let local_store = new_fs_store(local_dir.path().to_str().unwrap());
998 let write_cache = env
999 .create_write_cache(local_store.clone(), ReadableSize::mb(10))
1000 .await;
1001
1002 let region_id = RegionId::new(1024, 2);
1003 let file_id = FileId::random();
1004 let key = IndexKey::new(region_id, file_id, FileType::Parquet);
1005 let remote_path = format!("download-if-absent/{file_id}.parquet");
1006 let remote_data = Bytes::from_static(b"download-if-absent-test");
1007 remote_store
1008 .write(&remote_path, remote_data.clone())
1009 .await
1010 .unwrap();
1011
1012 let downloaded = write_cache
1013 .download_if_absent(key, &remote_path, &remote_store, remote_data.len() as u64)
1014 .await
1015 .unwrap();
1016
1017 assert!(downloaded);
1018 assert!(write_cache.file_cache().contains_key(&key));
1019
1020 let cached_data = local_store
1021 .read(&write_cache.file_cache().cache_file_path(key))
1022 .await
1023 .unwrap();
1024 assert_eq!(cached_data.to_vec(), remote_data.to_vec());
1025 }
1026}