vector/sinks/util/
partitioner.rs1use vector_lib::{event::Event, partition::Partitioner};
2
3use crate::{
4 internal_events::{KeyOutsideBasePrefixError, TemplateRenderingError},
5 template::{Template, TemplateRenderingError as TplRenderError},
6};
7
8pub(crate) fn render_key_with_fallback(
17 template: &Template,
18 event: &Event,
19 dead_letter: Option<&str>,
20) -> Option<String> {
21 match template.render_string(event) {
22 Ok(key) => Some(key),
23 Err(TplRenderError::Confined {
24 rendered_preview,
25 rendered_len,
26 message,
27 }) => {
28 emit!(KeyOutsideBasePrefixError {
29 key_preview: &rendered_preview,
30 key_len: rendered_len,
31 message: &message,
32 });
33 None
34 }
35 Err(error) => {
36 if let Some(fallback) = dead_letter {
37 emit!(TemplateRenderingError {
38 error,
39 field: Some("key_prefix"),
40 drop_event: false,
41 });
42 Some(fallback.to_owned())
43 } else {
44 emit!(TemplateRenderingError {
45 error,
46 field: Some("key_prefix"),
47 drop_event: true,
48 });
49 None
50 }
51 }
52 }
53}
54
55pub struct KeyPartitioner {
61 key_prefix_template: Template,
62 dead_letter_key_prefix: Option<String>,
63}
64
65impl KeyPartitioner {
66 pub const fn new(
67 key_prefix_template: Template,
68 dead_letter_key_prefix: Option<String>,
69 ) -> Self {
70 Self {
71 key_prefix_template,
72 dead_letter_key_prefix,
73 }
74 }
75}
76
77impl Partitioner for KeyPartitioner {
78 type Item = Event;
79 type Key = Option<String>;
80
81 fn partition(&self, item: &Self::Item) -> Self::Key {
82 render_key_with_fallback(
83 &self.key_prefix_template,
84 item,
85 self.dead_letter_key_prefix.as_deref(),
86 )
87 }
88}