Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 0 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

23 changes: 11 additions & 12 deletions src/meta/api/src/api_impl/schema_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,26 +171,26 @@ pub async fn get_history_table_metas(
Ok(tb_metas)
}

pub async fn construct_drop_table_txn_operations(
pub(crate) struct VersionedTable {
pub id: TableId,
pub meta: SeqV<TableMeta>,
}

pub(crate) async fn construct_drop_table_txn_operations(
kv_api: &(impl kvapi::KVApi<Error = MetaError> + ?Sized),
table_name: String,
tenant: &Tenant,
catalog_name: Option<String>,
table_id: u64,
existing: VersionedTable,
db_id: u64,
if_exists: bool,
if_delete: bool,
txn: &mut TxnRequest,
) -> Result<(u64, u64), KVAppError> {
let tbid = TableId { table_id };

// Check if table exists.
let (tb_meta_seq, tb_meta) = kv_api.get_pb_seq_and_value(&tbid).await?;
if tb_meta_seq == 0 {
return Err(KVAppError::AppError(AppError::UnknownTableId(
UnknownTableId::new(table_id, "drop_table_by_id failed to find valid tb_meta"),
)));
}
let VersionedTable { id: tbid, meta } = existing;
let table_id = tbid.table_id;
let tb_meta_seq = meta.seq;
let mut tb_meta = meta.data;

// Get db name, tenant name and related info for tx.
let table_id_to_name = TableIdToName { table_id };
Expand Down Expand Up @@ -232,7 +232,6 @@ pub async fn construct_drop_table_txn_operations(
"drop table by id"
);

let mut tb_meta = tb_meta.unwrap();
// drop a table with drop_on time
if tb_meta.drop_on.is_some() {
return if if_exists {
Expand Down
71 changes: 66 additions & 5 deletions src/meta/api/src/api_impl/table_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ use databend_common_meta_app::app_error::MultiStmtTxnCommitFailed;
use databend_common_meta_app::app_error::StreamAlreadyExists;
use databend_common_meta_app::app_error::StreamVersionMismatched;
use databend_common_meta_app::app_error::TableAlreadyExists;
use databend_common_meta_app::app_error::TableEngineMismatch;
use databend_common_meta_app::app_error::TableSnapshotExpired;
use databend_common_meta_app::app_error::TableVersionMismatched;
use databend_common_meta_app::app_error::UndropTableHasNoHistory;
Expand Down Expand Up @@ -128,6 +129,7 @@ use super::database_api::DatabaseApi;
use super::database_util::get_db_or_err;
use super::garbage_collection_api::ORPHAN_POSTFIX;
use super::garbage_collection_api::get_history_tables_for_gc;
use super::schema_api::VersionedTable;
use super::schema_api::build_upsert_table_deduplicated_label;
use super::schema_api::construct_drop_table_txn_operations;
use super::schema_api::get_db_by_id_or_err;
Expand Down Expand Up @@ -361,6 +363,24 @@ where
});
}
CreateOption::CreateOrReplace => {
// CREATE OR REPLACE must not replace an existing table with a
// different engine.
let existing_table_meta = self
.get_pb(&TableId::new(*id.data))
.await?
.ok_or_else(|| {
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
*id.data,
"create or replace failed to find existing table meta",
)))
})?;
TableEngineMismatch::ensure(
req.table_name(),
&existing_table_meta.engine,
&req.table_meta.engine,
)
.map_err(|e| KVAppError::AppError(e.into()))?;

if req.as_dropped {
// If the table is being created as a dropped table, we do not
// need to combine with drop_table_txn operations, just return
Expand All @@ -369,12 +389,16 @@ where

SeqV::new(id.seq, *id.data)
} else {
let existing = VersionedTable {
id: TableId::new(*id.data),
meta: existing_table_meta,
};
let (seq, id) = construct_drop_table_txn_operations(
self,
req.name_ident.table_name.clone(),
&req.name_ident.tenant,
req.catalog_name.clone(),
*id.data,
existing,
*seq_db_id.data,
true,
false,
Expand Down Expand Up @@ -605,12 +629,23 @@ where

let mut txn = TxnRequest::default();

let tbid = TableId::new(table_id);
let table_meta = self.get_pb(&tbid).await?.ok_or_else(|| {
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
table_id,
"drop_table_by_id failed to find valid tb_meta",
)))
})?;
let existing = VersionedTable {
id: tbid,
meta: table_meta,
};
let opt = construct_drop_table_txn_operations(
self,
req.table_name.clone(),
&req.tenant,
None,
table_id,
existing,
req.db_id,
req.if_exists,
true,
Expand Down Expand Up @@ -1153,6 +1188,12 @@ where
// get tb_meta of the last table id
let tbid = TableId { table_id };
let (tb_meta_seq, tb_meta) = self.get_pb_seq_and_value(&tbid).await?;
let mut tb_meta = tb_meta.ok_or_else(|| {
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
table_id,
"commit_table_meta",
)))
})?;

