(
&mut self,
)
| 3954 | } |
| 3955 | |
| 3956 | fn parse_create_kafka_sink_connection( |
| 3957 | &mut self, |
| 3958 | ) -> Result<CreateSinkConnection<Raw>, ParserError> { |
| 3959 | self.expect_keyword(CONNECTION)?; |
| 3960 | |
| 3961 | let connection = self.parse_raw_name()?; |
| 3962 | |
| 3963 | let options = if self.consume_token(&Token::LParen) { |
| 3964 | let options = self.parse_comma_separated(Parser::parse_kafka_sink_config_option)?; |
| 3965 | self.expect_token(&Token::RParen)?; |
| 3966 | options |
| 3967 | } else { |
| 3968 | vec![] |
| 3969 | }; |
| 3970 | |
| 3971 | // one token of lookahead: |
| 3972 | // * `KEY (` means we're parsing a list of columns for the key |
| 3973 | // * `KEY FORMAT` means there is no key, we'll parse a KeyValueFormat later |
| 3974 | let key = |
| 3975 | if self.peek_keyword(KEY) && self.peek_nth_token(1) != Some(Token::Keyword(FORMAT)) { |
| 3976 | let _ = self.expect_keyword(KEY); |
| 3977 | let key_columns = self.parse_parenthesized_column_list(Mandatory)?; |
| 3978 | |
| 3979 | let not_enforced = if self.peek_keywords(&[NOT, ENFORCED]) { |
| 3980 | self.expect_keywords(&[NOT, ENFORCED])?; |
| 3981 | true |
| 3982 | } else { |
| 3983 | false |
| 3984 | }; |
| 3985 | Some(SinkKey { |
| 3986 | key_columns, |
| 3987 | not_enforced, |
| 3988 | }) |
| 3989 | } else { |
| 3990 | None |
| 3991 | }; |
| 3992 | |
| 3993 | let headers = if self.parse_keyword(HEADERS) { |
| 3994 | Some(self.parse_identifier()?) |
| 3995 | } else { |
| 3996 | None |
| 3997 | }; |
| 3998 | |
| 3999 | Ok(CreateSinkConnection::Kafka { |
| 4000 | connection, |
| 4001 | options, |
| 4002 | key, |
| 4003 | headers, |
| 4004 | }) |
| 4005 | } |
| 4006 | |
| 4007 | fn parse_create_iceberg_sink_connection( |
| 4008 | &mut self, |
no test coverage detected