std::string error_msg; auto *ctx= get_duckdb_context(thd);
/* Safety net: flush if prepare() was not called (no 2PC).
This is a no-op when appenders were already flushed. */ if (ctx->flush_appenders(error_msg))
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_APPEND_ERROR, error_msg.c_str(), "DuckDB");
ctx->duckdb_trans_rollback(error_msg); return1;
}
std::string query= "DROP SCHEMA IF EXISTS " + quote_duckdb_identifier(db.name);
if (duckdb_register_trx(thd))
DBUG_VOID_RETURN; auto *ctx= get_duckdb_context(thd); auto query_result= myduck::duckdb_query(ctx->get_connection(), query);
DBUG_VOID_RETURN;
}
if (batch_state == myduck::BatchState::UNDEFINED)
{ if (dml_in_batch)
batch_state= myduck::BatchState::IN_INSERT_ONLY_BATCH; else
batch_state= myduck::BatchState::NOT_IN_BATCH;
ctx->set_batch_state(batch_state);
} return batch_state;
}
/* Build duckdb type map of blob type */ staticvoid build_duckdb_blob_map(Field **field_list, MY_BITMAP *map)
{ for (Field **f_ptr= field_list; *f_ptr != nullptr; f_ptr++)
{
Field *field= *f_ptr;
enum_field_types type= field->real_type();
if (type == MYSQL_TYPE_SET || type == MYSQL_TYPE_ENUM ||
type == MYSQL_TYPE_BIT || type == MYSQL_TYPE_GEOMETRY ||
type == MYSQL_TYPE_VARCHAR || type == MYSQL_TYPE_STRING ||
type == MYSQL_TYPE_TINY_BLOB || type == MYSQL_TYPE_BLOB ||
type == MYSQL_TYPE_MEDIUM_BLOB || type == MYSQL_TYPE_LONG_BLOB)
{ if (FieldConvertor::convert_type(field) == "BLOB")
bitmap_set_bit(map, field->field_index);
}
}
}
/* ----- DML operations ----- */
int ha_duckdb::write_row(const uchar *buf)
{
DBUG_ENTER("ha_duckdb::write_row"); int ret= 0;
THD *thd= ha_thd();
Duckdb_share::Autoinc_range cur= share->autoinc_range.load(); for (;;)
{ /* Fast path: claim [cur.next, cur.next + need) with one CAS. */ while (cur.end - cur.next >= need)
{
Duckdb_share::Autoinc_range claimed= {cur.next + need, cur.end}; if (share->autoinc_range.compare_exchange_weak(cur, claimed))
{
*first_value= cur.next;
*nb_reserved_values= nb_desired_values;
DBUG_VOID_RETURN;
} /* CAS failure reloaded cur. */
}
mysql_mutex_lock(&share->autoinc_refill_mutex);
cur= share->autoinc_range.load(); if (cur.end - cur.next >= need)
{ /* Someone refilled while we waited; back to the fast path. */
mysql_mutex_unlock(&share->autoinc_refill_mutex); continue;
}
ulonglong want= std::max(need, autoinc_cache_size);
DatabaseTableNames dt(table->s->normalized_path.str);
ulonglong first= reserve_autoinc_block(ha_thd(), dt.db_name.c_str(),
dt.table_name.c_str(), want); if (first == 0)
{
mysql_mutex_unlock(&share->autoinc_refill_mutex); /* The only failure signal the server understands. */
*first_value= ULONGLONG_MAX;
*nb_reserved_values= 0;
DBUG_VOID_RETURN;
} /* Publishthenewblockminusourownshare.Aplainstorecannotclash withconcurrentfast-pathclaims:thosestillcomefromtheoldrange, whichliesentirelybelow`first`becausethesequenceismonotonic. Theremainderoftheoldblockisdropped,notmerged:idsbetweenthe twoblocksmayalreadybelongtoaDuckDB-sidewriter.Theunusedids becomegaps,whichAUTO_INCREMENTpermits.
*/
share->autoinc_range.store({first + need, first + want});
mysql_mutex_unlock(&share->autoinc_refill_mutex);
auto *ctx= get_duckdb_context(thd);
query_result= myduck::duckdb_query(ctx->get_connection(), query); if (query_result->HasError())
{
my_error(ER_GET_ERRMSG, MYF(0), HA_ERR_INTERNAL_ERROR, query_result->GetError().c_str(), "DuckDB");
DBUG_RETURN(HA_ERR_INTERNAL_ERROR);
}
current_chunk.reset();
current_row_index= 0;
for (auto *field= table->field; *field; ++field)
bitmap_set_bit(table->write_set, (*field)->field_index);
DBUG_RETURN(0);
}
int ha_duckdb::rnd_end()
{
DBUG_ENTER("ha_duckdb::rnd_end");
query_result.reset();
current_chunk.reset();
DBUG_RETURN(0);
}
int ha_duckdb::rnd_next(uchar *buf)
{
DBUG_ENTER("ha_duckdb::rnd_next");
THD *thd= ha_thd();
if (!query_result)
DBUG_RETURN(HA_ERR_INTERNAL_ERROR);
memset(buf, 0, table->s->reclength);
/* fetch new chunk when current chunk is empty */ if (!current_chunk || current_row_index >= current_chunk->size())
{
current_chunk.reset();
current_chunk= query_result->Fetch();
if (!current_chunk)
{ if (query_result->HasError())
{
my_error(ER_GET_ERRMSG, MYF(0), HA_ERR_INTERNAL_ERROR,
query_result->GetError().c_str(), "DuckDB");
DBUG_RETURN(HA_ERR_INTERNAL_ERROR);
}
DBUG_RETURN(HA_ERR_END_OF_FILE);
}
current_row_index= 0;
}
/* store the fields of a tuple */ for (size_t col_idx= 0; col_idx < current_chunk->ColumnCount(); ++col_idx)
{
duckdb::Value value= current_chunk->GetValue(col_idx, current_row_index);
Field *field= table->field[col_idx];
store_duckdb_field_in_mysql_format(field, value, thd);
}
/* update NULL field tag */ if (table->s->null_bytes > 0)
{ if (table->null_flags)
memcpy(buf, table->null_flags, table->s->null_bytes); else
memset(buf, 0, table->s->null_bytes);
}
int ha_duckdb::rnd_pos(uchar *, uchar *)
{
DBUG_ENTER("ha_duckdb::rnd_pos");
DBUG_RETURN(HA_ERR_WRONG_COMMAND);
}
int ha_duckdb::info(uint flag)
{
DBUG_ENTER("ha_duckdb::info"); if (flag & HA_STATUS_VARIABLE)
{ /* Retrieve variable info, such as row counts and file lengths */
stats.records= records();
stats.deleted= 0; // stats.data_file_length = // stats.index_file_length = // stats.delete_length =
stats.check_time= 0; // stats.mrr_length_per_rec =
// stats.data_file_length may be unset for TIAMAT; avoid division by // garbage. if (stats.records == 0 || stats.data_file_length == 0)
stats.mean_rec_length= 0; else
stats.mean_rec_length= (ulong) (stats.data_file_length / stats.records);
}
DBUG_RETURN(0);
}
ha_rows ha_duckdb::records()
{
DBUG_ENTER("ha_tiamat::records"); // Optimizer may call records()/info() in contexts where ha_share isn't // initialized for this handler instance. Return a conservative estimate. if (stats.records)
DBUG_RETURN(stats.records);
DBUG_RETURN(10);
}
int ha_duckdb::extra(enum ha_extra_function operation)
{
DBUG_ENTER("ha_duckdb::extra");
THD *thd= ha_thd(); auto *ctx= get_duckdb_context(thd);
switch (operation)
{ case HA_EXTRA_BEGIN_COPY:
ctx->set_in_copy_ddl(true); break; case HA_EXTRA_END_COPY: case HA_EXTRA_ABORT_COPY:
ctx->set_in_copy_ddl(false); break; default: break;
}
DBUG_RETURN(0);
}
int ha_duckdb::delete_all_rows()
{
DBUG_ENTER("ha_duckdb::delete_all_rows"); int ret= 0;
THD *thd= ha_thd();
ret= duckdb_register_trx(thd); if (ret)
DBUG_RETURN(ret);
auto *ctx= get_duckdb_context(thd);
/* Discard any pending batch rows for this table */
DatabaseTableNames dt(table->s->normalized_path.str);
ctx->delete_appender(dt.db_name, dt.table_name);
/* Execute DELETE FROM "schema"."table" */
std::string query= "DELETE FROM " + quote_duckdb_identifier(dt.db_name) + "." + quote_duckdb_identifier(dt.table_name);
auto query_result= myduck::duckdb_query(ctx->get_connection(), query); if (query_result->HasError())
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_DML_ERROR, query_result->GetError().c_str(), "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
int ha_duckdb::direct_delete_rows_init()
{
DBUG_ENTER("ha_duckdb::direct_delete_rows_init");
DBUG_RETURN(0);
}
int ha_duckdb::direct_delete_rows(ha_rows *delete_rows)
{
DBUG_ENTER("ha_duckdb::direct_delete_rows"); int ret= 0;
THD *thd= ha_thd();
LEX_STRING *source_query= thd_query_string(thd); if (myduck::mariadb_query_has_unsafe_quote_escape(
thd, source_query->str, source_query->length))
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_DML_ERROR, "Unsafe MariaDB backslash quote escape in forwarded SQL", "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
ret= duckdb_register_trx(thd); if (ret)
DBUG_RETURN(ret);
auto *ctx= get_duckdb_context(thd);
/* Flush any pending batch rows so DuckDB sees consistent data */
std::string error_msg; if (ctx->flush_appenders(error_msg))
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_APPEND_ERROR, error_msg.c_str(), "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
/* Execute the original DELETE statement in DuckDB */
LEX_STRING *qs= thd_query_string(thd);
std::string query(qs->str, qs->length); auto result= myduck::duckdb_query(thd, query, true); if (result->HasError())
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_DML_ERROR, result->GetError().c_str(), "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
/* DuckDB returns a single row with the count of affected rows */ auto chunk= result->Fetch(); if (chunk && chunk->size() > 0)
*delete_rows= chunk->GetValue(0, 0).GetValue<int64_t>(); else
*delete_rows= 0;
int ha_duckdb::direct_update_rows_init(List<Item> *update_fields
__attribute__((unused)))
{
DBUG_ENTER("ha_duckdb::direct_update_rows_init");
DBUG_RETURN(0);
}
int ha_duckdb::direct_update_rows(ha_rows *update_rows, ha_rows *found_rows)
{
DBUG_ENTER("ha_duckdb::direct_update_rows"); int ret= 0;
THD *thd= ha_thd();
LEX_STRING *source_query= thd_query_string(thd); if (myduck::mariadb_query_has_unsafe_quote_escape(
thd, source_query->str, source_query->length))
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_DML_ERROR, "Unsafe MariaDB backslash quote escape in forwarded SQL", "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
ret= duckdb_register_trx(thd); if (ret)
DBUG_RETURN(ret);
auto *ctx= get_duckdb_context(thd);
/* Flush any pending batch rows so DuckDB sees consistent data */
std::string error_msg; if (ctx->flush_appenders(error_msg))
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_APPEND_ERROR, error_msg.c_str(), "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
/* Execute the original UPDATE statement in DuckDB */
LEX_STRING *qs= thd_query_string(thd);
std::string query(qs->str, qs->length); auto result= myduck::duckdb_query(thd, query, true); if (result->HasError())
{
my_error(ER_GET_ERRMSG, MYF(0), HA_DUCKDB_DML_ERROR, result->GetError().c_str(), "DuckDB");
DBUG_RETURN(HA_DUCKDB_DML_ERROR);
}
/* DuckDB returns a single row with the count of affected rows */ auto chunk= result->Fetch();
ha_rows affected= 0; if (chunk && chunk->size() > 0)
affected= chunk->GetValue(0, 0).GetValue<int64_t>();
*update_rows= affected;
*found_rows= affected;
srv_duckdb_status.duckdb_rows_update+= affected;
DBUG_RETURN(0);
}
int ha_duckdb::external_lock(THD *thd, int lock_type)
{
DBUG_ENTER("ha_duckdb::external_lock"); if (lock_type != F_UNLCK)
{ /* DuckDB does not support XA transactions. Reject DML early. */ if (myduck::reject_xa_if_active(thd))
DBUG_RETURN(HA_ERR_WRONG_COMMAND);
int ret= duckdb_register_trx(thd); if (ret)
DBUG_RETURN(ret);
}
DBUG_RETURN(0);
}
if (database_changed(table->s->db.str, altered_table->s->db.str))
DBUG_RETURN(HA_ALTER_INPLACE_NOT_SUPPORTED);
if (ha_alter_info->alter_info->flags & ALTER_COLUMN_ORDER)
DBUG_RETURN(HA_ALTER_INPLACE_NOT_SUPPORTED);
if (ha_alter_info->error_if_not_empty)
DBUG_RETURN(HA_ALTER_INPLACE_NOT_SUPPORTED);
/* Reject ALTER on tables without PK when require_primary_key is ON */ if (myduck::require_primary_key && table->s->primary_key == MAX_KEY)
{
my_error(ER_REQUIRES_PRIMARY_KEY, MYF(0));
DBUG_RETURN(HA_ALTER_ERROR);
}
using DDL_convertor= std::unique_ptr<AlterTableConvertor>; using DDL_convertors= std::vector<DDL_convertor>;
DDL_convertor convertor;
DDL_convertors convertors;
std::vector<std::string> statements; for (auto &conv : convertors)
{ if (!conv || conv->check())
DBUG_RETURN(true);
std::string sql= conv->translate(); if (!sql.empty())
statements.push_back(std::move(sql));
}
if (statements.empty())
DBUG_RETURN(false);
/* A single MariaDB ALTER TABLE can produce multiple DuckDB statements.
Execute the generated operations atomically on a dedicated connection. */ auto con= myduck::DuckdbManager::CreateConnection(); auto query_result= myduck::duckdb_query(*con, "BEGIN"); if (query_result->HasError())
{
my_error(ER_GET_ERRMSG, MYF(0), HA_ERR_GENERIC,
query_result->GetError().c_str(), "DuckDB");
DBUG_RETURN(true);
}
static MYSQL_THDVAR_BOOL(cross_engine_ryow, PLUGIN_VAR_RQCMDARG, "In cross-engine joins, read own uncommitted writes " "from non-DuckDB tables via a direct handler scan in " "the parent transaction (disables predicate pushdown " "and index access path for those tables)",
NULL, NULL, FALSE);
/* ---- THDVAR accessor functions (used from duckdb_context.cc) ---- */
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.