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