-
Notifications
You must be signed in to change notification settings - Fork 18
Expand file tree
/
Copy pathoutput_provider.rs
More file actions
120 lines (108 loc) · 4.57 KB
/
Copy pathoutput_provider.rs
File metadata and controls
120 lines (108 loc) · 4.57 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
use crate::runtime::output::Output;
use crate::runtime::task::OutputConfig;
pub struct OutputProvider;
impl OutputProvider {
pub fn from_output_configs(
output_configs: &[OutputConfig],
) -> Result<Vec<Box<dyn Output>>, Box<dyn std::error::Error + Send>> {
if output_configs.is_empty() {
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"Empty output configs list",
)) as Box<dyn std::error::Error + Send>);
}
const MAX_OUTPUTS: usize = 64;
if output_configs.len() > MAX_OUTPUTS {
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Too many outputs: {} (maximum is {})",
output_configs.len(),
MAX_OUTPUTS
),
)) as Box<dyn std::error::Error + Send>);
}
let mut outputs = Vec::new();
for (output_idx, output_config) in output_configs.iter().enumerate() {
let output = Self::from_output_config(output_config, output_idx)?;
outputs.push(output);
}
Ok(outputs)
}
fn from_output_config(
output_config: &OutputConfig,
output_idx: usize,
) -> Result<Box<dyn Output>, Box<dyn std::error::Error + Send>> {
match output_config {
OutputConfig::Kafka {
bootstrap_servers,
topic,
partition,
extra,
runtime: _,
} => {
use crate::runtime::output::output_runner::OutputRunner;
use crate::runtime::output::protocol::kafka::{
KafkaOutputProtocol, KafkaProducerConfig,
};
let servers: Vec<String> = bootstrap_servers
.split(',')
.map(|s| s.trim().to_string())
.filter(|s| !s.is_empty())
.collect();
if servers.is_empty() {
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Invalid bootstrap_servers in output config: empty or invalid (topic: {})",
topic
),
)) as Box<dyn std::error::Error + Send>);
}
let partition_opt = Some(*partition as i32);
let properties = extra.clone();
let kafka_config =
KafkaProducerConfig::new(servers, topic.clone(), partition_opt, properties);
let protocol = KafkaOutputProtocol::new(kafka_config);
let runtime = output_config.output_runtime_config();
Ok(Box::new(OutputRunner::new(protocol, output_idx, runtime)))
}
OutputConfig::Pulsar {
url,
topic,
extra,
runtime: _,
} => {
use crate::runtime::output::output_runner::OutputRunner;
use crate::runtime::output::protocol::pulsar::{
PulsarOutputProtocol, PulsarProducerConfig,
};
if url.is_empty() {
return Err(Box::new(std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!(
"Invalid pulsar url in output config: empty (topic: {})",
topic
),
)) as Box<dyn std::error::Error + Send>);
}
let pulsar_config =
PulsarProducerConfig::new(url.clone(), topic.clone(), extra.clone());
let protocol = PulsarOutputProtocol::new(pulsar_config);
let runtime = output_config.output_runtime_config();
Ok(Box::new(OutputRunner::new(protocol, output_idx, runtime)))
}
}
}
}