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::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
42/// Returns the first column whose type contains VARIANT (nested included). Callers use this to
43/// reject connectors that cannot encode variant values.
44pub fn first_variant_column(columns: &[ColumnCatalog]) -> Option<&ColumnCatalog> {
45    columns
46        .iter()
47        .find(|column| column.data_type().contains_variant())
48}
49
50/// Reject sinks that carry VARIANT columns at creation time unless their connector supports it.
51fn reject_variant_sink(sink_catalog: &SinkCatalog) -> MetaResult<()> {
52    // Sinks into tables and blackhole sinks do not serialize values with an external encoder;
53    // Iceberg uses the standard Parquet Variant physical layout.
54    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}