debug!(
ident :% =(&tbid),
Expand All @@ -1162,7 +1203,29 @@ where

{
// reset drop on time
let mut tb_meta = tb_meta.unwrap();
let mut txn_req = TxnRequest::default();

if let Some(prev_table_id) = req.prev_table_id {
let prev_tbid = TableId::new(prev_table_id);
let previous = self.get_pb(&prev_tbid).await?.ok_or_else(|| {
KVAppError::AppError(AppError::UnknownTableId(UnknownTableId::new(
prev_table_id,
"commit_table_meta previous table",
)))
})?;

TableEngineMismatch::ensure(
req.name_ident.table_name.as_str(),
&previous.engine,
&tb_meta.engine,
)
.map_err(|e| KVAppError::AppError(e.into()))?;

txn_req
.condition
.push(txn_cond_seq(&prev_tbid, Eq, previous.seq));
}

// undrop a table with no drop_on time
if tb_meta.drop_on.is_none() {
return Err(KVAppError::AppError(AppError::UndropTableWithNoDropTime(
Expand All @@ -1171,8 +1234,6 @@ where
}
tb_meta.drop_on = None;

let mut txn_req = TxnRequest::default();

txn_req.condition.extend([
// db has not to change, i.e., no new table is created.
// Renaming db is OK and does not affect the seq of db_meta.
Expand Down
38 changes: 38 additions & 0 deletions src/meta/app/src/app_error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -195,6 +195,38 @@ impl TableAlreadyExists {
}
}

#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[error(
"Cannot replace '{table_name}': existing table uses engine {existing_engine}, but the new table uses engine {new_engine}"
)]
pub struct TableEngineMismatch {
table_name: String,
existing_engine: String,
new_engine: String,
}

impl TableEngineMismatch {
pub fn new(
table_name: impl Into<String>,
existing_engine: impl Into<String>,
new_engine: impl Into<String>,
) -> Self {
Self {
table_name: table_name.into(),
existing_engine: existing_engine.into(),
new_engine: new_engine.into(),
}
}

pub fn ensure(table_name: &str, existing_engine: &str, new_engine: &str) -> Result<(), Self> {
if existing_engine.eq_ignore_ascii_case(new_engine) {
return Ok(());
}

Err(Self::new(table_name, existing_engine, new_engine))
}
}

#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[error("ViewAlreadyExists: {view_name} while {context}")]
pub struct ViewAlreadyExists {
Expand Down Expand Up @@ -895,6 +927,9 @@ pub enum AppError {
#[error(transparent)]
TableAlreadyExists(#[from] TableAlreadyExists),

#[error(transparent)]
TableEngineMismatch(#[from] TableEngineMismatch),

#[error(transparent)]
ViewAlreadyExists(#[from] ViewAlreadyExists),

Expand Down Expand Up @@ -1199,6 +1234,8 @@ impl AppErrorMessage for TableAlreadyExists {
}
}

impl AppErrorMessage for TableEngineMismatch {}

impl AppErrorMessage for ViewAlreadyExists {
fn message(&self) -> String {
format!("'{}' as view Already Exists", self.view_name)
Expand Down Expand Up @@ -1412,6 +1449,7 @@ impl From<AppError> for ErrorCode {
ErrorCode::MaterializedViewAlreadyExists(err.message())
}
AppError::TableAlreadyExists(err) => ErrorCode::TableAlreadyExists(err.message()),
AppError::TableEngineMismatch(err) => ErrorCode::TableEngineNotSupported(err.message()),
AppError::ViewAlreadyExists(err) => ErrorCode::ViewAlreadyExists(err.message()),
AppError::CreateTableWithDropTime(err) => {
ErrorCode::CreateTableWithDropTime(err.message())
Expand Down
Loading