Execute a SQL statement. ALTER TABLE is handled by Paimon directly; everything else is delegated to DataFusion.
(&self, sql: &str)
| 292 | /// Execute a SQL statement. ALTER TABLE is handled by Paimon directly; |
| 293 | /// everything else is delegated to DataFusion. |
| 294 | pub async fn sql(&self, sql: &str) -> DFResult<DataFrame> { |
| 295 | let is_create_table = looks_like_create_table(sql); |
| 296 | let (rewritten_sql, partition_keys) = if is_create_table { |
| 297 | extract_partition_by(sql)? |
| 298 | } else { |
| 299 | (sql.to_string(), vec![]) |
| 300 | }; |
| 301 | if contains_time_travel_keyword(&rewritten_sql) { |
| 302 | // Time-travel queries are not DDL; skip our own parsing and handle directly. |
| 303 | return self.handle_time_travel_query(&rewritten_sql).await; |
| 304 | } |
| 305 | |
| 306 | let statements = Parser::parse_sql(&GenericDialect {}, &rewritten_sql) |
| 307 | .map_err(|e| DataFusionError::Plan(format!("SQL parse error: {e}")))?; |
| 308 | |
| 309 | if statements.len() != 1 { |
| 310 | return Err(DataFusionError::Plan( |
| 311 | "Expected exactly one SQL statement".to_string(), |
| 312 | )); |
| 313 | } |
| 314 | |
| 315 | match &statements[0] { |
| 316 | Statement::CreateTable(create_table) => { |
| 317 | if create_table.temporary { |
| 318 | self.handle_create_temp_table(create_table).await |
| 319 | } else { |
| 320 | let (catalog, _catalog_name, _) = |
| 321 | self.resolve_catalog_and_table(&create_table.name)?; |
| 322 | self.handle_create_table(&catalog, create_table, partition_keys) |
| 323 | .await |
| 324 | } |
| 325 | } |
| 326 | Statement::AlterTable(alter_table) => { |
| 327 | let (catalog, _catalog_name, _) = |
| 328 | self.resolve_catalog_and_table(&alter_table.name)?; |
| 329 | self.handle_alter_table( |
| 330 | &catalog, |
| 331 | &alter_table.name, |
| 332 | &alter_table.operations, |
| 333 | alter_table.if_exists, |
| 334 | ) |
| 335 | .await |
| 336 | } |
| 337 | Statement::Merge(merge) => self.handle_merge_into(merge).await, |
| 338 | Statement::Update(update) => self.handle_update(update).await, |
| 339 | Statement::Delete(delete) => self.handle_delete(delete).await, |
| 340 | Statement::Insert(insert) |
| 341 | if insert.overwrite |
| 342 | && insert.partitioned.as_ref().is_some_and(|p| !p.is_empty()) => |
| 343 | { |
| 344 | self.handle_insert_overwrite_partition(insert).await |
| 345 | } |
| 346 | Statement::Set(Set::SingleAssignment { |
| 347 | variable, values, .. |
| 348 | }) => { |
| 349 | let key = variable.to_string(); |
| 350 | let key = key.trim_matches('\'').trim_matches('"'); |
| 351 | if let Some(paimon_key) = key.strip_prefix("paimon.") { |