risingwave_meta/stream/
sink.rs1use anyhow::Context;
16use risingwave_common::catalog::ColumnCatalog;
17use risingwave_connector::dispatch_sink;
18use risingwave_connector::sink::catalog::SinkCatalog;
19use risingwave_connector::sink::trivial::BLACKHOLE_SINK;
20use risingwave_connector::sink::{CONNECTOR_TYPE_KEY, Sink, SinkParam, build_sink};
21use risingwave_pb::catalog::PbSink;
22
23use crate::{MetaError, MetaResult};
24
25#[await_tree::instrument(boxed)]
26pub async fn validate_sink(prost_sink_catalog: &PbSink) -> MetaResult<()> {
27 let sink_catalog = SinkCatalog::from(prost_sink_catalog);
28 reject_variant_sink(&sink_catalog)?;
29 let param = SinkParam::try_from_sink_catalog(sink_catalog)?;
30
31 let sink = build_sink(param)?;
32 sink.validate_unknown_fields()?;
33
34 dispatch_sink!(
35 sink,
36 sink,
37 Ok(sink.validate().await.context("failed to validate sink")?)
38 )
39}
40
41pub fn first_variant_column(columns: &[ColumnCatalog]) -> Option<&ColumnCatalog> {
44 columns
45 .iter()
46 .find(|column| column.data_type().contains_variant())
47}
48
49fn reject_variant_sink(sink_catalog: &SinkCatalog) -> MetaResult<()> {
51 if sink_catalog.target_table.is_some() {
53 return Ok(());
54 }
55 if let Some(connector) = sink_catalog.properties.get(CONNECTOR_TYPE_KEY)
56 && connector.eq_ignore_ascii_case(BLACKHOLE_SINK)
57 {
58 return Ok(());
59 }
60
61 if let Some(column) = first_variant_column(sink_catalog.full_columns()) {
62 return Err(MetaError::invalid_parameter(format!(
63 "sinking VARIANT columns is not supported yet: column `{}`",
64 column.name_with_hidden()
65 )));
66 }
67
68 Ok(())
69}
70
71#[cfg(test)]
72mod tests {
73 use std::collections::BTreeMap;
74
75 use risingwave_common::catalog::{ColumnCatalog, ColumnDesc, ColumnId};
76 use risingwave_common::types::{DataType, StructType};
77 use risingwave_pb::catalog::{PbSink, SinkType as PbSinkType};
78
79 use super::*;
80
81 fn make_sink(connector: &str, data_type: DataType) -> SinkCatalog {
82 SinkCatalog::from(PbSink {
83 columns: vec![
84 ColumnCatalog::visible(ColumnDesc::named("payload", ColumnId::new(0), data_type))
85 .to_protobuf(),
86 ],
87 properties: BTreeMap::from([(CONNECTOR_TYPE_KEY.to_owned(), connector.to_owned())]),
88 sink_type: PbSinkType::AppendOnly as i32,
89 ..Default::default()
90 })
91 }
92
93 #[test]
94 fn test_reject_variant_sink() {
95 for connector in ["kafka", "doris", "jdbc", "mongodb"] {
96 for data_type in [
97 DataType::Variant,
98 DataType::list(DataType::Variant),
99 DataType::Struct(StructType::new(vec![("v", DataType::Variant)])),
100 ] {
101 let sink = make_sink(connector, data_type.clone());
102 let err = reject_variant_sink(&sink).unwrap_err();
103 assert!(
104 err.to_string()
105 .contains("sinking VARIANT columns is not supported yet: column `payload`"),
106 "{connector} / {data_type:?}: {err:?}"
107 );
108 }
109 }
110 }
111
112 #[test]
113 fn test_allow_sink_without_variant() {
114 let sink = make_sink("kafka", DataType::Jsonb);
115 reject_variant_sink(&sink).unwrap();
116 }
117
118 #[test]
119 fn test_allow_variant_for_non_encoding_sinks() {
120 let sink = make_sink(BLACKHOLE_SINK, DataType::Variant);
121 reject_variant_sink(&sink).unwrap();
122
123 let mut sink_into_table = make_sink("table", DataType::Variant);
124 sink_into_table.target_table = Some(1.into());
125 reject_variant_sink(&sink_into_table).unwrap();
126 }
127
128 fn column(name: &str, data_type: DataType) -> ColumnCatalog {
129 ColumnCatalog::visible(ColumnDesc::named(name, ColumnId::new(0), data_type))
130 }
131
132 #[test]
133 fn test_first_variant_column() {
134 assert!(first_variant_column(&[]).is_none());
135 assert!(
136 first_variant_column(&[column("a", DataType::Int32), column("b", DataType::Jsonb)])
137 .is_none()
138 );
139
140 for data_type in [
141 DataType::Variant,
142 DataType::list(DataType::Variant),
143 DataType::Struct(StructType::new(vec![("v", DataType::Variant)])),
144 ] {
145 let columns = [
146 column("a", DataType::Int32),
147 column("v", data_type.clone()),
148 column("z", DataType::Variant),
149 ];
150 let found = first_variant_column(&columns).expect("should find variant column");
151 assert_eq!(found.name(), "v", "{data_type:?}");
152 }
153 }
154}