Skip to main content

vector/sinks/doris/
config.rs

1//! Configuration for the `Doris` sink.
2
3use super::sink::DorisSink;
4
5use crate::{
6    codecs::EncodingConfigWithFraming,
7    http::{Auth, HttpClient},
8    sinks::{
9        doris::{
10            client::DorisSinkClient, common::DorisCommon, health::DorisHealthLogic,
11            retry::DorisRetryLogic, service::DorisService,
12        },
13        prelude::*,
14        util::{RealtimeSizeBasedDefaultBatchSettings, UriSerde, service::HealthConfig},
15    },
16    template::ConfinementConfig,
17};
18use futures;
19use futures_util::TryFutureExt;
20use std::collections::HashMap;
21use std::sync::Arc;
22
23/// Configuration for the `doris` sink.
24#[configurable_component(sink("doris", "Deliver log data to an Apache Doris database."))]
25#[derive(Clone, Debug)]
26#[serde(deny_unknown_fields)]
27pub struct DorisConfig {
28    /// A list of Doris endpoints to send logs to.
29    ///
30    /// The endpoint must contain an HTTP scheme, and may specify a
31    /// hostname or IP address and port.
32    #[serde(default)]
33    #[configurable(metadata(docs::examples = "http://127.0.0.1:8030"))]
34    pub endpoints: Vec<UriSerde>,
35
36    /// The database that contains the table data will be inserted into.
37    #[configurable(metadata(docs::examples = "mydatabase"))]
38    pub database: Template,
39
40    /// The table data is inserted into.
41    #[configurable(metadata(docs::examples = "mytable"))]
42    pub table: Template,
43
44    /// The prefix for Stream Load label.
45    /// The final label will be in format: `{label_prefix}_{database}_{table}_{timestamp}_{uuid}`.
46    #[configurable(metadata(docs::examples = "vector"))]
47    #[serde(default = "default_label_prefix")]
48    pub label_prefix: String,
49
50    /// Enable request logging.
51    #[serde(default, skip_serializing_if = "crate::serde::is_default")]
52    pub log_request: bool,
53
54    /// Custom HTTP headers to add to the request.
55    ///
56    /// These headers can be used to set Doris-specific Stream Load parameters:
57    /// - `format`: Data format (json, csv.)
58    /// - `read_json_by_line`: Whether to read JSON line by line
59    /// - `strip_outer_array`: Whether to strip outer array brackets
60    /// - Column mappings and transformations
61    ///
62    /// See [Doris Stream Load documentation](https://doris.apache.org/docs/data-operate/import/import-way/stream-load-manual)
63    /// for all available parameters.
64    #[serde(default)]
65    #[configurable(metadata(docs::additional_props_description = "An HTTP header value."))]
66    pub headers: HashMap<String, String>,
67
68    #[serde(flatten)]
69    pub encoding: EncodingConfigWithFraming,
70
71    /// Compression algorithm to use for HTTP requests.
72    #[serde(default)]
73    pub compression: Compression,
74
75    /// Number of retries attempted before failing.
76    #[serde(default = "default_max_retries")]
77    pub max_retries: isize,
78
79    #[configurable(derived)]
80    #[serde(default)]
81    pub batch: BatchConfig<RealtimeSizeBasedDefaultBatchSettings>,
82
83    #[configurable(derived)]
84    pub auth: Option<Auth>,
85
86    #[serde(default)]
87    #[configurable(derived)]
88    pub request: TowerRequestConfig,
89
90    #[configurable(derived)]
91    pub tls: Option<TlsConfig>,
92
93    /// Options for determining the health of Doris endpoints.
94    #[serde(default)]
95    #[configurable(derived)]
96    #[serde(rename = "distribution")]
97    pub endpoint_health: Option<HealthConfig>,
98
99    #[configurable(derived)]
100    #[serde(
101        default,
102        deserialize_with = "crate::serde::bool_or_struct",
103        skip_serializing_if = "crate::serde::is_default"
104    )]
105    pub acknowledgements: AcknowledgementsConfig,
106
107    #[configurable(derived)]
108    #[serde(flatten)]
109    pub confinement: ConfinementConfig,
110}
111
112fn default_label_prefix() -> String {
113    "vector".to_string()
114}
115
116const fn default_max_retries() -> isize {
117    -1
118}
119
120impl Default for DorisConfig {
121    fn default() -> Self {
122        Self {
123            endpoints: Vec::new(),
124            database: Template::try_from("").unwrap(),
125            table: Template::try_from("").unwrap(),
126            label_prefix: default_label_prefix(),
127            log_request: false,
128            headers: HashMap::new(),
129            encoding: (
130                Some(vector_lib::codecs::encoding::FramingConfig::NewlineDelimited),
131                vector_lib::codecs::JsonSerializerConfig::default(),
132            )
133                .into(),
134            compression: Compression::default(),
135            max_retries: default_max_retries(),
136            batch: BatchConfig::default(),
137            auth: None,
138            request: TowerRequestConfig::default(),
139            tls: None,
140            endpoint_health: None,
141            acknowledgements: AcknowledgementsConfig::default(),
142            confinement: ConfinementConfig::default(),
143        }
144    }
145}
146
147impl_generate_config_from_default!(DorisConfig);
148
149#[async_trait::async_trait]
150#[typetag::serde(name = "doris")]
151impl SinkConfig for DorisConfig {
152    async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
153        let endpoints = self.endpoints.clone();
154
155        if endpoints.is_empty() {
156            return Err("No endpoints configured.'.".into());
157        }
158        let commons = DorisCommon::parse_many(self).await?;
159        let common = commons[0].clone();
160
161        let client = HttpClient::new(common.tls_settings.clone(), &cx.proxy)?;
162
163        // Setup retry logic using the configured request settings
164        let request_settings = self.request.into_settings();
165
166        let health_config = self.endpoint_health.clone().unwrap_or_default();
167
168        let services_futures = commons
169            .iter()
170            .map(|common| {
171                let client_clone = client.clone();
172                let compression = self.compression;
173                let label_prefix = self.label_prefix.clone();
174                let headers = self.headers.clone();
175                let log_request = self.log_request;
176                let base_url = common.base_url.clone();
177                let auth = common.auth.clone();
178
179                async move {
180                    let endpoint = base_url.to_string();
181
182                    let doris_client = DorisSinkClient::new(
183                        client_clone,
184                        base_url,
185                        auth,
186                        compression,
187                        label_prefix,
188                        headers,
189                    )
190                    .await;
191
192                    let doris_client_safe = doris_client.into_thread_safe();
193
194                    let service = DorisService::new(doris_client_safe, log_request);
195
196                    Ok::<_, crate::Error>((endpoint, service))
197                }
198            })
199            .collect::<Vec<_>>();
200
201        // Wait for all futures to complete
202        let services_results = futures::future::join_all(services_futures).await;
203
204        // Filter out successful results
205        let services = services_results
206            .into_iter()
207            .filter_map(Result::ok)
208            .collect::<Vec<_>>();
209
210        let service = request_settings.distributed_service(
211            DorisRetryLogic {},
212            services,
213            health_config,
214            DorisHealthLogic,
215            1, // Buffer bound is hardcoded to 1 for sinks
216        );
217
218        // Create DorisSink with the configured service
219        let sink = DorisSink::new(service, self, &common)?;
220
221        let sink = VectorSink::from_event_streamsink(sink);
222
223        // Create a shared client instance to avoid repeated creation
224        let healthcheck_doris_client = {
225            let doris_client = DorisSinkClient::new(
226                client.clone(),
227                common.base_url.clone(),
228                common.auth.clone(),
229                self.compression,
230                self.label_prefix.clone(),
231                self.headers.clone(),
232            )
233            .await;
234            doris_client.into_thread_safe()
235        };
236
237        // Use the previously saved client for health check, no need to create a new instance
238        let healthcheck = futures::future::select_ok(commons.into_iter().map(move |common| {
239            let client = Arc::clone(&healthcheck_doris_client);
240            async move { common.healthcheck(client).await }.boxed()
241        }))
242        .map_ok(|((), _)| ())
243        .boxed();
244
245        self.confinement.set_confinement_gauge("sink", Self::NAME);
246        Ok((sink, healthcheck))
247    }
248
249    fn input(&self) -> Input {
250        Input::log()
251    }
252
253    fn acknowledgements(&self) -> &AcknowledgementsConfig {
254        &self.acknowledgements
255    }
256}
257
258#[cfg(test)]
259mod tests {
260    use super::*;
261
262    #[test]
263    fn generate_config() {
264        crate::test_util::test_generate_config::<DorisConfig>();
265    }
266
267    #[test]
268    fn test_default_values() {
269        assert_eq!(default_label_prefix(), "vector");
270        assert_eq!(default_max_retries(), -1);
271    }
272
273    #[test]
274    fn confinement_rejects_unconfined_database_template() {
275        let template = Template::try_from("{{ tenant }}").unwrap();
276        let config = ConfinementConfig::default();
277        let result = template.confine(&config, "doris", "database");
278        assert!(
279            result.is_err(),
280            "bare template with no literal prefix must be rejected"
281        );
282    }
283
284    #[test]
285    fn confinement_opt_out_allows_unconfined_database_template() {
286        let template = Template::try_from("{{ tenant }}").unwrap();
287        let config = ConfinementConfig {
288            dangerously_allow_unconfined_template_resolution: true,
289        };
290        let result = template.confine(&config, "doris", "database");
291        assert!(result.is_ok(), "opt-out must allow bare template");
292    }
293
294    #[test]
295    fn confinement_blocks_dotdot_escape_at_render() {
296        use crate::event::LogEvent;
297        use vrl::event_path;
298        let template = Template::try_from("mydb_{{ tenant }}").unwrap();
299        let config = ConfinementConfig::default();
300        let confined = template.confine(&config, "doris", "database").unwrap();
301        let mut event = LogEvent::default();
302        event.insert(event_path!("tenant"), "/../evil");
303        let result = confined.render_string(&crate::event::Event::Log(event));
304        assert!(
305            result.is_err(),
306            "dotdot escape in rendered value must be rejected by prefix check"
307        );
308    }
309}