Skip to main content

risingwave_meta/stream/
sink.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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
41/// Returns the first column whose type contains VARIANT (nested included). No sink connector can
42/// encode variant values yet, so callers reject such columns at creation time.
43pub fn first_variant_column(columns: &[ColumnCatalog]) -> Option<&ColumnCatalog> {
44    columns
45        .iter()
46        .find(|column| column.data_type().contains_variant())
47}
48
49/// Reject sinks that carry VARIANT columns at creation time.
50fn reject_variant_sink(sink_catalog: &SinkCatalog) -> MetaResult<()> {
51    // Sinks into tables and blackhole sinks do not serialize values with an external encoder.
52    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}