index/fulltext_index/create/
bloom_filter.rs1use std::collections::HashMap;
16use std::sync::Arc;
17use std::sync::atomic::AtomicUsize;
18
19use async_trait::async_trait;
20use common_error::ext::BoxedError;
21use puffin::puffin_manager::{PuffinWriter, PutOptions};
22use snafu::{OptionExt, ResultExt};
23use tokio_util::compat::{TokioAsyncReadCompatExt, TokioAsyncWriteCompatExt};
24
25use crate::bloom_filter::creator::BloomFilterCreator;
26use crate::external_provider::ExternalTempFileProvider;
27use crate::fulltext_index::create::FulltextIndexCreator;
28use crate::fulltext_index::error::{
29 AbortedSnafu, BiErrorsSnafu, BloomFilterFinishSnafu, ExternalSnafu, PuffinAddBlobSnafu, Result,
30 SerializeToJsonSnafu,
31};
32use crate::fulltext_index::tokenizer::{Analyzer, ChineseTokenizer, EnglishTokenizer};
33use crate::fulltext_index::{Config, KEY_FULLTEXT_CONFIG};
34
35const PIPE_BUFFER_SIZE_FOR_SENDING_BLOB: usize = 8192;
36
37pub struct BloomFilterFulltextIndexCreator {
39 inner: Option<BloomFilterCreator>,
40 analyzer: Analyzer,
41 config: Config,
42}
43
44impl BloomFilterFulltextIndexCreator {
45 pub fn new(
46 config: Config,
47 rows_per_segment: usize,
48 false_positive_rate: f64,
49 intermediate_provider: Arc<dyn ExternalTempFileProvider>,
50 global_memory_usage: Arc<AtomicUsize>,
51 global_memory_usage_threshold: Option<usize>,
52 ) -> Self {
53 let tokenizer = match config.analyzer {
54 crate::fulltext_index::Analyzer::English => Box::new(EnglishTokenizer) as _,
55 crate::fulltext_index::Analyzer::Chinese => Box::new(ChineseTokenizer) as _,
56 };
57 let analyzer = Analyzer::new(tokenizer, config.case_sensitive);
58
59 let inner = BloomFilterCreator::new(
60 rows_per_segment,
61 false_positive_rate,
62 intermediate_provider,
63 global_memory_usage,
64 global_memory_usage_threshold,
65 );
66 Self {
67 inner: Some(inner),
68 analyzer,
69 config,
70 }
71 }
72}
73
74#[async_trait]
75impl FulltextIndexCreator for BloomFilterFulltextIndexCreator {
76 async fn push_text(&mut self, text: &str) -> Result<()> {
77 let mut token_buf = Vec::new();
78 let hashes = self.analyzer.analyze_text_hashes(text, &mut token_buf);
79 self.inner
80 .as_mut()
81 .context(AbortedSnafu)?
82 .push_row_hashes(hashes)
83 .await
84 .map_err(BoxedError::new)
85 .context(ExternalSnafu)?;
86 Ok(())
87 }
88
89 async fn finish(
90 &mut self,
91 puffin_writer: &mut (impl PuffinWriter + Send),
92 blob_key: &str,
93 mut put_options: PutOptions,
94 ) -> Result<u64> {
95 put_options.compression = None;
98
99 let creator = self.inner.as_mut().context(AbortedSnafu)?;
100
101 let (tx, rx) = tokio::io::duplex(PIPE_BUFFER_SIZE_FOR_SENDING_BLOB);
102
103 let property_key = KEY_FULLTEXT_CONFIG.to_string();
104 let property_value = serde_json::to_string(&self.config).context(SerializeToJsonSnafu)?;
105
106 let (index_finish, puffin_add_blob) = futures::join!(
107 creator.finish(tx.compat_write()),
108 puffin_writer.put_blob(
109 blob_key,
110 rx.compat(),
111 put_options,
112 HashMap::from([(property_key, property_value)]),
113 )
114 );
115
116 match (
117 puffin_add_blob.context(PuffinAddBlobSnafu),
118 index_finish.context(BloomFilterFinishSnafu),
119 ) {
120 (Err(e1), Err(e2)) => BiErrorsSnafu {
121 first: Box::new(e1),
122 second: Box::new(e2),
123 }
124 .fail()?,
125
126 (Ok(_), e @ Err(_)) => e?,
127 (e @ Err(_), Ok(_)) => e.map(|_| ())?,
128 (Ok(written_bytes), Ok(_)) => {
129 return Ok(written_bytes);
130 }
131 }
132 Ok(0)
133 }
134
135 async fn abort(&mut self) -> Result<()> {
136 self.inner.take().context(AbortedSnafu)?;
137 Ok(())
138 }
139
140 fn memory_usage(&self) -> usize {
141 self.inner
142 .as_ref()
143 .map(|i| i.memory_usage())
144 .unwrap_or_default()
145 }
146}