YoushouldhavereceivedacopyoftheGNUGeneralPublicLicense alongwiththisprogram;ifnot,writetotheFreeSoftware
Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02111-1301 USA */
#define MYSQL_SERVER 1
/* For use of 'PRIu64': */ #define __STDC_FORMAT_MACROS
#include <my_global.h>
#include <inttypes.h>
/* The C++ file's header */ #include"./ha_rocksdb.h"
// Internal MariaDB APIs not exposed in any header. extern"C" { /** *Gettheuserthread'sbinaryloggingformat *@paramthduserthread *@returnValuetobeusedasindexintothebinlog_format_namesarray
*/ int thd_binlog_format(const MYSQL_THD thd);
class Rdb_open_tables_map { private: /* Hash table used to track the handlers of open tables */
std::unordered_map<std::string, Rdb_table_handler *> m_table_map;
/* The mutex used to protect the hash table */ mutable mysql_mutex_t m_mutex;
staticint rocksdb_create_checkpoint(
THD *const thd MY_ATTRIBUTE((__unused__)), struct st_mysql_sys_var *const var MY_ATTRIBUTE((__unused__)), void *const save MY_ATTRIBUTE((__unused__)), struct st_mysql_value *const value) { char buf[FN_REFLEN]; int len = sizeof(buf); constchar *const checkpoint_dir_raw = value->val_str(value, buf, &len); if (checkpoint_dir_raw) { if (rdb != nullptr) {
std::string checkpoint_dir = rdb_normalize_dir(checkpoint_dir_raw); // NO_LINT_DEBUG
sql_print_information("RocksDB: creating checkpoint in directory : %s\n",
checkpoint_dir.c_str());
rocksdb::Checkpoint *checkpoint; auto status = rocksdb::Checkpoint::Create(rdb, &checkpoint); // We can only return HA_EXIT_FAILURE/HA_EXIT_SUCCESS here which is why // the return code is ignored, but by calling into rdb_error_to_mysql, // it will call my_error for us, which will propogate up to the client. int rc __attribute__((__unused__)); if (status.ok()) {
status = checkpoint->CreateCheckpoint(checkpoint_dir.c_str()); delete checkpoint; if (status.ok()) { // NO_LINT_DEBUG
sql_print_information( "RocksDB: created checkpoint in directory : %s\n",
checkpoint_dir.c_str()); return HA_EXIT_SUCCESS;
} else {
rc = ha_rocksdb::rdb_error_to_mysql(status);
}
} else {
rc = ha_rocksdb::rdb_error_to_mysql(status);
}
}
} return HA_EXIT_FAILURE;
}
/* This method is needed to indicate that the
ROCKSDB_CREATE_CHECKPOINT command is not read-only */ staticvoid rocksdb_create_checkpoint_stub(THD *const thd, struct st_mysql_sys_var *const var, void *const var_ptr, constvoid *const save) {}
rocksdb::Status s;
s = rdb->CompactFiles(c_options, cf_handle, file_names, 1);
// Due to a race, it's possible for CompactFiles to collide // with auto compaction, causing an error to return // regarding file not found. In that case, retry. if (s.IsInvalidArgument()) { continue;
}
if (!s.ok() && !s.IsAborted()) {
rdb_handle_io_error(s, RDB_IO_ERROR_GENERAL); return HA_EXIT_FAILURE;
} break;
} if (i == max_attempts) {
num_errors++;
}
}
/* DBOptions contains Statistics and needs to be destructed last */ static std::unique_ptr<rocksdb::BlockBasedTableOptions> rocksdb_tbl_options =
std::unique_ptr<rocksdb::BlockBasedTableOptions>( new rocksdb::BlockBasedTableOptions()); static std::unique_ptr<rocksdb::DBOptions> rocksdb_db_options =
rdb_init_rocksdb_db_options();
/* This enum needs to be kept up to date with rocksdb::TxnDBWritePolicy */ staticconstchar *write_policy_names[] = {"write_committed", "write_prepared", "write_unprepared", NullS};
#if0// MARIAROCKS_NOT_YET : read-free replication is not supported /* This array needs to be kept up to date with myrocks::read_free_rpl_type */ staticconstchar *read_free_rpl_names[] = {"OFF", "PK_ONLY", "PK_SK", NullS};
/* This enum needs to be kept up to date with rocksdb::InfoLogLevel */ staticconstchar *info_log_level_names[] = {"debug_level", "info_level", "warn_level", "error_level", "fatal_level", NullS};
RDB_MUTEX_LOCK_CHECK(rdb_sysvars_mutex);
rocksdb_db_options->statistics->set_stats_level( static_cast<rocksdb::StatsLevel>(
*static_cast<const uint64_t *>(save))); // Actual stats level is defined at rocksdb dbopt::statistics::stats_level_ // so adjusting rocksdb_stats_level here to make sure it points to // the correct stats level.
rocksdb_stats_level = rocksdb_db_options->statistics->get_stats_level();
RDB_MUTEX_UNLOCK_CHECK(rdb_sysvars_mutex);
}
if (rocksdb_reset_stats) {
rocksdb::Status s = rdb->ResetStats();
// RocksDB will always return success. Let's document this assumption here // as well so that we'll get immediately notified when contract changes.
DBUG_ASSERT(s == rocksdb::Status::OK());
s = rocksdb_stats->Reset();
DBUG_ASSERT(s == rocksdb::Status::OK());
}
enum rocksdb_flush_log_at_trx_commit_type : unsignedint {
FLUSH_LOG_NEVER = 0,
FLUSH_LOG_SYNC,
FLUSH_LOG_BACKGROUND,
FLUSH_LOG_MAX /* must be last */
};
staticint rocksdb_validate_flush_log_at_trx_commit(
THD *const thd, struct st_mysql_sys_var *const var, /* in: pointer to system variable */ void *var_ptr, /* out: immediate result for update function */ struct st_mysql_value *const value /* in: incoming value */) { longlong new_value;
/* value is NULL */ if (value->val_int(value, &new_value)) { return HA_EXIT_FAILURE;
}
if (rocksdb_db_options->allow_mmap_writes && new_value != FLUSH_LOG_NEVER) { return HA_EXIT_FAILURE;
}
// TODO: 0 means don't wait at all, and we don't support it yet? static MYSQL_THDVAR_ULONG(lock_wait_timeout, PLUGIN_VAR_RQCMDARG, "Number of seconds to wait for lock", nullptr,
nullptr, /*default*/ 1, /*min*/ 1, /*max*/ RDB_MAX_LOCK_WAIT_SECONDS, 0);
static MYSQL_THDVAR_ULONG(deadlock_detect_depth, PLUGIN_VAR_RQCMDARG, "Number of transactions deadlock detection will " "traverse through before assuming deadlock",
nullptr, nullptr, /*default*/ RDB_DEADLOCK_DETECT_DEPTH, /*min*/ 2, /*max*/ ULONG_MAX, 0);
static MYSQL_THDVAR_BOOL(
commit_time_batch_for_recovery, PLUGIN_VAR_RQCMDARG, "TransactionOptions::commit_time_batch_for_recovery for RocksDB", nullptr,
nullptr, TRUE);
static MYSQL_THDVAR_BOOL(
trace_sst_api, PLUGIN_VAR_RQCMDARG, "Generate trace output in the log for each call to the SstFileWriter",
nullptr, nullptr, FALSE);
static MYSQL_THDVAR_BOOL(
bulk_load, PLUGIN_VAR_RQCMDARG, "Use bulk-load mode for inserts. This disables " "unique_checks and enables rocksdb_commit_in_the_middle",
rocksdb_check_bulk_load, nullptr, FALSE);
static MYSQL_THDVAR_BOOL(bulk_load_allow_sk, PLUGIN_VAR_RQCMDARG, "Allow bulk loading of sk keys during bulk-load. " "Can be changed only when bulk load is disabled", /* Intentionally reuse unsorted's check function */
rocksdb_check_bulk_load_allow_unsorted, nullptr, FALSE);
static MYSQL_THDVAR_BOOL(bulk_load_allow_unsorted, PLUGIN_VAR_RQCMDARG, "Allow unsorted input during bulk-load. " "Can be changed only when bulk load is disabled",
rocksdb_check_bulk_load_allow_unsorted, nullptr, FALSE);
static MYSQL_SYSVAR_BOOL(enable_bulk_load_api, rocksdb_enable_bulk_load_api,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Enables using SstFileWriter for bulk loading",
nullptr, nullptr, rocksdb_enable_bulk_load_api);
static MYSQL_SYSVAR_STR(git_hash, rocksdb_git_hash,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Git revision of the RocksDB library used by MyRocks",
nullptr, nullptr, ROCKSDB_GIT_HASH);
static MYSQL_THDVAR_STR(tmpdir, PLUGIN_VAR_OPCMDARG | PLUGIN_VAR_MEMALLOC, "Directory for temporary files during DDL operations",
nullptr, nullptr, "");
#define DEFAULT_SKIP_UNIQUE_CHECK_TABLES ".*" static MYSQL_THDVAR_STR(
skip_unique_check_tables, PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_MEMALLOC, "Skip unique constraint checking for the specified tables", nullptr,
nullptr, DEFAULT_SKIP_UNIQUE_CHECK_TABLES);
static MYSQL_THDVAR_BOOL(
commit_in_the_middle, PLUGIN_VAR_RQCMDARG, "Commit rows implicitly every rocksdb_bulk_load_size, on bulk load/insert, " "update and delete",
nullptr, nullptr, FALSE);
static MYSQL_THDVAR_BOOL(
blind_delete_primary_key, PLUGIN_VAR_RQCMDARG, "Deleting rows by primary key lookup, without reading rows (Blind Deletes)." " Blind delete is disabled if the table has secondary key",
nullptr, nullptr, FALSE);
#if0// MARIAROCKS_NOT_YET : read-free replication is not supported
// This is bound to succeed since we've already checked for bad patterns in // rocksdb_validate_read_free_rpl_tables
rdb_read_free_regex_handler.set_patterns(wlist);
// update all table defs struct Rdb_read_free_rpl_updater : public Rdb_tables_scanner { int add_table(Rdb_tbl_def *tdef) override {
tdef->check_and_set_read_free_rpl_table(); return HA_EXIT_SUCCESS;
}
} updater;
ddl_manager.scan_for_tables(&updater);
if (wlist == DEFAULT_READ_FREE_RPL_TABLES) { // If running SET var = DEFAULT, then rocksdb_validate_read_free_rpl_tables // isn't called, and memory is never allocated for the value. Allocate it // here.
*static_cast<constchar **>(var_ptr) = my_strdup(wlist, MYF(MY_WME));
} else { // Otherwise, we just reuse the value allocated from // rocksdb_validate_read_free_rpl_tables.
*static_cast<constchar **>(var_ptr) = wlist;
}
}
static MYSQL_SYSVAR_STR(
read_free_rpl_tables, rocksdb_read_free_rpl_tables,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_MEMALLOC /*| PLUGIN_VAR_ALLOCATED*/, "List of tables that will use read-free replication on the slave " "(i.e. not lookup a row during replication)",
rocksdb_validate_read_free_rpl_tables, rocksdb_update_read_free_rpl_tables,
DEFAULT_READ_FREE_RPL_TABLES);
static MYSQL_SYSVAR_ENUM(
read_free_rpl, rocksdb_read_free_rpl,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_MEMALLOC, "Use read-free replication on the slave (i.e. no row lookup during " "replication). Default is OFF, PK_SK will enable it on all tables with " "primary key. PK_ONLY will enable it on tables where the only key is the " "primary key (i.e. no secondary keys)",
nullptr, nullptr, read_free_rpl_type::OFF, &read_free_rpl_typelib); #endif
static MYSQL_THDVAR_BOOL(skip_bloom_filter_on_read, PLUGIN_VAR_RQCMDARG, "Skip using bloom filter for reads", nullptr, nullptr, FALSE);
static MYSQL_THDVAR_ULONG(max_row_locks, PLUGIN_VAR_RQCMDARG, "Maximum number of locks a transaction can have",
nullptr, nullptr, /*default*/ RDB_DEFAULT_MAX_ROW_LOCKS, /*min*/ 1, /*max*/ RDB_MAX_ROW_LOCKS, 0);
static MYSQL_THDVAR_ULONGLONG(
write_batch_max_bytes, PLUGIN_VAR_RQCMDARG, "Maximum size of write batch in bytes. 0 means no limit", nullptr, nullptr, /* default */ 0, /* min */ 0, /* max */ SIZE_T_MAX, 1);
static MYSQL_THDVAR_BOOL(
lock_scanned_rows, PLUGIN_VAR_RQCMDARG, "Take and hold locks on rows that are scanned but not updated", nullptr,
nullptr, FALSE);
static MYSQL_THDVAR_ULONG(bulk_load_size, PLUGIN_VAR_RQCMDARG, "Max #records in a batch for bulk-load mode", nullptr,
nullptr, /*default*/ RDB_DEFAULT_BULK_LOAD_SIZE, /*min*/ 1, /*max*/ RDB_MAX_BULK_LOAD_SIZE, 0);
static MYSQL_THDVAR_ULONGLONG(
merge_buf_size, PLUGIN_VAR_RQCMDARG, "Size to allocate for merge sort buffers written out to disk " "during inplace index creation",
nullptr, nullptr, /* default (64MB) */ RDB_DEFAULT_MERGE_BUF_SIZE, /* min (100B) */ RDB_MIN_MERGE_BUF_SIZE, /* max */ SIZE_T_MAX, 1);
static MYSQL_THDVAR_ULONGLONG(
merge_combine_read_size, PLUGIN_VAR_RQCMDARG, "Size that we have to work with during combine (reading from disk) phase " "of " "external sort during fast index creation",
nullptr, nullptr, /* default (1GB) */ RDB_DEFAULT_MERGE_COMBINE_READ_SIZE, /* min (100B) */ RDB_MIN_MERGE_COMBINE_READ_SIZE, /* max */ SIZE_T_MAX, 1);
static MYSQL_THDVAR_ULONGLONG(
merge_tmp_file_removal_delay_ms, PLUGIN_VAR_RQCMDARG, "Fast index creation creates a large tmp file on disk during index " "creation. Removing this large file all at once when index creation is " "complete can cause trim stalls on Flash. This variable specifies a " "duration to sleep (in milliseconds) between calling chsize() to truncate " "the file in chunks. The chunk size is the same as merge_buf_size",
nullptr, nullptr, /* default (0ms) */ RDB_DEFAULT_MERGE_TMP_FILE_REMOVAL_DELAY, /* min (0ms) */ RDB_MIN_MERGE_TMP_FILE_REMOVAL_DELAY, /* max */ SIZE_T_MAX, 1);
static MYSQL_THDVAR_INT(
manual_compaction_threads, PLUGIN_VAR_RQCMDARG, "How many rocksdb threads to run for manual compactions", nullptr, nullptr, /* default rocksdb.dboption max_subcompactions */ 0, /* min */ 0, /* max */ 128, 0);
static MYSQL_SYSVAR_UINT(max_latest_deadlocks, rocksdb_max_latest_deadlocks,
PLUGIN_VAR_RQCMDARG, "Maximum number of recent " "deadlocks to store",
nullptr, rocksdb_set_max_latest_deadlocks,
rocksdb::kInitialMaxDeadlocks, 0, UINT32_MAX, 0);
static MYSQL_SYSVAR_ENUM(
info_log_level, rocksdb_info_log_level, PLUGIN_VAR_RQCMDARG, "Filter level for info logs to be written mysqld error log. " "Valid values include 'debug_level', 'info_level', 'warn_level'" "'error_level' and 'fatal_level'",
nullptr, rocksdb_set_rocksdb_info_log_level,
rocksdb::InfoLogLevel::ERROR_LEVEL, &info_log_level_typelib);
static MYSQL_THDVAR_INT(
perf_context_level, PLUGIN_VAR_RQCMDARG, "Perf Context Level for rocksdb internal timer stat collection", nullptr,
nullptr, /* default */ rocksdb::PerfLevel::kUninitialized, /* min */ rocksdb::PerfLevel::kUninitialized, /* max */ rocksdb::PerfLevel::kOutOfBounds - 1, 0);
static MYSQL_SYSVAR_UINT(
wal_recovery_mode, rocksdb_wal_recovery_mode, PLUGIN_VAR_RQCMDARG, "DBOptions::wal_recovery_mode for RocksDB. Default is kAbsoluteConsistency",
nullptr, nullptr, /* default */ (uint)rocksdb::WALRecoveryMode::kAbsoluteConsistency, /* min */ (uint)rocksdb::WALRecoveryMode::kTolerateCorruptedTailRecords, /* max */ (uint)rocksdb::WALRecoveryMode::kSkipAnyCorruptedRecords, 0);
static MYSQL_SYSVAR_UINT(
stats_level, rocksdb_stats_level, PLUGIN_VAR_RQCMDARG, "Statistics Level for RocksDB. Default is 0 (kExceptHistogramOrTimers)",
nullptr, rocksdb_set_rocksdb_stats_level, /* default */ (uint)rocksdb::StatsLevel::kExceptHistogramOrTimers, /* min */ (uint)rocksdb::StatsLevel::kDisableAll, /* max */ (uint)rocksdb::StatsLevel::kAll, 0);
static MYSQL_SYSVAR_SIZE_T(compaction_readahead_size,
rocksdb_db_options->compaction_readahead_size,
PLUGIN_VAR_RQCMDARG, "DBOptions::compaction_readahead_size for RocksDB",
nullptr, nullptr,
rocksdb_db_options->compaction_readahead_size, /* min */ 0L, /* max */ SIZE_T_MAX, 0);
static MYSQL_SYSVAR_STR(
persistent_cache_path, rocksdb_persistent_cache_path,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Path for BlockBasedTableOptions::persistent_cache for RocksDB", nullptr,
nullptr, "");
static MYSQL_SYSVAR_ULONG(
persistent_cache_size_mb, rocksdb_persistent_cache_size_mb,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Size of cache in MB for BlockBasedTableOptions::persistent_cache " "for RocksDB",
nullptr, nullptr, rocksdb_persistent_cache_size_mb, /* min */ 0L, /* max */ ULONG_MAX, 0);
static MYSQL_SYSVAR_UINT64_T(
delete_obsolete_files_period_micros,
rocksdb_db_options->delete_obsolete_files_period_micros,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::delete_obsolete_files_period_micros for RocksDB", nullptr,
nullptr, rocksdb_db_options->delete_obsolete_files_period_micros, /* min */ 0, /* max */ LONGLONG_MAX, 0);
static MYSQL_SYSVAR_INT(max_background_jobs,
rocksdb_db_options->max_background_jobs,
PLUGIN_VAR_RQCMDARG, "DBOptions::max_background_jobs for RocksDB", nullptr,
rocksdb_set_max_background_jobs,
rocksdb_db_options->max_background_jobs, /* min */ -1, /* max */ MAX_BACKGROUND_JOBS, 0);
static MYSQL_SYSVAR_UINT(max_subcompactions,
rocksdb_db_options->max_subcompactions,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::max_subcompactions for RocksDB", nullptr,
nullptr, rocksdb_db_options->max_subcompactions, /* min */ 1, /* max */ MAX_SUBCOMPACTIONS, 0);
static MYSQL_SYSVAR_SIZE_T(max_log_file_size,
rocksdb_db_options->max_log_file_size,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::max_log_file_size for RocksDB", nullptr,
nullptr, rocksdb_db_options->max_log_file_size, /* min */ 0L, /* max */ SIZE_T_MAX, 0);
static MYSQL_SYSVAR_SIZE_T(log_file_time_to_roll,
rocksdb_db_options->log_file_time_to_roll,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::log_file_time_to_roll for RocksDB",
nullptr, nullptr,
rocksdb_db_options->log_file_time_to_roll, /* min */ 0L, /* max */ SIZE_T_MAX, 0);
static MYSQL_SYSVAR_SIZE_T(keep_log_file_num,
rocksdb_db_options->keep_log_file_num,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::keep_log_file_num for RocksDB", nullptr,
nullptr, rocksdb_db_options->keep_log_file_num, /* min */ 0L, /* max */ SIZE_T_MAX, 0);
static MYSQL_SYSVAR_UINT64_T(max_manifest_file_size,
rocksdb_db_options->max_manifest_file_size,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::max_manifest_file_size for RocksDB",
nullptr, nullptr,
rocksdb_db_options->max_manifest_file_size, /* min */ 0L, /* max */ ULONGLONG_MAX, 0);
static MYSQL_SYSVAR_INT(table_cache_numshardbits,
rocksdb_db_options->table_cache_numshardbits,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::table_cache_numshardbits for RocksDB",
nullptr, nullptr,
rocksdb_db_options->table_cache_numshardbits, // LRUCache limits this to 19 bits, anything greater // fails to create a cache and returns a nullptr /* min */ 0, /* max */ 19, 0);
static MYSQL_SYSVAR_UINT64_T(wal_ttl_seconds, rocksdb_db_options->WAL_ttl_seconds,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::WAL_ttl_seconds for RocksDB", nullptr,
nullptr, rocksdb_db_options->WAL_ttl_seconds, /* min */ 0L, /* max */ LONGLONG_MAX, 0);
static MYSQL_SYSVAR_UINT64_T(wal_size_limit_mb,
rocksdb_db_options->WAL_size_limit_MB,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::WAL_size_limit_MB for RocksDB", nullptr,
nullptr, rocksdb_db_options->WAL_size_limit_MB, /* min */ 0L, /* max */ LONGLONG_MAX, 0);
static MYSQL_SYSVAR_SIZE_T(manifest_preallocation_size,
rocksdb_db_options->manifest_preallocation_size,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "DBOptions::manifest_preallocation_size for RocksDB",
nullptr, nullptr,
rocksdb_db_options->manifest_preallocation_size, /* min */ 0L, /* max */ SIZE_T_MAX, 0);
// When pin_l0_filter_and_index_blocks_in_cache is true, RocksDB will use the // LRU cache, but will always keep the filter & index block's handle checked // out (=won't call ShardedLRUCache::Release), plus the parsed out objects // the LRU cache will never push flush them out, hence they're pinned. // // This fixes the mutex contention between :ShardedLRUCache::Lookup and // ShardedLRUCache::Release which reduced the QPS ratio (QPS using secondary // index / QPS using PK). static MYSQL_SYSVAR_BOOL(
pin_l0_filter_and_index_blocks_in_cache,
*reinterpret_cast<my_bool *>(
&rocksdb_tbl_options->pin_l0_filter_and_index_blocks_in_cache),
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "pin_l0_filter_and_index_blocks_in_cache for RocksDB", nullptr, nullptr, true);
static MYSQL_SYSVAR_STR(override_cf_options, rocksdb_override_cf_options,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "option overrides per cf for RocksDB", nullptr, nullptr, "");
static MYSQL_SYSVAR_STR(update_cf_options, rocksdb_update_cf_options,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_MEMALLOC /* psergey-merge: need this? : PLUGIN_VAR_ALLOCATED*/, "Option updates per column family for RocksDB",
rocksdb_validate_update_cf_options,
rocksdb_set_update_cf_options, nullptr);
static MYSQL_SYSVAR_UINT(flush_log_at_trx_commit,
rocksdb_flush_log_at_trx_commit, PLUGIN_VAR_RQCMDARG, "Sync on transaction commit. Similar to " "innodb_flush_log_at_trx_commit. 1: sync on commit, " "0,2: not sync on commit",
rocksdb_validate_flush_log_at_trx_commit, nullptr, /* default */ FLUSH_LOG_SYNC, /* min */ FLUSH_LOG_NEVER, /* max */ FLUSH_LOG_BACKGROUND, 0);
static MYSQL_THDVAR_BOOL(write_disable_wal, PLUGIN_VAR_RQCMDARG, "WriteOptions::disableWAL for RocksDB", nullptr,
nullptr, rocksdb::WriteOptions().disableWAL);
static MYSQL_THDVAR_BOOL(
write_ignore_missing_column_families, PLUGIN_VAR_RQCMDARG, "WriteOptions::ignore_missing_column_families for RocksDB", nullptr,
nullptr, rocksdb::WriteOptions().ignore_missing_column_families);
static MYSQL_THDVAR_BOOL(
unsafe_for_binlog, PLUGIN_VAR_RQCMDARG, "Allowing statement based binary logging which may break consistency",
nullptr, nullptr, FALSE);
static MYSQL_THDVAR_UINT(records_in_range, PLUGIN_VAR_RQCMDARG, "Used to override the result of records_in_range(). " "Set to a positive number to override",
nullptr, nullptr, 0, /* min */ 0, /* max */ INT_MAX, 0);
static MYSQL_THDVAR_UINT(force_index_records_in_range, PLUGIN_VAR_RQCMDARG, "Used to override the result of records_in_range() " "when FORCE INDEX is used",
nullptr, nullptr, 0, /* min */ 0, /* max */ INT_MAX, 0);
static MYSQL_SYSVAR_UINT(
debug_optimizer_n_rows, rocksdb_debug_optimizer_n_rows,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY | PLUGIN_VAR_NOSYSVAR, "Test only to override rocksdb estimates of table size in a memtable",
nullptr, nullptr, 0, /* min */ 0, /* max */ INT_MAX, 0);
static MYSQL_SYSVAR_UINT(force_compute_memtable_stats_cachetime,
rocksdb_force_compute_memtable_stats_cachetime,
PLUGIN_VAR_RQCMDARG, "Time in usecs to cache memtable estimates", nullptr,
nullptr, /* default */ 60 * 1000 * 1000, /* min */ 0, /* max */ INT_MAX, 0);
static MYSQL_SYSVAR_BOOL(
debug_optimizer_no_zero_cardinality,
rocksdb_debug_optimizer_no_zero_cardinality, PLUGIN_VAR_RQCMDARG, "In case if cardinality is zero, overrides it with some value", nullptr,
nullptr, TRUE);
static MYSQL_SYSVAR_BOOL(signal_drop_index_thread,
rocksdb_signal_drop_index_thread, PLUGIN_VAR_RQCMDARG, "Wake up drop index thread", nullptr,
rocksdb_drop_index_wakeup_thread, FALSE);
static MYSQL_SYSVAR_BOOL(
enable_ttl, rocksdb_enable_ttl, PLUGIN_VAR_RQCMDARG, "Enable expired TTL records to be dropped during compaction", nullptr,
nullptr, TRUE);
static MYSQL_SYSVAR_BOOL(
enable_ttl_read_filtering, rocksdb_enable_ttl_read_filtering,
PLUGIN_VAR_RQCMDARG, "For tables with TTL, expired records are skipped/filtered out during " "processing and in query results. Disabling this will allow these records " "to be seen, but as a result rows may disappear in the middle of " "transactions as they are dropped during compaction. Use with caution",
nullptr, nullptr, TRUE);
static MYSQL_SYSVAR_INT(
debug_ttl_rec_ts, rocksdb_debug_ttl_rec_ts, PLUGIN_VAR_RQCMDARG, "For debugging purposes only. Overrides the TTL of records to " "now() + debug_ttl_rec_ts. The value can be +/- to simulate " "a record inserted in the past vs a record inserted in the 'future'. " "A value of 0 denotes that the variable is not set. This variable is a " "no-op in non-debug builds",
nullptr, nullptr, 0, /* min */ -3600, /* max */ 3600, 0);
static MYSQL_SYSVAR_INT(
debug_ttl_snapshot_ts, rocksdb_debug_ttl_snapshot_ts, PLUGIN_VAR_RQCMDARG, "For debugging purposes only. Sets the snapshot during compaction to " "now() + debug_set_ttl_snapshot_ts. The value can be +/- to simulate " "a snapshot in the past vs a snapshot created in the 'future'. " "A value of 0 denotes that the variable is not set. This variable is a " "no-op in non-debug builds",
nullptr, nullptr, 0, /* min */ -3600, /* max */ 3600, 0);
static MYSQL_SYSVAR_INT(
debug_ttl_read_filter_ts, rocksdb_debug_ttl_read_filter_ts,
PLUGIN_VAR_RQCMDARG, "For debugging purposes only. Overrides the TTL read filtering time to " "time + debug_ttl_read_filter_ts. A value of 0 denotes that the variable " "is not set. This variable is a no-op in non-debug builds",
nullptr, nullptr, 0, /* min */ -3600, /* max */ 3600, 0);
static MYSQL_SYSVAR_BOOL(
debug_ttl_ignore_pk, rocksdb_debug_ttl_ignore_pk, PLUGIN_VAR_RQCMDARG, "For debugging purposes only. If true, compaction filtering will not occur " "on PK TTL data. This variable is a no-op in non-debug builds",
nullptr, nullptr, FALSE);
static MYSQL_SYSVAR_UINT(
max_manual_compactions, rocksdb_max_manual_compactions, PLUGIN_VAR_RQCMDARG, "Maximum number of pending + ongoing number of manual compactions",
nullptr, nullptr, /* default */ 10, /* min */ 0, /* max */ UINT_MAX, 0);
static MYSQL_SYSVAR_BOOL(
rollback_on_timeout, rocksdb_rollback_on_timeout, PLUGIN_VAR_OPCMDARG, "Whether to roll back the complete transaction or a single statement on " "lock wait timeout (a single statement by default)",
NULL, NULL, FALSE);
static MYSQL_SYSVAR_UINT(
debug_manual_compaction_delay, rocksdb_debug_manual_compaction_delay,
PLUGIN_VAR_RQCMDARG, "For debugging purposes only. Sleeping specified seconds " "for simulating long running compactions",
nullptr, nullptr, 0, /* min */ 0, /* max */ UINT_MAX, 0);
static MYSQL_SYSVAR_BOOL(
reset_stats, rocksdb_reset_stats, PLUGIN_VAR_RQCMDARG, "Reset the RocksDB internal statistics without restarting the DB", nullptr,
rocksdb_set_reset_stats, FALSE);
static MYSQL_SYSVAR_UINT(io_write_timeout, rocksdb_io_write_timeout_secs,
PLUGIN_VAR_RQCMDARG, "Timeout for experimental I/O watchdog", nullptr,
rocksdb_set_io_write_timeout, /* default */ 0, /* min */ 0L, /* max */ UINT_MAX, 0);
static MYSQL_SYSVAR_BOOL(enable_2pc, rocksdb_enable_2pc, PLUGIN_VAR_RQCMDARG, "Enable two phase commit for MyRocks", nullptr,
nullptr, TRUE);
static MYSQL_SYSVAR_BOOL(strict_collation_check, rocksdb_strict_collation_check,
PLUGIN_VAR_RQCMDARG, "Enforce case sensitive collation for MyRocks indexes",
nullptr, nullptr, TRUE);
static MYSQL_SYSVAR_STR(strict_collation_exceptions,
rocksdb_strict_collation_exceptions,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_MEMALLOC, "List of tables (using regex) that are excluded " "from the case sensitive collation enforcement",
nullptr, rocksdb_set_collation_exception_list, "");
static MYSQL_SYSVAR_BOOL(collect_sst_properties, rocksdb_collect_sst_properties,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Enables collecting SST file properties on each flush",
nullptr, nullptr, rocksdb_collect_sst_properties);
static MYSQL_SYSVAR_BOOL(
force_flush_memtable_now, rocksdb_force_flush_memtable_now_var,
PLUGIN_VAR_RQCMDARG, "Forces memstore flush which may block all write requests so be careful",
rocksdb_force_flush_memtable_now, rocksdb_force_flush_memtable_now_stub, FALSE);
static MYSQL_SYSVAR_BOOL(
force_flush_memtable_and_lzero_now,
rocksdb_force_flush_memtable_and_lzero_now_var, PLUGIN_VAR_RQCMDARG, "Acts similar to force_flush_memtable_now, but also compacts all L0 files",
rocksdb_force_flush_memtable_and_lzero_now,
rocksdb_force_flush_memtable_and_lzero_now_stub, FALSE);
static MYSQL_SYSVAR_UINT(
seconds_between_stat_computes, rocksdb_seconds_between_stat_computes,
PLUGIN_VAR_RQCMDARG, "Sets a number of seconds to wait between optimizer stats recomputation. " "Only changed indexes will be refreshed",
nullptr, nullptr, rocksdb_seconds_between_stat_computes, /* min */ 0L, /* max */ UINT_MAX, 0);
static MYSQL_SYSVAR_LONGLONG(compaction_sequential_deletes,
rocksdb_compaction_sequential_deletes,
PLUGIN_VAR_RQCMDARG, "RocksDB will trigger compaction for the file if " "it has more than this number sequential deletes " "per window",
nullptr, rocksdb_set_compaction_options,
DEFAULT_COMPACTION_SEQUENTIAL_DELETES, /* min */ 0L, /* max */ MAX_COMPACTION_SEQUENTIAL_DELETES, 0);
static MYSQL_SYSVAR_LONGLONG(
compaction_sequential_deletes_window,
rocksdb_compaction_sequential_deletes_window, PLUGIN_VAR_RQCMDARG, "Size of the window for counting rocksdb_compaction_sequential_deletes",
nullptr, rocksdb_set_compaction_options,
DEFAULT_COMPACTION_SEQUENTIAL_DELETES_WINDOW, /* min */ 0L, /* max */ MAX_COMPACTION_SEQUENTIAL_DELETES_WINDOW, 0);
static MYSQL_SYSVAR_LONGLONG(
compaction_sequential_deletes_file_size,
rocksdb_compaction_sequential_deletes_file_size, PLUGIN_VAR_RQCMDARG, "Minimum file size required for compaction_sequential_deletes", nullptr,
rocksdb_set_compaction_options, 0L, /* min */ -1L, /* max */ LLONG_MAX, 0);
static MYSQL_SYSVAR_BOOL(
print_snapshot_conflict_queries, rocksdb_print_snapshot_conflict_queries,
PLUGIN_VAR_RQCMDARG, "Logging queries that got snapshot conflict errors into *.err log", nullptr,
nullptr, rocksdb_print_snapshot_conflict_queries);
static MYSQL_THDVAR_INT(checksums_pct, PLUGIN_VAR_RQCMDARG, "How many percentages of rows to be checksummed",
nullptr, nullptr, RDB_MAX_CHECKSUMS_PCT, /* min */ 0, /* max */ RDB_MAX_CHECKSUMS_PCT, 0);
static MYSQL_THDVAR_BOOL(store_row_debug_checksums, PLUGIN_VAR_RQCMDARG, "Include checksums when writing index/table records",
nullptr, nullptr, false/* default value */);
static MYSQL_THDVAR_BOOL(verify_row_debug_checksums, PLUGIN_VAR_RQCMDARG, "Verify checksums when reading index/table records",
nullptr, nullptr, false/* default value */);
static MYSQL_THDVAR_BOOL(master_skip_tx_api, PLUGIN_VAR_RQCMDARG, "Skipping holding any lock on row access. " "Not effective on slave",
nullptr, nullptr, false);
static MYSQL_SYSVAR_UINT(
validate_tables, rocksdb_validate_tables,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Verify all .frm files match all RocksDB tables (0 means no verification, " "1 means verify and fail on error, and 2 means verify but continue",
nullptr, nullptr, 1/* default value */, 0 /* min value */, 2/* max value */, 0);
static MYSQL_SYSVAR_UINT(
ignore_datadic_errors, rocksdb_ignore_datadic_errors,
PLUGIN_VAR_RQCMDARG | PLUGIN_VAR_READONLY, "Ignore MyRocks' data directory errors. " "(CAUTION: Use only to start the server and perform repairs. Do NOT use " "for regular operation)",
nullptr, nullptr, 0/* default value */, 0 /* min value */, 1/* max value */, 0);
static MYSQL_SYSVAR_UINT(
table_stats_sampling_pct, rocksdb_table_stats_sampling_pct,
PLUGIN_VAR_RQCMDARG, "Percentage of entries to sample when collecting statistics about table " "properties. Specify either 0 to sample everything or percentage ["
STRINGIFY_ARG(RDB_TBL_STATS_SAMPLE_PCT_MIN) ".."
STRINGIFY_ARG(RDB_TBL_STATS_SAMPLE_PCT_MAX) "]. By default " STRINGIFY_ARG(RDB_DEFAULT_TBL_STATS_SAMPLE_PCT) "% of entries are sampled",
nullptr, rocksdb_set_table_stats_sampling_pct, /* default */
RDB_DEFAULT_TBL_STATS_SAMPLE_PCT, /* everything */ 0, /* max */ RDB_TBL_STATS_SAMPLE_PCT_MAX, 0);
static MYSQL_SYSVAR_UINT(
stats_recalc_rate, rocksdb_stats_recalc_rate, PLUGIN_VAR_RQCMDARG, "The number of indexes per second to recalculate statistics for. 0 to " "disable background recalculation",
nullptr, nullptr, 0/* default value */, 0 /* min value */,
UINT_MAX /* max value */, 0);
static MYSQL_SYSVAR_BOOL(
large_prefix, rocksdb_large_prefix, PLUGIN_VAR_RQCMDARG, "Support large index prefix length of 3072 bytes. If off, the maximum " "index prefix length is 767",
nullptr, nullptr, FALSE);
static MYSQL_SYSVAR_BOOL(
allow_to_start_after_corruption, rocksdb_allow_to_start_after_corruption,
PLUGIN_VAR_OPCMDARG | PLUGIN_VAR_READONLY, "Allow server still to start successfully even if RocksDB corruption is " "detected",
nullptr, nullptr, FALSE);
static MYSQL_SYSVAR_BOOL(error_on_suboptimal_collation,
rocksdb_error_on_suboptimal_collation,
PLUGIN_VAR_OPCMDARG | PLUGIN_VAR_READONLY, "Raise an error instead of warning if a sub-optimal " "collation is used",
nullptr, nullptr, TRUE);
static MYSQL_SYSVAR_BOOL(
enable_insert_with_update_caching,
rocksdb_enable_insert_with_update_caching, PLUGIN_VAR_OPCMDARG, "Whether to enable optimization where we cache the read from a failed " "insertion attempt in INSERT ON DUPLICATE KEY UPDATE",
nullptr, nullptr, TRUE);
staticint rocksdb_compact_column_family(THD *const thd, struct st_mysql_sys_var *const var, void *const var_ptr, struct st_mysql_value *const value) { char buff[STRING_BUFFER_USUAL_SIZE]; int len = sizeof(buff);
DBUG_ASSERT(value != nullptr);
if (constchar *const cf = value->val_str(value, buff, &len)) { auto cfh = cf_manager.get_cf(cf); if (cfh != nullptr && rdb != nullptr) { int mc_id = rdb_mc_thread.request_manual_compaction(
cfh, nullptr, nullptr, THDVAR(thd, manual_compaction_threads)); if (mc_id == -1) {
my_error(ER_INTERNAL_ERROR, MYF(0), "Can't schedule more manual compactions. " "Increase rocksdb_max_manual_compactions or stop issuing " "more manual compactions."); return HA_EXIT_FAILURE;
} elseif (mc_id < 0) { return HA_EXIT_FAILURE;
} // NO_LINT_DEBUG
sql_print_information("RocksDB: Manual compaction of column family: %s\n",
cf); // Checking thd state every short cycle (100ms). This is for allowing to // exiting this function without waiting for CompactRange to finish. do {
my_sleep(100000);
} while (!thd->killed &&
!rdb_mc_thread.is_manual_compaction_finished(mc_id));
if (thd->killed) { // This cancels if requested compaction state is INITED. // TODO(yoshinorim): Cancel running compaction as well once // it is supported in RocksDB.
rdb_mc_thread.clear_manual_compaction_request(mc_id, true);
}
}
} return HA_EXIT_SUCCESS;
}
// If the owning Rdb_transaction gets destructed we need to not reference // it anymore. void detach() { m_owning_tx = nullptr; }
};
#ifdef MARIAROCKS_NOT_YET // ER_LOCK_WAIT_TIMEOUT error also has a reason in facebook/mysql-5.6 #endif
String timeout_message(constchar *command, constchar *name1, constchar *name2)
{
String msg;
msg.append(STRING_WITH_LEN("Timeout on "));
msg.append(command, strlen(command));
msg.append(STRING_WITH_LEN(": "));
msg.append(name1, strlen(name1)); if (name2 && name2[0])
{
msg.append('.');
msg.append(name2, strlen(name2));
} return msg;
}
/* This is the base class for transactions when interacting with rocksdb.
*/ class Rdb_transaction { protected:
ulonglong m_write_count = 0;
ulonglong m_insert_count = 0;
ulonglong m_update_count = 0;
ulonglong m_delete_count = 0;
ulonglong m_lock_count = 0;
std::unordered_map<GL_INDEX_ID, ulonglong> m_auto_incr_map;
// This should be used only when updating binlog information. virtual rocksdb::WriteBatchBase *get_write_batch() = 0; virtualbool commit_no_binlog() = 0; virtual rocksdb::Iterator *get_iterator( const rocksdb::ReadOptions &options,
rocksdb::ColumnFamilyHandle *column_family) = 0;
// Iterate through the merge map merging all keys into data dictionary.
rocksdb::Status s; for (auto &it : m_auto_incr_map) {
s = dict_manager.put_auto_incr_val(wb, it.first, it.second); if (!s.ok()) { return s;
}
}
m_auto_incr_map.clear(); return s;
}
private: // The Rdb_sst_info structures we are currently loading. In a partitioned // table this can have more than one entry
std::vector<std::shared_ptr<Rdb_sst_info>> m_curr_bulk_load;
std::string m_curr_bulk_load_tablename;
/* External merge sorts for bulk load: key ID -> merge sort instance */
std::unordered_map<GL_INDEX_ID, Rdb_index_merge> m_key_merge;
public: int get_key_merge(GL_INDEX_ID kd_gl_id, rocksdb::ColumnFamilyHandle *cf,
Rdb_index_merge **key_merge) { int res; auto it = m_key_merge.find(kd_gl_id); if (it == m_key_merge.end()) {
m_key_merge.emplace(
std::piecewise_construct, std::make_tuple(kd_gl_id),
std::make_tuple(
get_rocksdb_tmpdir(), THDVAR(get_thd(), merge_buf_size),
THDVAR(get_thd(), merge_combine_read_size),
THDVAR(get_thd(), merge_tmp_file_removal_delay_ms), cf));
it = m_key_merge.find(kd_gl_id); if ((res = it->second.init()) != 0) { return res;
}
}
*key_merge = &it->second; return HA_EXIT_SUCCESS;
}
/* Finish bulk loading for all table handlers belongs to one connection */ int finish_bulk_load(bool *is_critical_error = nullptr, int print_client_error = true) {
Ensure_cleanup cleanup([&]() { // Always clear everything regardless of success/failure
m_curr_bulk_load.clear();
m_curr_bulk_load_tablename.clear();
m_key_merge.clear();
});
int rc = 0; if (is_critical_error) {
*is_critical_error = true;
}
// PREPARE phase: finish all on-going bulk loading Rdb_sst_info and // collect all Rdb_sst_commit_info containing (SST files, cf) int rc2 = 0;
std::vector<Rdb_sst_info::Rdb_sst_commit_info> sst_commit_list;
sst_commit_list.reserve(m_curr_bulk_load.size());
for (auto &sst_info : m_curr_bulk_load) {
Rdb_sst_info::Rdb_sst_commit_info commit_info;
// Commit the list of SST files and move it to the end of // sst_commit_list, effectively transfer the ownership over
rc2 = sst_info->finish(&commit_info, print_client_error); if (rc2 && rc == 0) { // Don't return yet - make sure we finish all the SST infos
rc = rc2;
}
// Make sure we have work to do - we might be losing the race if (rc2 == 0 && commit_info.has_work()) {
sst_commit_list.emplace_back(std::move(commit_info));
DBUG_ASSERT(!commit_info.has_work());
}
}
if (rc) { return rc;
}
// MERGING Phase: Flush the index_merge sort buffers into SST files in // Rdb_sst_info and collect all Rdb_sst_commit_info containing // (SST files, cf) if (!m_key_merge.empty()) {
Ensure_cleanup malloc_cleanup([]() { /* Explicitlytelljemalloctocleanupanyunuseddirtypagesatthis point. Seehttps://reviews.facebook.net/D63723 for more details.
*/
purge_all_jemalloc_arenas();
});
rocksdb::Slice merge_key;
rocksdb::Slice merge_val; for (auto it = m_key_merge.begin(); it != m_key_merge.end(); it++) {
GL_INDEX_ID index_id = it->first;
std::shared_ptr<const Rdb_key_def> keydef =
ddl_manager.safe_find(index_id);
std::string table_name = ddl_manager.safe_get_table_name(index_id);
// Unable to find key definition or table name since the // table could have been dropped. // TODO(herman): there is a race here between dropping the table // and detecting a drop here. If the table is dropped while bulk // loading is finishing, these keys being added here may // be missed by the compaction filter and not be marked for // removal. It is unclear how to lock the sql table from the storage // engine to prevent modifications to it while bulk load is occurring. if (keydef == nullptr) { if (is_critical_error) { // We used to set the error but simply ignores it. This follows // current behavior and we should revisit this later
*is_critical_error = false;
} return HA_ERR_KEY_NOT_FOUND;
} elseif (table_name.empty()) { if (is_critical_error) { // We used to set the error but simply ignores it. This follows // current behavior and we should revisit this later
*is_critical_error = false;
} return HA_ERR_NO_SUCH_TABLE;
} const std::string &index_name = keydef->get_name();
Rdb_index_merge &rdb_merge = it->second;
// Rdb_sst_info expects a denormalized table name in the form of // "./database/table"
std::replace(table_name.begin(), table_name.end(), '.', '/');
table_name = "./" + table_name; auto sst_info = std::make_shared<Rdb_sst_info>(
rdb, table_name, index_name, rdb_merge.get_cf(),
*rocksdb_db_options, THDVAR(get_thd(), trace_sst_api));
while ((rc2 = rdb_merge.next(&merge_key, &merge_val)) == 0) { if ((rc2 = sst_info->put(merge_key, merge_val)) != 0) {
rc = rc2;
// Don't return yet - make sure we finish the sst_info break;
}
}
// -1 => no more items if (rc2 != -1 && rc != 0) {
rc = rc2;
}
Rdb_sst_info::Rdb_sst_commit_info commit_info;
rc2 = sst_info->finish(&commit_info, print_client_error); if (rc2 != 0 && rc == 0) { // Only set the error from sst_info->finish if finish failed and we // didn't fail before. In other words, we don't have finish's // success mask earlier failures
rc = rc2;
}
if (rc) { return rc;
}
if (commit_info.has_work()) {
sst_commit_list.emplace_back(std::move(commit_info));
DBUG_ASSERT(!commit_info.has_work());
}
}
}
// Early return in case we lost the race completely and end up with no // work at all if (sst_commit_list.size() == 0) { return rc;
}
// INGEST phase: Group all Rdb_sst_commit_info by cf (as they might // have the same cf across different indexes) and call out to RocksDB // to ingest all SST files in one atomic operation
rocksdb::IngestExternalFileOptions options;
options.move_files = true;
options.snapshot_consistency = false;
options.allow_global_seqno = false;
options.allow_blocking_flush = false;
const rocksdb::Status s = rdb->IngestExternalFiles(args); if (THDVAR(m_thd, trace_sst_api)) { // NO_LINT_DEBUG
sql_print_information( "SST Tracing: IngestExternalFile '%zu' files returned %s", file_count,
s.ok() ? "ok" : "not ok");
}
if (!s.ok()) { if (print_client_error) {
Rdb_sst_info::report_error_msg(s, nullptr);
} return HA_ERR_ROCKSDB_BULK_LOAD;
}
// COMMIT phase: mark everything as completed. This avoids SST file // deletion kicking in. Otherwise SST files would get deleted if this // entire operation is aborted for (auto &commit_info : sst_commit_list) {
commit_info.commit();
}
rocksdb::Iterator *get_iterator(
rocksdb::ColumnFamilyHandle *const column_family, bool skip_bloom_filter, bool fill_cache, const rocksdb::Slice &eq_cond_lower_bound, const rocksdb::Slice &eq_cond_upper_bound, bool read_current = false, bool create_snapshot = true) { // Make sure we are not doing both read_current (which implies we don't // want a snapshot) and create_snapshot which makes sure we create // a snapshot
DBUG_ASSERT(column_family != nullptr);
DBUG_ASSERT(!read_current || !create_snapshot);
if (create_snapshot) acquire_snapshot(true);
rocksdb::ReadOptions options = m_read_opts;
if (skip_bloom_filter) {
options.total_order_seek = true;
options.iterate_lower_bound = &eq_cond_lower_bound;
options.iterate_upper_bound = &eq_cond_upper_bound;
} else { // With this option, Iterator::Valid() returns false if key // is outside of the prefix bloom filter range set at Seek(). // Must not be set to true if not using bloom filter.
options.prefix_same_as_start = true;
}
options.fill_cache = fill_cache; if (read_current) {
options.snapshot = nullptr;
} return get_iterator(options, column_family);
}
/* Calledwhena"top-level"statementinsideatransactioncompletes successfullyanditschangesbecomepartofthetransaction'schanges.
*/ int make_stmt_savepoint_permanent() { // Take another RocksDB savepoint only if we had changes since the last // one. This is very important for long transactions doing lots of // SELECTs. if (m_writes_at_last_savepoint != m_write_count) {
rocksdb::WriteBatchBase *batch = get_write_batch();
rocksdb::Status status = rocksdb::Status::NotFound(); while ((status = batch->PopSavePoint()) == rocksdb::Status::OK()) {
}
if (status != rocksdb::Status::NotFound()) { return HA_EXIT_FAILURE;
}
private: void release_tx(void) { // We are done with the current active transaction object. Preserve it // for later reuse.
DBUG_ASSERT(m_rocksdb_reuse_tx == nullptr);
m_rocksdb_reuse_tx = m_rocksdb_tx;
m_rocksdb_tx = nullptr;
}
bool prepare(const rocksdb::TransactionName &name) override {
rocksdb::Status s;
s = m_rocksdb_tx->SetName(name); if (!s.ok()) {
rdb_handle_io_error(s, RDB_IO_ERROR_TX_COMMIT); returnfalse;
}
s = merge_auto_incr_map(m_rocksdb_tx->GetWriteBatch()->GetWriteBatch()); if (!s.ok()) {
rdb_handle_io_error(s, RDB_IO_ERROR_TX_COMMIT); returnfalse;
}
s = m_rocksdb_tx->Prepare(); if (!s.ok()) {
rdb_handle_io_error(s, RDB_IO_ERROR_TX_COMMIT); returnfalse;
} returntrue;
}
bool commit_no_binlog() override { bool res = false;
rocksdb::Status s;
s = merge_auto_incr_map(m_rocksdb_tx->GetWriteBatch()->GetWriteBatch()); if (!s.ok()) {
rdb_handle_io_error(s, RDB_IO_ERROR_TX_COMMIT);
res = true; goto error;
}
release_snapshot();
s = m_rocksdb_tx->Commit(); if (!s.ok()) {
rdb_handle_io_error(s, RDB_IO_ERROR_TX_COMMIT);
res = true; goto error;
}
on_commit();
error:
on_rollback(); /* Save the transaction object to be reused */
release_tx();
rocksdb::Status get(rocksdb::ColumnFamilyHandle *const column_family, const rocksdb::Slice &key,
rocksdb::PinnableSlice *const value) const override { // clean PinnableSlice right before Get() for multiple gets per statement // the resources after the last Get in a statement are cleared in // handler::reset call
value->Reset();
global_stats.queries[QUERIES_POINT].inc(); return m_rocksdb_tx->Get(m_read_opts, column_family, key, value);
}
if (value != nullptr) {
value->Reset();
}
rocksdb::Status s; // If snapshot is null, pass it to GetForUpdate and snapshot is // initialized there. Snapshot validation is skipped in that case. if (m_read_opts.snapshot == nullptr || do_validate) {
s = m_rocksdb_tx->GetForUpdate(
m_read_opts, column_family, key, value, exclusive,
m_read_opts.snapshot ? do_validate : false);
} else { // If snapshot is set, and if skipping validation, // call GetForUpdate without validation and set back old snapshot auto saved_snapshot = m_read_opts.snapshot;
m_read_opts.snapshot = nullptr;
s = m_rocksdb_tx->GetForUpdate(m_read_opts, column_family, key, value,
exclusive, false);
m_read_opts.snapshot = saved_snapshot;
} return s;
}
Forhookingtostartofstatementthatisitsowntransaction,see ha_rocksdb::external_lock().
*/ void start_stmt() override { // Set the snapshot to delayed acquisition (SetSnapshotOnNextOperation)
acquire_snapshot(false);
}
/* Thismustbecalledwhenlaststatementisrolledback,butthetransaction continues
*/ void rollback_stmt() override { /* TODO: here we must release the locks taken since the start_stmt() call */ if (m_rocksdb_tx) { const rocksdb::Snapshot *const org_snapshot = m_rocksdb_tx->GetSnapshot();
rollback_to_stmt_savepoint();
const rocksdb::Snapshot *const cur_snapshot = m_rocksdb_tx->GetSnapshot(); if (org_snapshot != cur_snapshot) { if (org_snapshot != nullptr) m_snapshot_timestamp = 0;
explicit Rdb_transaction_impl(THD *const thd)
: Rdb_transaction(thd), m_rocksdb_tx(nullptr) { // Create a notifier that can be called when a snapshot gets generated.
m_notifier = std::make_shared<Rdb_snapshot_notifier>(this);
}
~Rdb_transaction_impl() override {
rollback();
// Theoretically the notifier could outlive the Rdb_transaction_impl // (because of the shared_ptr), so let it know it can't reference // the transaction anymore.
m_notifier->detach();
// Free any transaction memory that is still hanging around. delete m_rocksdb_reuse_tx;
DBUG_ASSERT(m_rocksdb_tx == nullptr);
}
};
/* This is a rocksdb write batch. This class doesn't hold or wait on any transactionlocks(skipsrocksdbtransactionAPI)thusgivingbetter performance.
Currentlythisisonlyusedforreplicationthreadswhichareguaranteed tobenon-conflicting.Anyfurtherusageofthisclassshouldcompletely bethoughtthoroughly.
*/ class Rdb_writebatch_impl : public Rdb_transaction {
rocksdb::WriteBatchWithIndex *m_batch;
rocksdb::WriteOptions write_opts; // Called after commit/rollback. void reset() {
m_batch->Clear();
m_read_opts = rocksdb::ReadOptions();
m_ddl_transaction = false;
}
void release_lock(rocksdb::ColumnFamilyHandle *const column_family, const std::string &rowkey) override { // Nothing to do here since we don't hold any row locks.
}
/** Foraslave,prepare()updatestheslave_gtid_infotablewhichtracksthe replicationprogress.
*/ staticint rocksdb_prepare(THD* thd, bool prepare_tx)
{ bool async=false; // This is "ASYNC_COMMIT" feature which is only present in webscalesql
Rdb_transaction *tx = get_tx_from_thd(thd); if (!tx->can_prepare()) { return HA_EXIT_FAILURE;
} if (prepare_tx ||
(!my_core::thd_test_options(thd, OPTION_NOT_AUTOCOMMIT | OPTION_BEGIN))) { /* We were instructed to prepare the whole transaction, or
this is an SQL statement end and autocommit is on */
staticvoid rocksdb_checkpoint_request(void *cookie)
{ const rocksdb::Status s= rdb->FlushWAL(true); //TODO: what to do on error? if (s.ok())
{
rocksdb_wal_group_syncs++;
commit_checkpoint_notify_ha(cookie);
}
}
/* @paramall:TRUE-committhetransaction FALSE-SQLstatementended
*/ staticvoid rocksdb_commit_ordered(THD* thd, bool all)
{ // Same assert as InnoDB has
DBUG_ASSERT(all || (!thd_test_options(thd, OPTION_NOT_AUTOCOMMIT |
OPTION_BEGIN)));
Rdb_transaction *tx = get_tx_from_thd(thd); if (!tx->is_two_phase()) { /* ordered_commitissupposedlyslowerasitisdonesequentially inordertopreservecommitorder.
/* note: h->external_lock(F_UNLCK) is called after this function is called) */
Rdb_transaction *tx = get_tx_from_thd(thd);
/* this will trigger saving of perf_context information */
Rdb_perf_context_guard guard(tx, rocksdb_perf_context_level(thd));
if (tx != nullptr) { if (commit_tx || (!my_core::thd_test_options(
thd, OPTION_NOT_AUTOCOMMIT | OPTION_BEGIN))) { /* Thiswillnotaddanythingtocommit_latency_stats,andthisiscorrect right?
*/ if (tx->commit_ordered_done)
{
thd_wakeup_subsequent_commits(thd, 0);
DBUG_RETURN((tx->commit_ordered_res? HA_ERR_INTERNAL_ERROR: 0));
}
/* Wegethere -ForaCOMMITstatementthatfinishesamulti-statementtransaction -Forastatementthathasitsowntransaction
*/ if (thd->slave_thread)
{ // An attempt to make parallel slave performant (not fully successful, // see MDEV-15372):
// First, commit without syncing. This establishes the commit order
tx->set_sync(false); bool tx_had_writes = tx->get_write_count()? true : false ; if (tx->commit()) {
DBUG_RETURN(HA_ERR_ROCKSDB_COMMIT_FAILED);
}
thd_wakeup_subsequent_commits(thd, 0);
if (tx_had_writes && rocksdb_flush_log_at_trx_commit == FLUSH_LOG_SYNC)
{
rocksdb::Status s= rdb->FlushWAL(true); if (!s.ok())
DBUG_RETURN(HA_ERR_INTERNAL_ERROR);
}
} else
{ /* Not a slave thread */ if (tx->commit()) {
DBUG_RETURN(HA_ERR_ROCKSDB_COMMIT_FAILED);
}
}
} else { /* Wegetherewhencommittingastatementwithinatransaction.
*/
tx->make_stmt_savepoint_permanent();
}
if (my_core::thd_tx_isolation(thd) <= ISO_READ_COMMITTED) { // For READ_COMMITTED, we release any existing snapshot so that we will // see any changes that occurred since the last statement.
tx->release_snapshot();
}
}
// `Add()` is implemented in a thread-safe manner.
commit_latency_stats->Add(timer.ElapsedNanos() / 1000);
if (my_core::thd_tx_isolation(thd) <= ISO_READ_COMMITTED) { // For READ_COMMITTED, we release any existing snapshot so that we will // see any changes that occurred since the last statement.
tx->release_snapshot();
}
} return HA_EXIT_SUCCESS;
}
// Calculate how much space we will need int len = vsnprintf(nullptr, 0, format, args);
va_end(args);
if (len < 0) {
res = std::string("<format error>");
} elseif (len == 0) { // Shortcut for an empty string
res = std::string("");
} else { // For short enough output use a static buffer char *buff = static_buff;
std::unique_ptr<char[]> dynamic_buff = nullptr;
len++; // Add one for null terminator
// for longer output use an allocated buffer if (static_cast<uint>(len) > sizeof(static_buff)) {
dynamic_buff.reset(newchar[len]);
buff = dynamic_buff.get();
}
// Now re-do the vsnprintf with the buffer which is now large enough
(void)vsnprintf(buff, len, format, args_copy);
// Convert to a std::string. Note we could have created a std::string // large enough and then converted the buffer to a 'char*' and created // the output in place. This would probably work but feels like a hack. // Since this isn't code that needs to be super-performant we are going // with this 'safer' method.
res = std::string(buff);
}
va_end(args_copy);
return res;
}
class Rdb_snapshot_status : public Rdb_tx_list_walker { private:
std::string m_data;
if (stat_type == HA_ENGINE_STATUS) {
DBUG_ASSERT(rdb != nullptr);
std::string str;
/* Global DB Statistics */ if (rocksdb_stats) {
str = rocksdb_stats->ToString();
// Use the same format as internal RocksDB statistics entries to make // sure that output will look unified.
DBUG_ASSERT(commit_latency_stats != nullptr);
// Retrieve additional stalling related numbers from RocksDB and append // them to the buffer meant for displaying detailed statistics. The intent // here is to avoid adding another row to the query output because of // just two numbers. // // NB! We're replacing hyphens with underscores in output to better match // the existing naming convention. if (rdb->GetIntProperty("rocksdb.is-write-stopped", &v)) {
snprintf(buf, sizeof(buf), "rocksdb.is_write_stopped COUNT : %llu\n", (ulonglong)v);
str.append(buf);
}
/* Show the background thread status */
std::vector<rocksdb::ThreadStatus> thread_list;
rocksdb::Status s = rdb->GetEnv()->GetThreadList(&thread_list);
if (!s.ok()) { // NO_LINT_DEBUG
sql_print_error("RocksDB: Returned error (%s) from GetThreadList.\n",
s.ToString().c_str());
res |= true;
} else { /* For each background thread retrieved, print out its information */ for (auto &it : thread_list) { /* Only look at background threads. Ignore user threads, if any. */ if (it.thread_type > rocksdb::ThreadStatus::LOW_PRIORITY) { continue;
}
DBUG_ASSERT(op == snapshot_operation::SNAPSHOT_CREATE ||
op == snapshot_operation::SNAPSHOT_ATTACH);
// case: if binlogs are available get binlog file/pos and gtid info if (op == snapshot_operation::SNAPSHOT_CREATE && mysql_bin_log_is_open()) {
mysql_bin_log_lock_commits(ss_info);
}
if (op == snapshot_operation::SNAPSHOT_ATTACH) {
explicit_snapshot = Rdb_explicit_snapshot::get(ss_info->snapshot_id); if (!explicit_snapshot) {
my_printf_error(ER_UNKNOWN_ERROR, "Snapshot %llu does not exist", MYF(0),
ss_info->snapshot_id);
error = HA_EXIT_FAILURE;
}
} #endif
// case: all good till now if (error == HA_EXIT_SUCCESS) {
tx = get_or_create_tx(thd);
Rdb_perf_context_guard guard(tx, rocksdb_perf_context_level(thd));
#ifdef MARIADB_NOT_YET if (explicit_snapshot) {
tx->m_explicit_snapshot = explicit_snapshot;
} #endif
#ifdef MARIADB_NOT_YET // case: an explicit snapshot was not assigned to this transaction if (!tx->m_explicit_snapshot) {
tx->m_explicit_snapshot =
Rdb_explicit_snapshot::create(ss_info, rdb, tx->m_read_opts.snapshot); if (!tx->m_explicit_snapshot) {
my_printf_error(ER_UNKNOWN_ERROR, "Could not create snapshot", MYF(0));
error = HA_EXIT_FAILURE;
}
} #endif
}
#ifdef MARIADB_NOT_YET // case: unlock the binlog if (op == snapshot_operation::SNAPSHOT_CREATE && mysql_bin_log_is_open()) {
mysql_bin_log_unlock_commits(ss_info);
}
// copy over the snapshot details to pass to the upper layers if (tx->m_explicit_snapshot) {
*ss_info = tx->m_explicit_snapshot->ss_info;
ss_info->op = op;
} #endif
return error;
} #endif
/* Dummy SAVEPOINT support. This is needed for long running transactions *likemysqldump(https://bugs.mysql.com/bug.php?id=71017). *CurrentSAVEPOINTdoesnotcorrectlyhandleROLLBACKanddoesnotreturn *errors.Thisneedstobeaddressedinfutureversions(Issue#96).
*/ staticint rocksdb_savepoint(THD *const thd, void *const savepoint) { return HA_EXIT_SUCCESS;
}
if (rdb_normalize_tablename(it, &str) != HA_EXIT_SUCCESS) { /* Function needs to return void because of the interface and we've *detectedanerrorwhichshouldn'thappen.There'snowaytolet *callerknowthatsomethingfailed.
*/
SHIP_ASSERT(false); return;
}
if (rdb_split_normalized_tablename(str, &dbname, &tablename, &partname)) { continue;
}
is_partition = (partname.size() != 0);
table_handler = rdb_open_tables.get_table_handler(it.c_str()); if (table_handler == nullptr) { continue;
}
// If we're starting from scratch and there are no options saved yet then this // is a valid case. Therefore we can't compare the current set of options to // anything. if (status.IsNotFound()) { return rocksdb::Status::OK();
}
if (!status.ok()) { return status;
}
if (loaded_cf_descs.size() != cf_descr.size()) { return rocksdb::Status::NotSupported( "Mismatched size of column family " "descriptors.");
}
// Please see RocksDB documentation for more context about why we need to set // user-defined functions and pointer-typed options manually. for (size_t i = 0; i < loaded_cf_descs.size(); i++) {
loaded_cf_descs[i].options.compaction_filter =
cf_descr[i].options.compaction_filter;
loaded_cf_descs[i].options.compaction_filter_factory =
cf_descr[i].options.compaction_filter_factory;
loaded_cf_descs[i].options.comparator = cf_descr[i].options.comparator;
loaded_cf_descs[i].options.memtable_factory =
cf_descr[i].options.memtable_factory;
loaded_cf_descs[i].options.merge_operator =
cf_descr[i].options.merge_operator;
loaded_cf_descs[i].options.prefix_extractor =
cf_descr[i].options.prefix_extractor;
loaded_cf_descs[i].options.table_factory =
cf_descr[i].options.table_factory;
}
// This is the essence of the function - determine if it's safe to open the // database or not.
status = CheckOptionsCompatibility(dbpath, rocksdb::Env::Default(), main_opts,
loaded_cf_descs,
rocksdb_ignore_unknown_options);
staticvoid rocksdb_update_optimizer_costs(OPTIMIZER_COSTS *costs)
{ /* See optimizer_costs.txt for how these are calculated */
costs->row_next_find_cost= 0.00015161;
costs->row_lookup_cost= 0.00150453;
costs->key_next_find_cost= 0.00025108;
costs->key_lookup_cost= 0.00079369;
costs->row_copy_cost= 0.00006087;
}
if (rocksdb_ignore_datadic_errors)
{
sql_print_information( "CAUTION: Running with rocksdb_ignore_datadic_errors=1. " " This should only be used to perform repairs");
}
if (rdb_check_rocksdb_corruption()) { // NO_LINT_DEBUG
sql_print_error( "RocksDB: There was a corruption detected in RockDB files. " "Check error log emitted earlier for more details."); if (rocksdb_allow_to_start_after_corruption) { // NO_LINT_DEBUG
sql_print_information( "RocksDB: Remove rocksdb_allow_to_start_after_corruption to prevent " "server operating if RocksDB corruption is detected.");
} else { // NO_LINT_DEBUG
sql_print_error( "RocksDB: The server will exit normally and stop restart " "attempts. Remove %s file from data directory and " "start mariadbd manually.",
rdb_corruption_marker_file_name().c_str()); exit(0);
}
}
// Validate the assumption about the size of ROCKSDB_SIZEOF_HIDDEN_PK_COLUMN.
static_assert(sizeof(longlong) == 8, "Assuming that longlong is 8 bytes.");
if (rocksdb_db_options->max_open_files > (long)open_files_limit) { // NO_LINT_DEBUG
sql_print_information( "RocksDB: rocksdb_max_open_files should not be " "greater than the open_files_limit, effective value " "of rocksdb_max_open_files is being set to " "open_files_limit / 2.");
rocksdb_db_options->max_open_files = open_files_limit / 2;
} elseif (rocksdb_db_options->max_open_files == -2) {
rocksdb_db_options->max_open_files = open_files_limit / 2;
}
#if0// MARIAROCKS_NOT_YET : read-free replication is not supported
rdb_read_free_regex_handler.set_patterns(DEFAULT_READ_FREE_RPL_TABLES); #endif
if (rocksdb_db_options->allow_mmap_reads &&
rocksdb_db_options->use_direct_reads) { // allow_mmap_reads implies !use_direct_reads and RocksDB will not open if // mmap_reads and direct_reads are both on. (NO_LINT_DEBUG)
sql_print_error( "RocksDB: Can't enable both use_direct_reads " "and allow_mmap_reads\n");
DBUG_RETURN(HA_EXIT_FAILURE);
}
if (!check_status.ok()) { // NO_LINT_DEBUG
sql_print_error( "RocksDB: Unable to use direct io in rocksdb-datadir:" "(%s)",
check_status.getState());
DBUG_RETURN(HA_EXIT_FAILURE);
}
}
if (rocksdb_db_options->allow_mmap_writes &&
rocksdb_db_options->use_direct_io_for_flush_and_compaction) { // See above comment for allow_mmap_reads. (NO_LINT_DEBUG)
sql_print_error( "RocksDB: Can't enable both " "use_direct_io_for_flush_and_compaction and " "allow_mmap_writes\n");
DBUG_RETURN(HA_EXIT_FAILURE);
}
if (rocksdb_db_options->allow_mmap_writes &&
rocksdb_flush_log_at_trx_commit != FLUSH_LOG_NEVER) { // NO_LINT_DEBUG
sql_print_error( "RocksDB: rocksdb_flush_log_at_trx_commit needs to be 0 " "to use allow_mmap_writes");
DBUG_RETURN(HA_EXIT_FAILURE);
}
// sst_file_manager will move deleted rocksdb sst files to trash_dir // to be deleted in a background thread.
std::string trash_dir = std::string(rocksdb_datadir) + "/trash";
rocksdb_db_options->sst_file_manager.reset(NewSstFileManager(
rocksdb_db_options->env, myrocks_logger, trash_dir,
rocksdb_sst_mgr_rate_bytes_per_sec, true/* delete_existing_trash */));
std::vector<std::string> cf_names;
rocksdb::Status status;
status = rocksdb::DB::ListColumnFamilies(*rocksdb_db_options, rocksdb_datadir,
&cf_names); if (!status.ok()) { /* Whenwestartonanemptydatadir,ListColumnFamiliesreturnsIOError, andRocksDBdoesn'tprovideanywaytocheckwhatkindoferroritwas. Checkingsystemerrnohappenstoworkrightnow.
*/ if (status.IsIOError() #ifndef _WIN32
&& errno == ENOENT #endif
) {
sql_print_information("RocksDB: Got ENOENT when listing column families");
status =
check_rocksdb_options_compatibility(rocksdb_datadir, main_opts, cf_descr);
// We won't start if we'll determine that there's a chance of data corruption // because of incompatible options. if (!status.ok()) {
rdb_log_status_error(
status, "Compatibility check against existing database options failed");
DBUG_RETURN(HA_EXIT_FAILURE);
}
status = rocksdb::TransactionDB::Open(
main_opts, tx_db_options, rocksdb_datadir, cf_descr, &cf_handles, &rdb);
if (dict_manager.init(rdb, &cf_manager)) { // NO_LINT_DEBUG
sql_print_error("RocksDB: Failed to initialize data dictionary.");
DBUG_RETURN(HA_EXIT_FAILURE);
}
if (binlog_manager.init(&dict_manager)) { // NO_LINT_DEBUG
sql_print_error("RocksDB: Failed to initialize binlog manager.");
DBUG_RETURN(HA_EXIT_FAILURE);
}
if (ddl_manager.init(&dict_manager, &cf_manager, rocksdb_validate_tables)) { // NO_LINT_DEBUG
sql_print_error("RocksDB: Failed to initialize DDL manager.");
if (rocksdb_ignore_datadic_errors)
{
sql_print_error("RocksDB: rocksdb_ignore_datadic_errors=1, " "trying to continue");
} else
DBUG_RETURN(HA_EXIT_FAILURE);
}
// Creating an instance of HistogramImpl should only happen after RocksDB // has been successfully initialized.
commit_latency_stats = new rocksdb::HistogramImpl();
// Construct a list of directories which will be monitored by I/O watchdog // to make sure that we won't lose write access to them.
std::vector<std::string> directories;
// 1. Data directory.
directories.push_back(mysql_real_data_home);
if (rocksdb_pause_background_work)
rdb->ContinueBackgroundWork();
// signal the drop index thread to stop
rdb_drop_idx_thread.signal(true);
// Flush all memtables for not losing data, even if WAL is disabled.
rocksdb_flush_all_memtables();
// Stop all rocksdb background work
CancelAllBackgroundWork(rdb->GetBaseDB(), true);
// Signal the background thread to stop and to persist all stats collected // from background flushes and compactions. This will add more keys to a new // memtable, but since the memtables were just flushed, it should not trigger // a flush that can stall due to background threads being stopped. As long // as these keys are stored in a WAL file, they can be retrieved on restart.
rdb_bg_thread.signal(true);
// Wait for the background thread to finish. auto err = rdb_bg_thread.join(); if (err != 0) { // We'll log the message and continue because we're shutting down and // continuation is the optimal strategy. // NO_LINT_DEBUG
sql_print_error("RocksDB: Couldn't stop the background thread: (errno=%d)",
err);
}
// Wait for the drop index thread to finish.
err = rdb_drop_idx_thread.join(); if (err != 0) { // NO_LINT_DEBUG
sql_print_error("RocksDB: Couldn't stop the index thread: (errno=%d)", err);
}
// signal the manual compaction thread to stop
rdb_mc_thread.signal(true); // Wait for the manual compaction thread to finish.
err = rdb_mc_thread.join(); if (err != 0) { // NO_LINT_DEBUG
sql_print_error( "RocksDB: Couldn't stop the manual compaction thread: (errno=%d)", err);
}
// Disown the cache data since we're shutting down. // This results in memory leaks but it improved the shutdown time. // Don't disown when running under valgrind #ifndef HAVE_valgrind if (rocksdb_tbl_options->block_cache) {
rocksdb_tbl_options->block_cache->DisownData();
} #endif/* HAVE_valgrind */
#ifndef DBUG_OFF // simulate that RocksDB has reported corrupted data staticvoid dbug_change_status_to_corrupted(rocksdb::Status *status) {
*status = rocksdb::Status::Corruption();
} #endif
// If the iterator is not valid it might be because of EOF but might be due // to IOError or corruption. The good practice is always check it. // https://github.com/facebook/rocksdb/wiki/Iterator#error-handling staticinlinebool is_valid(rocksdb::Iterator *scan_it) { if (scan_it->Valid()) { returntrue;
} else {
rocksdb::Status s = scan_it->status();
DBUG_EXECUTE_IF("rocksdb_return_status_corrupted",
dbug_change_status_to_corrupted(&s);); if (s.IsIOError() || s.IsCorruption()) { if (s.IsCorruption()) {
rdb_persist_corruption_marker();
}
rdb_handle_io_error(s, RDB_IO_ERROR_GENERAL);
} returnfalse;
}
}
// First, look up the table in the hash map.
RDB_MUTEX_LOCK_CHECK(m_mutex); constauto it = m_table_map.find(table_name_str); if (it != m_table_map.end()) { // Found it
table_handler = it->second;
} else { char *tmp_name;
// Since we did not find it in the hash map, attempt to create and add it // to the hash map. if (!(table_handler = reinterpret_cast<Rdb_table_handler *>(my_multi_malloc(
PSI_INSTRUMENT_ME,
MYF(MY_WME | MY_ZEROFILL), &table_handler, sizeof(*table_handler),
&tmp_name, table_name_str.length() + 1, NullS)))) { // Allocating a new Rdb_table_handler and a new table name failed.
RDB_MUTEX_UNLOCK_CHECK(m_mutex); return nullptr;
}
if (use_datadic && dict_manager.get_auto_incr_val(
m_tbl_def->get_autoincr_gl_index_id(), &auto_incr)) {
update_auto_incr_val(auto_incr);
}
// If we find nothing in the data dictionary, or if we are in debug mode, // then call index_last to get the last value. // // This is needed when upgrading from a server that did not support // persistent auto_increment, of if the table is empty. // // For debug mode, we are just verifying that the data dictionary value is // greater than or equal to the maximum value in the table. if (auto_incr == 0 || validate_last) {
auto_incr = load_auto_incr_value_from_index();
update_auto_incr_val(auto_incr);
}
// If we failed to find anything from the data dictionary and index, then // initialize auto_increment to 1. if (m_tbl_def->m_auto_incr_val == 0) {
update_auto_incr_val(1);
}
}
// Do a lookup. We only need index column, so it should be index-only. // (another reason to make it index-only is that table->read_set is not set // appropriately and non-index-only lookup will not read the value) constbool save_keyread_only = m_keyread_only;
m_keyread_only = true;
m_converter->set_is_key_requested(true);
void ha_rocksdb::update_auto_incr_val(ulonglong val) {
ulonglong auto_incr_val = m_tbl_def->m_auto_incr_val; while (
auto_incr_val < val &&
!m_tbl_def->m_auto_incr_val.compare_exchange_weak(auto_incr_val, val)) { // Do nothing - just loop until auto_incr_val is >= val or we successfully // set it
}
}
void ha_rocksdb::update_auto_incr_val_from_field() {
Field *field;
ulonglong new_val, max_val;
field = table->key_info[table->s->next_number_index].key_part[0].field;
max_val = rdb_get_int_col_max_value(field);
MY_BITMAP *const old_map =
dbug_tmp_use_all_columns(table, &table->read_set);
new_val = field->val_int(); // don't increment if we would wrap around if (new_val != max_val) {
new_val++;
}
// Only update if positive value was set for auto_incr column. if (new_val <= max_val) {
Rdb_transaction *const tx = get_or_create_tx(table->in_use);
tx->set_auto_incr(m_tbl_def->get_autoincr_gl_index_id(), new_val);
// Update the in memory auto_incr value in m_tbl_def.
update_auto_incr_val(new_val);
}
}
longlong hidden_pk_id = 1; // Do a lookup. if (!index_last(table->record[0])) { /* DecodePKfieldfromthekey
*/ auto err = read_hidden_pk_id_from_rowkey(&hidden_pk_id); if (err) { if (is_new_snapshot) {
tx->release_snapshot();
} return err;
}
hidden_pk_id++;
}
longlong old = m_tbl_def->m_hidden_pk_val; while (old < hidden_pk_id &&
!m_tbl_def->m_hidden_pk_val.compare_exchange_weak(old, hidden_pk_id)) {
}
/* Get PK value from m_tbl_def->m_hidden_pk_info. */
longlong ha_rocksdb::update_hidden_pk_val() {
DBUG_ASSERT(has_hidden_pk(table)); const longlong new_val = m_tbl_def->m_hidden_pk_val++; return new_val;
}
/* Get the id of the hidden pk id from m_last_rowkey */ int ha_rocksdb::read_hidden_pk_id_from_rowkey(longlong *const hidden_pk_id) {
DBUG_ASSERT(table != nullptr);
DBUG_ASSERT(has_hidden_pk(table));
// Get hidden primary key from old key slice
Rdb_string_reader reader(&rowkey_slice); if ((!reader.read(Rdb_key_def::INDEX_NUMBER_SIZE))) { return HA_ERR_ROCKSDB_CORRUPT_DATA;
}
constint length= 8; /* was Field_longlong::PACK_LENGTH in FB MySQL tree */ const uchar *from = reinterpret_cast<const uchar *>(reader.read(length)); if (from == nullptr) { /* Mem-comparable image doesn't have enough bytes */ return HA_ERR_ROCKSDB_CORRUPT_DATA;
}
DBUG_ASSERT(table_handler != nullptr);
DBUG_ASSERT(table_handler->m_ref_count > 0); if (!--table_handler->m_ref_count) { // Last reference was released. Tear down the hash entry. constauto ret MY_ATTRIBUTE((__unused__)) =
m_table_map.erase(std::string(table_handler->m_table_name));
DBUG_ASSERT(ret == 1); // the hash entry must actually be found and deleted
my_core::thr_lock_delete(&table_handler->m_thr_lock);
my_free(table_handler);
}
/* Hide record if it has expired before the current snapshot time. */
uint64 read_filter_ts = 0; #ifndef DBUG_OFF
read_filter_ts += rdb_dbug_set_ttl_read_filter_ts(); #endif bool is_hide_ttl =
ts + kd.m_ttl_duration + read_filter_ts <= static_cast<uint64>(curr_ts); if (is_hide_ttl) {
update_row_stats(ROWS_FILTERED);
void dbug_modify_rec_varchar12(rocksdb::PinnableSlice *on_disk_rec) {
std::string res; // The record is NULL-byte followed by VARCHAR(10). // Put the NULL-byte
res.append("\0", 1); // Then, add a valid VARCHAR(12) value.
res.append("\xC", 1);
res.append("123456789ab", 12);
/* Sometimes, we may use m_sk_packed_tuple for storing packed PK */
max_packed_sk_len = pack_key_len; for (uint i = 0; i < table_arg->s->keys; i++) { /* Primary key was processed above */ if (i == table_arg->s->primary_key) continue;
// TODO: move this into get_table_handler() ??
kd_arr[i]->setup(table_arg, tbl_def_arg);
m_tbl_def = ddl_manager.find(fullname); if (m_tbl_def == nullptr) {
my_error(ER_INTERNAL_ERROR, MYF(0), "Attempt to open a table that is not present in RocksDB-SE data " "dictionary");
DBUG_RETURN(HA_ERR_ROCKSDB_INVALID_TABLE);
} if (m_tbl_def->m_key_count != table->s->keys + has_hidden_pk(table)? 1:0)
{
sql_print_error("MyRocks: DDL mismatch: .frm file has %u indexes, " "MyRocks has %u (%s hidden pk)",
table->s->keys, m_tbl_def->m_key_count,
has_hidden_pk(table)? "1" : "no");
if (rocksdb_ignore_datadic_errors)
{
sql_print_error("MyRocks: rocksdb_ignore_datadic_errors=1, " "trying to continue");
} else
{
my_error(ER_INTERNAL_ERROR, MYF(0), "MyRocks: DDL mismatch. Check the error log for details");
DBUG_RETURN(HA_ERR_ROCKSDB_INVALID_TABLE);
}
}
/* Load auto_increment value only once on first use. */ if (table->found_next_number_field && m_tbl_def->m_auto_incr_val == 0) {
load_auto_incr_value();
}
/* Load hidden pk only once on first use. */ if (has_hidden_pk(table) && m_tbl_def->m_hidden_pk_val == 0 &&
(err = load_hidden_pk_value()) != HA_EXIT_SUCCESS) {
free_key_buffers();
DBUG_RETURN(err);
}
/* Index block size in MyRocks: used by MySQL in query optimization */
stats.block_size = rocksdb_tbl_options->block_size;
#ifdef MARIAROCKS_NOT_YET // MDEV-10976 #endif /* Determine at open whether we should skip unique checks for this table */
set_skip_unique_check_tables(THDVAR(ha_thd(), skip_unique_check_tables));
if (m_table_handler != nullptr) {
rdb_open_tables.release_table_handler(m_table_handler);
m_table_handler = nullptr;
}
// These are needed to suppress valgrind errors in rocksdb.partition
m_last_rowkey.free();
m_sk_tails.free();
m_sk_tails_old.free();
m_pk_unpack_info.free();
DBUG_RETURN(HA_EXIT_SUCCESS);
}
staticconstchar *rdb_error_messages[] = { "Table must have a PRIMARY KEY.", "Specifying DATA DIRECTORY for an individual table is not supported.", "Specifying INDEX DIRECTORY for an individual table is not supported.", "RocksDB commit failed.", "Failure during bulk load operation.", "Found data corruption.", "CRC checksum mismatch.", "Invalid table.", "Could not access RocksDB properties.", "File I/O error during merge/sort operation.", "RocksDB status: not found.", "RocksDB status: corruption.", "RocksDB status: not supported.", "RocksDB status: invalid argument.", "RocksDB status: io error.", "RocksDB status: no space.", "RocksDB status: merge in progress.", "RocksDB status: incomplete.", "RocksDB status: shutdown in progress.", "RocksDB status: timed out.", "RocksDB status: aborted.", "RocksDB status: lock limit reached.", "RocksDB status: busy.", "RocksDB status: deadlock.", "RocksDB status: expired.", "RocksDB status: try again.",
};
static_assert((sizeof(rdb_error_messages) / sizeof(rdb_error_messages[0])) ==
((HA_ERR_ROCKSDB_LAST - HA_ERR_ROCKSDB_FIRST) + 1), "Number of error messages doesn't match number of error codes");
//psergey-merge: do we need this in MariaDB: we have get_error_messages //below... #if0 staticconstchar *rdb_get_error_message(int nr) { return rdb_error_messages[nr - HA_ERR_ROCKSDB_FIRST];
} #endif
// We can be called with the values which are < HA_ERR_FIRST because most // MySQL internal functions will just return HA_EXIT_FAILURE in case of // an error.
/* MyRocks supports only the following collations for indexed columns */ staticconst std::set<uint> RDB_INDEX_COLLATIONS = {
COLLATION_BINARY, COLLATION_UTF8_BIN, COLLATION_LATIN1_BIN};
staticbool rdb_is_index_collation_supported( const my_core::Field *const field) { const my_core::enum_field_types type = field->real_type(); /* Handle [VAR](CHAR|BINARY) or TEXT|BLOB */ if (type == MYSQL_TYPE_VARCHAR || type == MYSQL_TYPE_STRING ||
type == MYSQL_TYPE_BLOB) {
if ((err = Rdb_key_def::extract_ttl_col(table_arg, tbl_def_arg, &ttl_column,
&ttl_field_offset))) {
DBUG_RETURN(err);
}
/* We don't currently support TTL on tables with hidden primary keys. */ if (ttl_duration > 0 && has_hidden_pk(table_arg)) {
my_error(ER_RDB_TTL_UNSUPPORTED, MYF(0));
DBUG_RETURN(HA_EXIT_FAILURE);
}
/* Thefirstloopcheckstheindexparametersandcreates columnfamiliesifnecessary.
*/ for (uint i = 0; i < tbl_def_arg->m_key_count; i++) {
rocksdb::ColumnFamilyHandle *cf_handle;
if (!is_hidden_pk(i, table_arg, tbl_def_arg) &&
tbl_def_arg->base_tablename().find(tmp_file_prefix) != 0) { if (!tsys_set)
{
tsys_set= true;
my_core::filename_to_tablename(tbl_def_arg->base_tablename().c_str(),
tablename_sys, sizeof(tablename_sys));
}
for (uint part = 0; part < table_arg->key_info[i].ext_key_parts;
part++)
{ /* MariaDB: disallow NOPAD collations */ if (rdb_field_uses_nopad_collation(
table_arg->key_info[i].key_part[part].field))
{
my_error(ER_MYROCKS_CANT_NOPAD_COLLATION, MYF(0));
DBUG_RETURN(HA_EXIT_FAILURE);
}
if (rocksdb_strict_collation_check &&
!rdb_is_index_collation_supported(
table_arg->key_info[i].key_part[part].field) &&
!rdb_collation_exceptions->matches(tablename_sys)) {
char buf[1024];
my_snprintf(buf, sizeof(buf), "Indexed column %s.%s uses a collation that does not " "allow index-only access in secondary key and has " "reduced disk space efficiency in primary key.",
tbl_def_arg->full_tablename().c_str(),
table_arg->key_info[i].key_part[part].field->field_name.str);
// Internal consistency check to make sure that data in TABLE and // Rdb_tbl_def structures matches. Either both are missing or both are // specified. Yes, this is critical enough to make it into SHIP_ASSERT.
SHIP_ASSERT(IF_PARTITIONING(!table_arg->part_info,true) == tbl_def_arg->base_partition().empty());
// Generate the name for the column family to use. bool per_part_match_found = false;
std::string cf_name =
generate_cf_name(i, table_arg, tbl_def_arg, &per_part_match_found);
// Prevent create from using the system column family. if (cf_name == DEFAULT_SYSTEM_CF_NAME) {
my_error(ER_WRONG_ARGUMENTS, MYF(0), "column family not valid for storing index data.");
DBUG_RETURN(HA_EXIT_FAILURE);
}
// Here's how `get_or_create_cf` will use the input parameters: // // `cf_name` - will be used as a CF name.
cf_handle = cf_manager.get_or_create_cf(rdb, cf_name);
const GL_INDEX_ID gl_index_id = okd.get_gl_index_id(); struct Rdb_index_info index_info; if (!dict_manager.get_index_info(gl_index_id, &index_info)) { // NO_LINT_DEBUG
sql_print_error( "RocksDB: Could not get index information " "for Index Number (%u,%u), table %s",
gl_index_id.cf_id, gl_index_id.index_id,
old_tbl_def_arg->full_tablename().c_str());
DBUG_RETURN(HA_EXIT_FAILURE);
}
for (i = 0; i < tbl_def_arg->m_key_count; i++) {
new_key_pos[get_key_name(i, table_arg, tbl_def_arg)] = i;
}
for (i = 0; i < old_tbl_def_arg->m_key_count; i++) { if (is_hidden_pk(i, old_table_arg, old_tbl_def_arg)) {
old_key_pos[old_key_descr[i]->m_name] = i; continue;
}
if (compare_keys(old_key, new_key) && !unique_to_non_unique) { continue;
}
/* Check to make sure key parts match. */ if (compare_key_parts(old_key, new_key)) { continue;
}
old_key_pos[old_key->name.str] = i;
}
DBUG_RETURN(old_key_pos);
}
/* Check to see if two keys are identical. */ int ha_rocksdb::compare_keys(const KEY *const old_key, const KEY *const new_key) const {
DBUG_ENTER_FUNC();
/* Check index name. */ if (strcmp(old_key->name.str, new_key->name.str) != 0) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
/* If index algorithms are different then keys are different. */ if (old_key->algorithm != new_key->algorithm) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
/* Check that the key is identical between old and new tables. */ if ((old_key->flags ^ new_key->flags) & HA_KEYFLAG_MASK) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
/* Check index comment. (for column family changes) */
std::string old_comment(old_key->comment.str, old_key->comment.length);
std::string new_comment(new_key->comment.str, new_key->comment.length); if (old_comment.compare(new_comment) != 0) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
DBUG_RETURN(HA_EXIT_SUCCESS);
}
/* Check two keys to ensure that key parts within keys match */ int ha_rocksdb::compare_key_parts(const KEY *const old_key, const KEY *const new_key) const {
DBUG_ENTER_FUNC();
/* Skip if key parts do not match, as it is a different key */ if (new_key->user_defined_key_parts != old_key->user_defined_key_parts) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
/* Check to see that key parts themselves match */ for (uint i = 0; i < old_key->user_defined_key_parts; i++) { if (strcmp(old_key->key_part[i].field->field_name.str,
new_key->key_part[i].field->field_name.str) != 0) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
/* Check if prefix index key part length has changed */ if (old_key->key_part[i].length != new_key->key_part[i].length) {
DBUG_RETURN(HA_EXIT_FAILURE);
}
}
// Use PRIMARY_FORMAT_VERSION_UPDATE1 here since it is the same value as // SECONDARY_FORMAT_VERSION_UPDATE1 so it doesn't matter if this is a // primary key or secondary key.
DBUG_EXECUTE_IF("MYROCKS_LEGACY_VARBINARY_FORMAT", {
kv_version = Rdb_key_def::PRIMARY_FORMAT_VERSION_UPDATE1;
});
while (*str != '\0') { // Scan from our current pos looking for 'FOREIGN'
str = rdb_find_in_string(str, "FOREIGN", &success); if (!success) { returnfalse;
}
// Skip past the found "FOREIGN'
str = rdb_check_next_token(&my_charset_bin, str, "FOREIGN", &success);
DBUG_ASSERT(success);
if (!my_isspace(&my_charset_bin, *str)) { returnfalse;
}
// See if the next token is 'KEY'
str = rdb_check_next_token(&my_charset_bin, str, "KEY", &success); if (!success) { continue;
}
// See if the next token is '('
str = rdb_check_next_token(&my_charset_bin, str, "(", &success); if (!success) { // There is an optional index id after 'FOREIGN KEY', skip it
str = rdb_skip_id(&my_charset_bin, str);
// Now check for '(' again
str = rdb_check_next_token(&my_charset_bin, str, "(", &success);
}
// If we have found 'FOREIGN KEY [<word>] (' we can be confident we have // a foreign key clause. return success;
}
// We never found a valid foreign key clause returnfalse;
}
/* Create table/key descriptions and put them into the data dictionary */
m_tbl_def = new Rdb_tbl_def(table_name);
uint n_keys = table_arg->s->keys;
/* Ifnoprimarykeyfound,createahiddenPKandplaceitinsidetable definition
*/ if (has_hidden_pk(table_arg)) {
n_keys += 1; // reset hidden pk id // the starting valid value for hidden pk is 1
m_tbl_def->m_hidden_pk_val = 1;
}
m_key_descr_arr = new std::shared_ptr<Rdb_key_def>[n_keys];
m_tbl_def->m_key_count = n_keys;
m_tbl_def->m_key_descr_arr = m_key_descr_arr;
if (create_info->data_file_name) { // DATA DIRECTORY is used to create tables under a specific location // outside the MySQL data directory. We don't support this for MyRocks. // The `rocksdb_datadir` setting should be used to configure RocksDB data // directory.
my_error(WARN_OPTION_IGNORED, ME_NOTE, "DATA DIRECTORY");
}
if (create_info->index_file_name && table_arg->s->keys) { // Similar check for INDEX DIRECTORY as well.
my_error(WARN_OPTION_IGNORED, ME_NOTE, "INDEX DIRECTORY");
}
int err; /* Constructdbname.tablenameourselves,becausepartitioning passesstringslike"./test/t14#P#p0"forindividualpartitions, whiletable_arg->s->table_namehasnoneofthat.
*/
std::string str;
err = rdb_normalize_tablename(name, &str); if (err != HA_EXIT_SUCCESS) {
DBUG_RETURN(err);
}
// FOREIGN KEY isn't supported yet
THD *const thd = _current_thd(); if (contains_foreign_key(thd)) {
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "FOREIGN KEY for the RocksDB storage engine");
DBUG_RETURN(HA_ERR_UNSUPPORTED);
}
// Check whether Data Dictionary contain information
Rdb_tbl_def *tbl = ddl_manager.find(str); if (tbl != nullptr) { if (thd->lex->sql_command == SQLCOM_TRUNCATE) {
err = delete_table(tbl); if (err != HA_EXIT_SUCCESS) {
DBUG_RETURN(err);
}
} else {
my_error(ER_METADATA_INCONSISTENCY, MYF(0), str.c_str(), name);
DBUG_RETURN(HA_ERR_ROCKSDB_CORRUPT_DATA);
}
}
// The below adds/clears hooks in RocksDB sync points. There's no reason for // this code to be in ::create() but it needs to be somewhere where it is // away from any tight loops and where one can invoke it from mtr:
DBUG_EXECUTE_IF("rocksdb_enable_delay_commits",
{ auto syncpoint= rocksdb::SyncPoint::GetInstance();
syncpoint->SetCallBack("DBImpl::WriteImpl:BeforeLeaderEnters",
[&](void* /*arg*/) {my_sleep(500);} );
syncpoint->EnableProcessing();
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN, ER_WRONG_ARGUMENTS, "enable_delay_commits_mode ON");
Rdb_field_packing dummy1;
res = dummy1.setup(nullptr, key_info->key_part[part].field, inx, part,
key_info->key_part[part].length);
if (res && all_parts) { for (uint i = 0; i < part; i++) {
Field *field; if ((field = key_info->key_part[i].field)) {
Rdb_field_packing dummy; if (!dummy.setup(nullptr, field, inx, i,
key_info->key_part[i].length)) { /* Cannot do index-only reads for this column */
res = false; break;
}
}
}
}
switch (find_flag) { case HA_READ_KEY_EXACT:
rc = read_key_exact(kd, m_scan_it, full_key_match, key_slice,
ttl_filter_ts); break; case HA_READ_BEFORE_KEY:
*move_forward = false;
rc = read_before_key(kd, full_key_match, key_slice, ttl_filter_ts); if (rc == 0 && !kd.covers_key(m_scan_it->key())) { /* The record we've got is not from this index */
rc = HA_ERR_KEY_NOT_FOUND;
} break; case HA_READ_AFTER_KEY: case HA_READ_KEY_OR_NEXT:
rc = read_after_key(kd, key_slice, ttl_filter_ts); if (rc == 0 && !kd.covers_key(m_scan_it->key())) { /* The record we've got is not from this index */
rc = HA_ERR_KEY_NOT_FOUND;
} break; case HA_READ_KEY_OR_PREV: case HA_READ_PREFIX: /* This flag is not used by the SQL layer, so we don't support it yet. */
rc = HA_ERR_UNSUPPORTED; break; case HA_READ_PREFIX_LAST: case HA_READ_PREFIX_LAST_OR_PREV:
*move_forward = false; /* Findthelastrecordwiththespecifiedindexprefixlookup. -HA_READ_PREFIX_LASTrequiresthattherecordhasthe prefix=lookup(iftherearenosuchrecords, HA_ERR_KEY_NOT_FOUNDshouldbereturned). -HA_READ_PREFIX_LAST_OR_PREVhasnosuchrequirement.Ifthereareno recordswithprefix=lookup,weshouldreturnthelastrecord beforethat.
*/
rc = read_before_key(kd, full_key_match, key_slice, ttl_filter_ts); if (rc == 0) { const rocksdb::Slice &rkey = m_scan_it->key(); if (!kd.covers_key(rkey)) { /* The record we've got is not from this index */
rc = HA_ERR_KEY_NOT_FOUND;
} elseif (find_flag == HA_READ_PREFIX_LAST) {
uint size = kd.pack_index_tuple(table, m_pack_buffer,
m_sk_packed_tuple, m_record_buffer,
key, keypart_map);
rocksdb::Slice lookup_tuple( reinterpret_cast<char *>(m_sk_packed_tuple), size);
// We need to compare the key we've got with the original search // prefix. if (!kd.value_matches_prefix(rkey, lookup_tuple)) {
rc = HA_ERR_KEY_NOT_FOUND;
}
}
} break; default:
DBUG_ASSERT(0); break;
}
if (m_lock_rows != RDB_LOCK_NONE) { /* We need to put a lock and re-read */
rc = get_row_by_rowid(buf, m_pk_packed_tuple, pk_size);
} else { /* Unpack from the row we've read */ const rocksdb::Slice &value = m_scan_it->value();
rc = convert_record_from_storage_format(&rkey, &value, buf);
}
return rc;
}
int ha_rocksdb::read_row_from_secondary_key(uchar *const buf, const Rdb_key_def &kd, bool move_forward) { int rc = 0;
uint pk_size= 0;
/* Get the key columns and primary key value */ const rocksdb::Slice &rkey = m_scan_it->key(); const rocksdb::Slice &value = m_scan_it->value();
@details m_scan_itpointsattheindexkey-valuepairthatweshouldreadthe(pk,row) pairfor.
*/ int ha_rocksdb::secondary_index_read(constint keyno, uchar *const buf) {
DBUG_ASSERT(table != nullptr); #ifdef MARIAROCKS_NOT_YET
stats.rows_requested++; #endif /* Use STATUS_NOT_FOUND when record not found or some error occurred */
table->status = STATUS_NOT_FOUND;
if (is_valid(m_scan_it)) {
rocksdb::Slice key = m_scan_it->key();
/* Check if we've ran out of records of this index */ if (m_key_descr_arr[keyno]->covers_key(key)) { int rc = 0;
// TODO: We could here check if we have ran out of range we're scanning const uint size = m_key_descr_arr[keyno]->get_primary_key_tuple(
table, *m_pk_descr, &key, m_pk_packed_tuple); if (size == RDB_INVALID_KEY_LEN) { return HA_ERR_ROCKSDB_CORRUPT_DATA;
}
if (!start_key) { // Read first record
result = ha_index_first(table->record[0]);
} else { #ifdef MARIAROCKS_NOT_YET if (is_using_prohibited_gap_locks(
is_using_full_unique_key(active_index, start_key->keypart_map,
start_key->flag))) {
DBUG_RETURN(HA_ERR_LOCK_DEADLOCK);
} #endif
increment_statistics(&SSV::ha_read_key_count);
result =
index_read_map_impl(table->record[0], start_key->key,
start_key->keypart_map, start_key->flag, end_key);
} if (result) {
DBUG_RETURN((result == HA_ERR_KEY_NOT_FOUND) ? HA_ERR_END_OF_FILE : result);
}
Rdb_transaction *const tx = get_or_create_tx(table->in_use); constbool is_new_snapshot = !tx->has_snapshot(); // Loop as long as we get a deadlock error AND we end up creating the // snapshot here (i.e. it did not exist prior to this) for (;;) {
DEBUG_SYNC(thd, "rocksdb.check_flags_rmi_scan"); if (thd && thd->killed) {
rc = HA_ERR_QUERY_INTERRUPTED; break;
} /* Thiswillopentheiteratorandpositionitatarecordthat'sequalor greaterthanthelookuptuple.
*/
setup_scan_iterator(kd, &slice, use_all_keys, eq_cond_len);
if (!kd.covers_key(rkey)) {
table->status = STATUS_NOT_FOUND; return HA_ERR_END_OF_FILE;
}
if (m_sk_match_prefix) { const rocksdb::Slice prefix((constchar *)m_sk_match_prefix,
m_sk_match_length); if (!kd.value_matches_prefix(rkey, prefix)) {
table->status = STATUS_NOT_FOUND; return HA_ERR_END_OF_FILE;
}
}
const rocksdb::Slice value = m_scan_it->value(); int err = kd.unpack_record(table, buf, &rkey, &value,
m_converter->get_verify_row_debug_checksums()); if (err != HA_EXIT_SUCCESS) { return err;
}
const check_result_t icp_status= handler_index_cond_check(this); if (icp_status == CHECK_NEG) {
rocksdb_smart_next(!move_forward, m_scan_it); continue; /* Get the next (or prev) index tuple */
} elseif (icp_status == CHECK_OUT_OF_RANGE ||
icp_status == CHECK_ABORTED_BY_USER) { /* We have walked out of range we are scanning */
table->status = STATUS_NOT_FOUND; return HA_ERR_END_OF_FILE;
} else/* icp_status == CHECK_POS */
{ /* Index Condition is satisfied. We have rc==0, proceed to fetch the
* row. */ break;
}
}
} return HA_EXIT_SUCCESS;
}
// Only when debugging: don't use snapshot when reading // Rdb_transaction *tx= get_or_create_tx(table->in_use); // tx->snapshot= nullptr;
bool save_verify_row_debug_checksums =
m_converter->get_verify_row_debug_checksums();
m_converter->set_verify_row_debug_checksums(true); /* For each secondary index, check that we can get a PK value from it */ // NO_LINT_DEBUG
sql_print_verbose_info("CHECKTABLE %s: Checking table %s", table_name,
table_name);
ha_rows UNINIT_VAR(row_checksums_at_start); // set/used iff first_index==true
ha_rows row_checksums = ha_rows(-1); bool first_index = true;
for (uint keyno = 0; keyno < table->s->keys; keyno++) { if (keyno != pk) {
extra(HA_EXTRA_KEYREAD);
ha_index_init(keyno, true);
ha_rows rows = 0;
ha_rows checksums = 0; if (first_index) {
row_checksums_at_start = m_converter->get_row_checksums_checked();
} int res; // NO_LINT_DEBUG
sql_print_verbose_info("CHECKTABLE %s: Checking index %s", table_name,
table->key_info[keyno].name.str); while (1) { if (!rows) {
res = index_first(table->record[0]);
} else {
res = index_next(table->record[0]);
}
// do nothing - we already have the result in m_retrieved_record and // already taken the lock
s = rocksdb::Status::OK();
} else {
s = get_for_update(tx, m_pk_descr->get_cf(), key_slice,
&m_retrieved_record);
}
if (!s.IsNotFound() && !s.ok()) {
DBUG_RETURN(tx->set_status_error(table->in_use, s, *m_pk_descr, m_tbl_def,
m_table_handler));
}
found = !s.IsNotFound();
table->status = STATUS_NOT_FOUND; if (found) { /* If we found the record, but it's expired, pretend we didn't find it. */ if (!skip_ttl_check && m_pk_descr->has_ttl() &&
should_hide_ttl_rec(*m_pk_descr, m_retrieved_record,
tx->m_snapshot_timestamp)) {
DBUG_RETURN(HA_ERR_KEY_NOT_FOUND);
}
constbool is_new_snapshot = !tx->has_snapshot(); // Loop as long as we get a deadlock error AND we end up creating the // snapshot here (i.e. it did not exist prior to this) for (;;) {
setup_scan_iterator(kd, &index_key, false, key_start_matching_bytes);
m_scan_it->Seek(index_key);
m_skip_scan_it_next_call = true;
rc = index_next_with_direction(buf, true); if (!should_recreate_snapshot(rc, is_new_snapshot)) { break; /* exit the loop */
}
// release the snapshot and iterator so they will be regenerated
tx->release_snapshot();
release_scan_iterator();
}
bool is_new_snapshot = !tx->has_snapshot(); // Loop as long as we get a deadlock error AND we end up creating the // snapshot here (i.e. it did not exist prior to this) for (;;) {
setup_scan_iterator(kd, &index_key, false, key_end_matching_bytes);
m_scan_it->SeekForPrev(index_key);
m_skip_scan_it_next_call = false;
/* Returns true if given index number is a primary key */ bool ha_rocksdb::is_pk(const uint index, const TABLE *const table_arg, const Rdb_tbl_def *const tbl_def_arg) {
DBUG_ASSERT(table_arg->s != nullptr);
return index == table_arg->s->primary_key ||
is_hidden_pk(index, table_arg, tbl_def_arg);
}
// When creating CF-s the caller needs to know if there was a custom CF name // specified for a given partition.
*per_part_match_found = false;
// Index comment is used to define the column family name specification(s). // If there was no comment, we get an emptry string, and it means "use the // default column family". constchar *const comment = get_key_comment(index, table_arg, tbl_def_arg);
if (IF_PARTITIONING(table_arg->part_info,nullptr) != nullptr && !*per_part_match_found) { // At this point we tried to search for a custom CF name for a partition, // but none was specified. Therefore default one will be used. return"";
}
// If we didn't find any partitioned/non-partitioned qualifiers, return the // comment itself. NOTE: this currently handles returning the cf name // specified in the index comment in the case of no partitions, which doesn't // use any qualifiers at the moment. (aka its a special case) if (cf_name.empty() && !key_comment.empty()) { return key_comment;
}
/* Note:"buf==table->record[0]"iscopiedfrominnodb.Iamnotawareof anyusecaseswherethisconditionisnottrue.
*/ if (table->next_number_field && buf == table->record[0]) { int err; if ((err = update_auto_increment())) {
DBUG_RETURN(err);
}
}
// clear cache at beginning of write for INSERT ON DUPLICATE // we may get multiple write->fail->read->update if there are multiple // values from INSERT
m_dup_pk_found = false;
if (key_found && row_info.old_data == nullptr && m_insert_with_update) { // In INSERT ON DUPLICATE KEY UPDATE ... case, if the insert failed // due to a duplicate key, remember the last key and skip the check // next time
m_dup_pk_found = true;
#ifndef DBUG_OFF // save it for sanity checking later
m_dup_pk_retrieved_record.copy(m_retrieved_record.data(),
m_retrieved_record.size(), &my_charset_bin); #endif
}
int ha_rocksdb::bulk_load_key(Rdb_transaction *const tx, const Rdb_key_def &kd, const rocksdb::Slice &key, const rocksdb::Slice &value, bool sort) {
DBUG_ENTER_FUNC(); int res;
THD *thd = ha_thd(); if (thd && thd->killed) {
DBUG_RETURN(HA_ERR_QUERY_INTERRUPTED);
}
rocksdb::ColumnFamilyHandle *cf = kd.get_cf();
// In the case of unsorted inserts, m_sst_info allocated here is not // used to store the keys. It is still used to indicate when tables // are switched. if (m_sst_info == nullptr || m_sst_info->is_done()) {
m_sst_info.reset(new Rdb_sst_info(rdb, m_table_handler->m_table_name,
kd.get_name(), cf, *rocksdb_db_options,
THDVAR(ha_thd(), trace_sst_api)));
res = tx->start_bulk_load(this, m_sst_info); if (res != HA_EXIT_SUCCESS) {
DBUG_RETURN(res);
}
}
DBUG_ASSERT(m_sst_info);
if (sort) {
Rdb_index_merge *key_merge;
DBUG_ASSERT(cf != nullptr);
res = tx->get_key_merge(kd.get_gl_index_id(), cf, &key_merge); if (res == HA_EXIT_SUCCESS) {
res = key_merge->add(key, value);
}
} else {
res = m_sst_info->put(key, value);
}
DBUG_RETURN(res);
}
int ha_rocksdb::finalize_bulk_load(bool print_client_error) {
DBUG_ENTER_FUNC();
int res = HA_EXIT_SUCCESS;
/* Skip if there are no possible ongoing bulk loads */ if (m_sst_info) { if (m_sst_info->is_done()) {
m_sst_info.reset();
DBUG_RETURN(res);
}
Rdb_sst_info::Rdb_sst_commit_info commit_info;
// Wrap up the current work in m_sst_info and get ready to commit // This transfer the responsibility of commit over to commit_info
res = m_sst_info->finish(&commit_info, print_client_error); if (res == 0) { // Make sure we have work to do - under race condition we could lose // to another thread and end up with no work if (commit_info.has_work()) {
rocksdb::IngestExternalFileOptions opts;
opts.move_files = true;
opts.snapshot_consistency = false;
opts.allow_global_seqno = false;
opts.allow_blocking_flush = false;
const rocksdb::Status s = rdb->IngestExternalFile(
commit_info.get_cf(), commit_info.get_committed_files(), opts); if (!s.ok()) { if (print_client_error) {
Rdb_sst_info::report_error_msg(s, nullptr);
}
res = HA_ERR_ROCKSDB_BULK_LOAD;
} else { // Mark the list of SST files as committed, otherwise they'll get // cleaned up when commit_info destructs
commit_info.commit();
}
}
}
m_sst_info.reset();
}
DBUG_RETURN(res);
}
if (table->found_next_number_field) {
update_auto_incr_val_from_field();
}
int rc = HA_EXIT_SUCCESS;
rocksdb::Slice value_slice; /* Prepare the new record to be written into RocksDB */ if ((rc = m_converter->encode_value_slice(
m_pk_descr, row_info.new_pk_slice, row_info.new_pk_unpack_info,
!row_info.old_pk_slice.empty(), should_store_row_debug_checksums(),
m_ttl_bytes, &m_ttl_bytes_updated, &value_slice))) { return rc;
}
@param[in]row_infoholdallrowdata,suchasoldkey/newkey @param[in]pk_changedwhetherprimarykeyischanged @return HA_EXIT_SUCCESSOK OtherHA_ERRerrorcode(canbeSE-specific)
*/ int ha_rocksdb::update_write_indexes(conststruct update_row_info &row_info, constbool pk_changed) { int rc; bool bulk_load_sk;
// The PK must be updated first to pull out the TTL value.
rc = update_write_pk(*m_pk_descr, row_info, pk_changed); if (rc != HA_EXIT_SUCCESS) { return rc;
}
// Update the remaining indexes. Allow bulk loading only if // allow_sk is enabled
bulk_load_sk = rocksdb_enable_bulk_load_api &&
THDVAR(table->in_use, bulk_load) &&
THDVAR(table->in_use, bulk_load_allow_sk); for (uint key_id = 0; key_id < m_tbl_def->m_key_count; key_id++) { if (is_pk(key_id, table, m_tbl_def)) { continue;
}
const rocksdb::Slice eq_cond(slice->data(), eq_cond_len); // The size of m_scan_it_lower_bound (and upper) is technically // max_packed_sk_len as calculated in ha_rocksdb::alloc_key_buffers. Rather // than recalculating that number, we pass in the max of eq_cond_len and // Rdb_key_def::INDEX_NUMBER_SIZE which is guaranteed to be smaller than // max_packed_sk_len, hence ensuring no buffer overrun. // // See ha_rocksdb::setup_iterator_bounds on how the bound_len parameter is // used. if (check_bloom_and_set_bounds(
ha_thd(), kd, eq_cond, use_all_keys,
std::max(eq_cond_len, (uint)Rdb_key_def::INDEX_NUMBER_SIZE),
m_scan_it_lower_bound, m_scan_it_upper_bound,
&m_scan_it_lower_bound_slice, &m_scan_it_upper_bound_slice)) {
skip_bloom = false;
}
// when this table is being updated, decode all fields
m_converter->setup_field_decoders(table->read_set,
m_lock_rows == RDB_LOCK_WRITE);
if (scan) {
m_rnd_scan_is_new_snapshot = !tx->has_snapshot();
setup_iterator_for_rnd_scan();
} else { /* We don't need any preparations for rnd_pos() calls. */
}
// If m_lock_rows is on then we will be doing a get_for_update when accessing // the index, so don't acquire the snapshot right away. Otherwise acquire // the snapshot immediately.
tx->acquire_snapshot(m_lock_rows == RDB_LOCK_NONE);
int rc; for (;;) {
rc = rnd_next_with_direction(buf, true); if (!should_recreate_snapshot(rc, m_rnd_scan_is_new_snapshot)) { break; /* exit the loop */
} // release the snapshot and iterator and then regenerate them
Rdb_transaction *tx = get_or_create_tx(table->in_use);
tx->release_snapshot();
release_scan_iterator();
setup_iterator_for_rnd_scan();
}
m_rnd_scan_is_new_snapshot = false;
if (rc == HA_ERR_KEY_NOT_FOUND) rc = HA_ERR_END_OF_FILE;
for (;;) {
DEBUG_SYNC(thd, "rocksdb.check_flags_rnwd"); if (thd && thd->killed) {
rc = HA_ERR_QUERY_INTERRUPTED; break;
}
if (m_skip_scan_it_next_call) {
m_skip_scan_it_next_call = false;
} else { if (move_forward) {
m_scan_it->Next(); /* this call cannot fail */
} else {
m_scan_it->Prev(); /* this call cannot fail */
}
}
if (!is_valid(m_scan_it)) {
rc = HA_ERR_END_OF_FILE; break;
}
/* check if we're out of this table */ const rocksdb::Slice key = m_scan_it->key(); if (!m_pk_descr->covers_key(key)) {
rc = HA_ERR_END_OF_FILE; break;
}
if (m_lock_rows != RDB_LOCK_NONE) { /* Locktherowwe'vejustread.
if (m_pk_descr->has_ttl() &&
should_hide_ttl_rec(*m_pk_descr, m_scan_it->value(),
tx->m_snapshot_timestamp)) { continue;
}
const rocksdb::Status s =
get_for_update(tx, m_pk_descr->get_cf(), key, &m_retrieved_record); if (s.IsNotFound() &&
should_skip_invalidated_record(HA_ERR_KEY_NOT_FOUND)) { continue;
}
if (!s.ok()) {
DBUG_RETURN(tx->set_status_error(table->in_use, s, *m_pk_descr,
m_tbl_def, m_table_handler));
}
// If we called get_for_update() use the value from that call not from // the iterator as it may be stale since we don't have a snapshot // when m_lock_rows is not RDB_LOCK_NONE.
m_last_rowkey.copy(key.data(), key.size(), &my_charset_bin);
rc = convert_record_from_storage_format(&key, buf);
} else { // Use the value from the iterator
rocksdb::Slice value = m_scan_it->value();
if (m_pk_descr->has_ttl() &&
should_hide_ttl_rec(
*m_pk_descr, value,
get_or_create_tx(table->in_use)->m_snapshot_timestamp)) { continue;
}
m_start_range= NULL; // when this table is being updated, decode all fields
m_converter->setup_field_decoders(table->read_set,
m_lock_rows == RDB_LOCK_WRITE);
if (!m_keyread_only) {
m_key_descr_arr[idx]->get_lookup_bitmap(table, &m_lookup_bitmap);
}
// If m_lock_rows is not RDB_LOCK_NONE then we will be doing a get_for_update // when accessing the index, so don't acquire the snapshot right away. // Otherwise acquire the snapshot immediately.
tx->acquire_snapshot(m_lock_rows == RDB_LOCK_NONE);
active_index = idx;
DBUG_RETURN(HA_EXIT_SUCCESS);
}
/** @return HA_EXIT_SUCCESSOK
*/ int ha_rocksdb::index_end() {
DBUG_ENTER_FUNC();
/** @return HA_EXIT_SUCCESSOK otherHA_ERRerrorcode(canbeSE-specific)
*/ int ha_rocksdb::truncate() {
DBUG_ENTER_FUNC();
DBUG_ASSERT(m_tbl_def != nullptr);
// Save table name to use later
std::string table_name = m_tbl_def->full_tablename();
// Delete current table int err = delete_table(m_tbl_def); if (err != HA_EXIT_SUCCESS) {
DBUG_RETURN(err);
}
// Reset auto_increment_value to 1 if auto-increment feature is enabled // By default, the starting valid value for auto_increment_value is 1
DBUG_RETURN(create_table(
table_name, table,
table->found_next_number_field ? 1 : 0/* auto_increment_value */));
}
const uint index = pk_index(table, m_tbl_def);
rocksdb::Status s =
delete_or_singledelete(index, tx, m_pk_descr->get_cf(), key_slice); if (!s.ok()) {
DBUG_RETURN(tx->set_status_error(table->in_use, s, *m_pk_descr, m_tbl_def,
m_table_handler));
} else {
bytes_written = key_slice.size();
}
longlong hidden_pk_id = 0; if (m_tbl_def->m_key_count > 1 && has_hidden_pk(table)) { int err = read_hidden_pk_id_from_rowkey(&hidden_pk_id); if (err) {
DBUG_RETURN(err);
}
}
// Delete the record for every secondary index for (uint i = 0; i < m_tbl_def->m_key_count; i++) { if (!is_pk(i, table, m_tbl_def)) { int packed_size; const Rdb_key_def &kd = *m_key_descr_arr[i];
packed_size = kd.pack_record(table, m_pack_buffer, buf, m_sk_packed_tuple,
nullptr, false, hidden_pk_id);
rocksdb::Slice secondary_key_slice( reinterpret_cast<constchar *>(m_sk_packed_tuple), packed_size); /* Deleting on secondary key doesn't need any locks: */
tx->get_indexed_write_batch()->SingleDelete(kd.get_cf(),
secondary_key_slice);
bytes_written += secondary_key_slice.size();
}
}
// if number of records is hardcoded, we do not want to force computation // of memtable cardinalities if (stats.records == 0 || (rocksdb_force_compute_memtable_stats &&
rocksdb_debug_optimizer_n_rows == 0)) { // First, compute SST files stats
uchar buf[Rdb_key_def::INDEX_NUMBER_SIZE * 2]; auto r = get_range(pk_index(table, m_tbl_def), buf);
uint64_t sz = 0;
uint8_t include_flags = rocksdb::DB::INCLUDE_FILES; // recompute SST files stats only if records count is 0 if (stats.records == 0) {
rdb->GetApproximateSizes(m_pk_descr->get_cf(), &r, 1, &sz,
include_flags);
stats.records += sz / ROCKSDB_ASSUMED_KEY_VALUE_DISK_SIZE;
stats.data_file_length += sz;
} // Second, compute memtable stats. This call is expensive, so cache // values computed for some time.
uint64_t cachetime = rocksdb_force_compute_memtable_stats_cachetime;
uint64_t time = (cachetime == 0) ? 0 : my_interval_timer() / 1000; if (cachetime == 0 ||
time > m_table_handler->m_mtcache_last_update + cachetime) {
uint64_t memtableCount;
uint64_t memtableSize;
// the stats below are calculated from skiplist which is a probabilistic // data structure, so the results vary between test runs // it also can return 0 for quite a large tables which means that // cardinality for memtable only indxes will be reported as 0
rdb->GetApproximateMemTableStats(m_pk_descr->get_cf(), r,
&memtableCount, &memtableSize);
// Atomically update all of these fields at the same time if (cachetime > 0) { if (m_table_handler->m_mtcache_lock.fetch_add( 1, std::memory_order_acquire) == 0) {
m_table_handler->m_mtcache_count = memtableCount;
m_table_handler->m_mtcache_size = memtableSize;
m_table_handler->m_mtcache_last_update = time;
}
m_table_handler->m_mtcache_lock.fetch_sub(1,
std::memory_order_release);
}
stats.records += memtableCount;
stats.data_file_length += memtableSize;
} else { // Cached data is still valid, so use it instead
stats.records += m_table_handler->m_mtcache_count;
stats.data_file_length += m_table_handler->m_mtcache_size;
}
// Do like InnoDB does. stats.records=0 confuses the optimizer if (stats.records == 0 && !(flag & (HA_STATUS_TIME | HA_STATUS_OPEN))) {
stats.records++;
}
}
if (rocksdb_debug_optimizer_n_rows > 0)
stats.records = rocksdb_debug_optimizer_n_rows;
if (flag & HA_STATUS_CONST) {
ref_length = m_pk_descr->max_storage_fmt_length();
for (uint i = 0; i < m_tbl_def->m_key_count; i++) { if (is_hidden_pk(i, table, m_tbl_def)) { continue;
}
KEY *const k = &table->key_info[i]; for (uint j = 0; j < k->ext_key_parts; j++) { const Rdb_index_stats &k_stats = m_key_descr_arr[i]->m_stats;
uint x;
if (k_stats.m_distinct_keys_per_prefix.size() > j &&
k_stats.m_distinct_keys_per_prefix[j] > 0) {
x = k_stats.m_rows / k_stats.m_distinct_keys_per_prefix[j]; /* Ifthenumberofrowsislessthanthenumberofprefixes(dueto sampling),theaveragenumberofrowswiththesameprefixis1.
*/ if (x == 0) {
x = 1;
}
} else {
x = 0;
} if (x > stats.records) x = stats.records; if ((x == 0 && rocksdb_debug_optimizer_no_zero_cardinality) ||
rocksdb_debug_optimizer_n_rows > 0) { // Fake cardinality implementation. For example, (idx1, idx2, idx3) // index /* MakeMariaRocksbehavethesamewayasMyRocksdoes: 1.SQLlayerthinksthatuniquesecondaryindexesarenotextended withPKcolumns(bothinMySQLandMariaDB) 2.MariaDBalsothinksthatindexeswithpartially-coveredcolumns arenotextendedwithPKcolumns.Usethesamenumberof keypartsthatMyRockswoulduse.
*/
uint ext_key_parts2; if (k->flags & HA_NOSAME)
ext_key_parts2= k->ext_key_parts; // This is #1 else
ext_key_parts2= m_key_descr_arr[i]->get_key_parts(); // This is #2.
// will have rec_per_key for (idx1)=4, (idx1,2)=2, and (idx1,2,3)=1. // rec_per_key for the whole index is 1, and multiplied by 2^n if // n suffix columns of the index are not used.
x = 1 << (ext_key_parts2 - j - 1);
}
k->rec_per_key[j] = x;
}
}
longlong hidden_pk_id = 0; if (has_hidden_pk(table) && read_hidden_pk_id_from_rowkey(&hidden_pk_id)) {
DBUG_ASSERT(false); // should never reach here
}
/* The following function was copied from ha_blackhole::store_lock: */
THR_LOCK_DATA **ha_rocksdb::store_lock(THD *const thd, THR_LOCK_DATA **to, enum thr_lock_type lock_type) {
DBUG_ENTER_FUNC();
/* Then, tell the SQL layer what kind of locking it should use: */ if (lock_type != TL_IGNORE && m_db_lock.type == TL_UNLOCK) { /* Hereiswherewegetintothegutsofarowlevellock. IfTL_UNLOCKisset IfwearenotdoingaLOCKTABLEorDISCARD/IMPORT TABLESPACE,thenallowmultiplewriters
*/
for (;;) { // The stop flag might be set by shutdown command // after drop_index_thread releases signal_mutex // (i.e. while executing expensive Seek()). To prevent drop_index_thread // from entering long cond_timedwait, checking if stop flag // is true or not is needed, with drop_index_interrupt_mutex held. if (m_stop) { break;
}
timespec ts; int sec= dict_manager.is_drop_index_empty()
? 24 * 60 * 60// no filtering
: 60; // filtering
set_timespec(ts,sec);
constauto ret MY_ATTRIBUTE((__unused__)) =
mysql_cond_timedwait(&m_signal_cond, &m_signal_mutex, &ts); if (m_stop) { break;
} // make sure, no program error is returned
DBUG_ASSERT(ret == 0 || ret == ETIMEDOUT);
RDB_MUTEX_UNLOCK_CHECK(m_signal_mutex);
// If the user changed the database part of the name then validate that the // 'to' database exists. if (from_db != to_db && !rdb_database_exists(to_db)) { // If we return a RocksDB specific error code here we get // "error: 206 - Unknown error 206". InnoDB gets // "error -1 - Unknown error -1" so let's match them.
DBUG_RETURN(-1);
}
// this function is needed only for online alter-table
DBUG_RETURN(COMPATIBLE_DATA_NO);
}
/** @return HA_EXIT_SUCCESSOK
*/ int ha_rocksdb::extra(enum ha_extra_function operation) {
DBUG_ENTER_FUNC();
switch (operation) { case HA_EXTRA_KEYREAD:
m_keyread_only = true; break; case HA_EXTRA_NO_KEYREAD:
m_keyread_only = false; break; case HA_EXTRA_FLUSH: /* Ifthetablehasblobs,thentheyarepartofm_retrieved_record. Thiscallinvalidatesthem.
*/
m_retrieved_record.Reset(); break; case HA_EXTRA_INSERT_WITH_UPDATE: // INSERT ON DUPLICATE KEY UPDATE if (rocksdb_enable_insert_with_update_caching) {
m_insert_with_update = true;
} break; case HA_EXTRA_NO_IGNORE_DUP_KEY: // PAIRED with HA_EXTRA_INSERT_WITH_UPDATE or HA_EXTRA_WRITE_CAN_REPLACE // that indicates the end of REPLACE / INSERT ON DUPLICATE KEY
m_insert_with_update = false; break;
for (uint i = 0; i < table->s->keys; i++) {
uchar buf[Rdb_key_def::INDEX_NUMBER_SIZE * 2]; auto range = get_range(i, buf); const rocksdb::Status s = rdb->CompactRange(getCompactRangeOptions(),
m_key_descr_arr[i]->get_cf(),
&range.start, &range.limit); if (!s.ok()) {
DBUG_RETURN(rdb_error_to_mysql(s));
}
}
// find per column family key ranges which need to be queried
std::unordered_map<rocksdb::ColumnFamilyHandle *, std::vector<rocksdb::Range>>
ranges;
std::unordered_map<GL_INDEX_ID, Rdb_index_stats> stats;
std::vector<uchar> buf(to_recalc.size() * 2 * Rdb_key_def::INDEX_NUMBER_SIZE);
uchar r_buf[Rdb_key_def::INDEX_NUMBER_SIZE * 2]; auto r = myrocks::get_range(*kd, r_buf);
uint64_t memtableCount;
uint64_t memtableSize;
rdb->GetApproximateMemTableStats(kd->get_cf(), r, &memtableCount,
&memtableSize); if (memtableCount < (uint64_t)stat.m_rows / 10) { // skip tables that already have enough stats from SST files to reduce // overhead and avoid degradation of big tables stats by sampling from // relatively tiny (less than 10% of full data set) memtable dataset continue;
}
std::unique_ptr<rocksdb::Iterator> it =
std::unique_ptr<rocksdb::Iterator>(
rdb->NewIterator(read_opts, kd->get_cf()));
cardinality_collector.Reset(); for (it->Seek(first_index_key); is_valid(it.get()); it->Next()) { const rocksdb::Slice key = it->key(); if (!kd->covers_key(key)) { break; // end of this index
}
stat.m_rows++;
if (table) { if (calculate_stats_for_table() != HA_EXIT_SUCCESS) {
DBUG_RETURN(HA_ADMIN_FAILED);
}
}
// A call to ::info is needed to repopulate some SQL level structs. This is // necessary for online analyze because we cannot rely on another ::open // call to call info for us. if (info(HA_STATUS_CONST | HA_STATUS_VARIABLE) != HA_EXIT_SUCCESS) {
DBUG_RETURN(HA_ADMIN_FAILED);
}
Field *field;
ulonglong new_val, max_val;
field = table->key_info[table->s->next_number_index].key_part[0].field;
max_val = rdb_get_int_col_max_value(field);
// Local variable reference to simplify code below auto &auto_incr = m_tbl_def->m_auto_incr_val;
if (inc == 1) {
DBUG_ASSERT(off == 1); // Optimization for the standard case where we are always simply // incrementing from the last position
// Use CAS operation in a loop to atomically get the next auto // increment value while ensuring that we don't wrap around to a negative // number. // // We set auto_incr to the min of max_val and new_val + 1. This means that // if we're at the maximum, we should be returning the same value for // multiple rows, resulting in duplicate key errors (as expected). // // If we return values greater than the max, the SQL layer will "truncate" // the value anyway, but it means that we store invalid values into // auto_incr that will be visible in SHOW CREATE TABLE.
new_val = auto_incr; while (new_val != std::numeric_limits<ulonglong>::max()) { if (auto_incr.compare_exchange_weak(new_val,
std::min(new_val + 1, max_val))) { break;
}
}
} else { // The next value can be more complicated if either 'inc' or 'off' is not 1
ulonglong last_val = auto_incr;
if (last_val > max_val) {
new_val = std::numeric_limits<ulonglong>::max();
} else { // Loop until we can correctly update the atomic value do {
DBUG_ASSERT(last_val > 0); // Calculate the next value in the auto increment series: offset // + N * increment where N is 0, 1, 2, ... // // For further information please visit: // http://dev.mysql.com/doc/refman/5.7/en/replication-options-master.html // // The following is confusing so here is an explanation: // To get the next number in the sequence above you subtract out the // offset, calculate the next sequence (N * increment) and then add the // offset back in. // // The additions are rearranged to avoid overflow. The following is // equivalent to (last_val - 1 + inc - off) / inc. This uses the fact // that (a+b)/c = a/c + b/c + (a%c + b%c)/c. To show why: // // (a+b)/c // = (a - a%c + a%c + b - b%c + b%c) / c // = (a - a%c) / c + (b - b%c) / c + (a%c + b%c) / c // = a/c + b/c + (a%c + b%c) / c // // Now, substitute a = last_val - 1, b = inc - off, c = inc to get the // following statement.
ulonglong n =
(last_val - 1) / inc + ((last_val - 1) % inc + inc - off) / inc;
// Check if n * inc + off will overflow. This can only happen if we have // an UNSIGNED BIGINT field. if (n > (std::numeric_limits<ulonglong>::max() - off) / inc) {
DBUG_ASSERT(max_val == std::numeric_limits<ulonglong>::max()); // The 'last_val' value is already equal to or larger than the largest // value in the sequence. Continuing would wrap around (technically // the behavior would be undefined). What should we do? // We could: // 1) set the new value to the last possible number in our sequence // as described above. The problem with this is that this // number could be smaller than a value in an existing row. // 2) set the new value to the largest possible number. This number // may not be in our sequence, but it is guaranteed to be equal // to or larger than any other value already inserted. // // For now I'm going to take option 2. // // Returning ULLONG_MAX from get_auto_increment will cause the SQL // layer to fail with ER_AUTOINC_READ_FAILED. This means that due to // the SE API for get_auto_increment, inserts will fail with // ER_AUTOINC_READ_FAILED if the column is UNSIGNED BIGINT, but // inserts will fail with ER_DUP_ENTRY for other types (or no failure // if the column is in a non-unique SK).
new_val = std::numeric_limits<ulonglong>::max();
auto_incr = new_val; // Store the largest value into auto_incr break;
}
new_val = n * inc + off;
// Attempt to store the new value (plus 1 since m_auto_incr_val contains // the next available value) into the atomic value. If the current // value no longer matches what we have in 'last_val' this will fail and // we will repeat the loop (`last_val` will automatically get updated // with the current value). // // See above explanation for inc == 1 for why we use std::min.
} while (!auto_incr.compare_exchange_weak(
last_val, std::min(new_val + 1, max_val)));
}
}
/* We don't support unique keys on table w/ no primary keys */ if ((ha_alter_info->handler_flags &
ALTER_ADD_UNIQUE_INDEX) &&
has_hidden_pk(altered_table)) {
DBUG_RETURN(my_core::HA_ALTER_INPLACE_NOT_SUPPORTED);
}
/* We only support changing auto_increment for table options. */ if ((ha_alter_info->handler_flags & ALTER_CHANGE_CREATE_OPTION) &&
!(ha_alter_info->create_info->used_fields & HA_CREATE_USED_AUTO)) {
DBUG_RETURN(my_core::HA_ALTER_INPLACE_NOT_SUPPORTED);
}
/* FIXME: MDEV-16099 Use alter algorithm=nocopy or algorithm=instant
for non-InnoDB engine */
DBUG_RETURN(my_core::HA_ALTER_INPLACE_COPY_LOCK);
}
/* Ifaddinguniqueindex,allocatespecialbuffersforduplicatechecking.
*/ int err; if ((err = alloc_key_buffers(
altered_table, ctx->m_new_tdef,
ha_alter_info->handler_flags &
ALTER_ADD_UNIQUE_INDEX))) {
my_error(ER_OUT_OF_RESOURCES, MYF(0));
DBUG_RETURN(err);
}
/* Populate all new secondary keys by scanning the primary key. */ if ((err = inplace_populate_sk(altered_table, ctx->m_added_indexes))) {
my_error(ER_SK_POPULATE_DURING_ALTER, MYF(0));
DBUG_RETURN(HA_EXIT_FAILURE);
}
}
/* Scan each record in the primary key in order */ for (res = index_first(table->record[0]); res == 0;
res = index_next(table->record[0])) {
longlong hidden_pk_id = 0; if (hidden_pk_exists &&
(res = read_hidden_pk_id_from_rowkey(&hidden_pk_id))) { // NO_LINT_DEBUG
sql_print_error("Error retrieving hidden pk id.");
ha_index_end();
DBUG_RETURN(res);
}
/* Here,res==-1meansthatwearefinished,while>0meansanerror occurred.
*/ if (res > 0) { // NO_LINT_DEBUG
sql_print_error("Error while bulk loading keys in external merge sort.");
DBUG_RETURN(res);
}
bool is_critical_error;
res = tx->finish_bulk_load(&is_critical_error); if (res && is_critical_error) { // NO_LINT_DEBUG
sql_print_error("Error finishing bulk load.");
DBUG_RETURN(res);
}
}
/* Explicitlytelljemalloctocleanupanyunuseddirtypagesatthispoint. Seehttps://reviews.facebook.net/D63723 for more details.
*/
purge_all_jemalloc_arenas();
Forpartitionedtables,arollbackcalltothisfunction(commit==false) isdoneforeachpartition.Asuccessfulcommitcallonlyexecutesonce forallpartitions.
*/ if (!commit) { /* If ctx has not been created yet, nothing to do here */ if (!ctx0) {
DBUG_RETURN(HA_EXIT_SUCCESS);
}
/* CannotcalldestructorforRdb_tbl_defdirectlybecausewedon'twantto erasethemappingsinsidetheddl_manager,astheold_key_descrisstill usingthem.
*/ if (ctx0->m_new_key_descr) { /* Delete the new key descriptors */ for (uint i = 0; i < ctx0->m_new_tdef->m_key_count; i++) {
ctx0->m_new_key_descr[i] = nullptr;
}
if (dict_manager.commit(batch)) { /* Shouldneverreachhere.WeassumeMyRockswillabortifcommitfails.
*/
DBUG_ASSERT(0);
}
dict_manager.unlock();
/* Mark ongoing create indexes as finished/remove from data dictionary */
dict_manager.finish_indexes_operation(
create_index_ids, Rdb_key_def::DDL_CREATE_INDEX_ONGOING);
for (;;) { // Wait until the next timeout or until we receive a signal to stop the // thread. Request to stop the thread should only be triggered when the // storage engine is being unloaded.
RDB_MUTEX_LOCK_CHECK(m_signal_mutex); constauto ret MY_ATTRIBUTE((__unused__)) =
mysql_cond_timedwait(&m_signal_cond, &m_signal_mutex, &ts_next_sync);
// Check that we receive only the expected error codes.
DBUG_ASSERT(ret == 0 || ret == ETIMEDOUT); constbool local_stop = m_stop; constbool local_save_stats = m_save_stats;
reset();
RDB_MUTEX_UNLOCK_CHECK(m_signal_mutex);
if (local_stop) { // If we're here then that's because condition variable was signaled by // another thread and we're shutting down. Break out the loop to make // sure that shutdown thread can proceed. break;
}
// This path should be taken only when the timer expired.
DBUG_ASSERT(ret == ETIMEDOUT);
if (local_save_stats) {
ddl_manager.persist_stats();
}
// Set the next timestamp for mysql_cond_timedwait() (which ends up calling // pthread_cond_timedwait()) to wait on.
set_timespec(ts_next_sync, WAKE_UP_INTERVAL);
// Flush the WAL. Sync it for both background and never modes to copy // InnoDB's behavior. For mode never, the wal file isn't even written, // whereas background writes to the wal file, but issues the syncs in a // background thread. if (rdb && (rocksdb_flush_log_at_trx_commit != FLUSH_LOG_SYNC) &&
!rocksdb_db_options->allow_mmap_writes) { const rocksdb::Status s = rdb->FlushWAL(true); if (!s.ok()) {
rdb_handle_io_error(s, RDB_IO_ERROR_BG_THREAD);
}
} // Recalculate statistics for indexes. if (rocksdb_stats_recalc_rate) {
std::unordered_map<GL_INDEX_ID, std::shared_ptr<const Rdb_key_def>>
to_recalc;
if (rdb_indexes_to_recalc.empty()) { struct Rdb_index_collector : public Rdb_tables_scanner { int add_table(Rdb_tbl_def *tdef) override { for (uint i = 0; i < tdef->m_key_count; i++) {
rdb_indexes_to_recalc.push_back(
tdef->m_key_descr_arr[i]->get_gl_index_id());
} return HA_EXIT_SUCCESS;
}
} collector;
ddl_manager.scan_for_tables(&collector);
}
constauto ret MY_ATTRIBUTE((__unused__)) =
mysql_cond_timedwait(&m_signal_cond, &m_signal_mutex, &ts); if (m_stop) { break;
} // make sure, no program error is returned
DBUG_ASSERT(ret == 0 || ret == ETIMEDOUT);
RDB_MUTEX_UNLOCK_CHECK(m_signal_mutex);
RDB_MUTEX_LOCK_CHECK(m_mc_mutex); // Grab the first item and proceed, if not empty. if (m_requests.empty()) {
RDB_MUTEX_UNLOCK_CHECK(m_mc_mutex);
RDB_MUTEX_LOCK_CHECK(m_signal_mutex); continue;
}
Manual_compaction_request &mcr = m_requests.begin()->second;
DBUG_ASSERT(mcr.cf != nullptr);
DBUG_ASSERT(mcr.state == Manual_compaction_request::INITED);
mcr.state = Manual_compaction_request::RUNNING;
RDB_MUTEX_UNLOCK_CHECK(m_mc_mutex);
DBUG_ASSERT(mcr.state == Manual_compaction_request::RUNNING); // NO_LINT_DEBUG
sql_print_information("Manual Compaction id %d cf %s started.", mcr.mc_id,
mcr.cf->GetName().c_str());
rocksdb_manual_compactions_running++; if (rocksdb_debug_manual_compaction_delay > 0) {
my_sleep(rocksdb_debug_manual_compaction_delay * 1000000);
} // CompactRange may take a very long time. On clean shutdown, // it is cancelled by CancelAllBackgroundWork, then status is // set to shutdownInProgress. const rocksdb::Status s = rdb->CompactRange(
getCompactRangeOptions(mcr.concurrency), mcr.cf, mcr.start, mcr.limit);
rocksdb_manual_compactions_running--; if (s.ok()) { // NO_LINT_DEBUG
sql_print_information("Manual Compaction id %d cf %s ended.", mcr.mc_id,
mcr.cf->GetName().c_str());
} else { // NO_LINT_DEBUG
sql_print_information("Manual Compaction id %d cf %s aborted. %s",
mcr.mc_id, mcr.cf->GetName().c_str(), s.getState()); if (!s.IsShutdownInProgress()) {
rdb_handle_io_error(s, RDB_IO_ERROR_BG_THREAD);
} else {
DBUG_ASSERT(m_requests.size() == 1);
}
}
rocksdb_manual_compactions_processed++;
clear_manual_compaction_request(mcr.mc_id, false);
RDB_MUTEX_LOCK_CHECK(m_signal_mutex);
}
clear_all_manual_compaction_requests();
DBUG_ASSERT(m_requests.empty());
RDB_MUTEX_UNLOCK_CHECK(m_signal_mutex);
mysql_mutex_destroy(&m_mc_mutex);
}
void Rdb_manual_compaction_thread::clear_manual_compaction_request( int mc_id, bool init_only) { bool erase = true;
RDB_MUTEX_LOCK_CHECK(m_mc_mutex); auto it = m_requests.find(mc_id); if (it != m_requests.end()) { if (init_only) {
Manual_compaction_request mcr = it->second; if (mcr.state != Manual_compaction_request::INITED) {
erase = false;
}
} if (erase) {
m_requests.erase(it);
}
} else { // Current code path guarantees that erasing by the same mc_id happens // at most once. INITED state may be erased by a thread that requested // the compaction. RUNNING state is erased by mc thread only.
DBUG_ASSERT(0);
}
RDB_MUTEX_UNLOCK_CHECK(m_mc_mutex);
}
constchar *get_rdb_io_error_string(const RDB_IO_ERROR_TYPE err_type) { // If this assertion fails then this means that a member has been either added // to or removed from RDB_IO_ERROR_TYPE enum and this function needs to be // changed to return the appropriate value.
static_assert(RDB_IO_ERROR_LAST == 4, "Please handle all the error types.");
switch (err_type) { case RDB_IO_ERROR_TYPE::RDB_IO_ERROR_TX_COMMIT: return"RDB_IO_ERROR_TX_COMMIT"; case RDB_IO_ERROR_TYPE::RDB_IO_ERROR_DICT_COMMIT: return"RDB_IO_ERROR_DICT_COMMIT"; case RDB_IO_ERROR_TYPE::RDB_IO_ERROR_BG_THREAD: return"RDB_IO_ERROR_BG_THREAD"; case RDB_IO_ERROR_TYPE::RDB_IO_ERROR_GENERAL: return"RDB_IO_ERROR_GENERAL"; default:
DBUG_ASSERT(false); return"(unknown)";
}
}
// In case of core dump generation we want this function NOT to be optimized // so that we can capture as much data as possible to debug the root cause // more efficiently. #ifdef __GNUC__ #endif void rdb_handle_io_error(const rocksdb::Status status, const RDB_IO_ERROR_TYPE err_type) { if (status.IsIOError()) { /* skip dumping core if write failed and we are allowed to do so */ #ifdef MARIAROCKS_NOT_YET if (skip_core_dump_on_error) {
opt_core_file = false;
} #endif switch (err_type) { case RDB_IO_ERROR_TX_COMMIT: case RDB_IO_ERROR_DICT_COMMIT: {
rdb_log_status_error(status, "failed to write to WAL"); /* NO_LINT_DEBUG */
sql_print_error("MyRocks: aborting on WAL write error.");
abort(); break;
} case RDB_IO_ERROR_BG_THREAD: {
rdb_log_status_error(status, "BG thread failed to write to RocksDB"); /* NO_LINT_DEBUG */
sql_print_error("MyRocks: aborting on BG write error.");
abort(); break;
} case RDB_IO_ERROR_GENERAL: {
rdb_log_status_error(status, "failed on I/O"); /* NO_LINT_DEBUG */
sql_print_error("MyRocks: aborting on I/O error.");
abort(); break;
} default:
DBUG_ASSERT(0); break;
}
} elseif (status.IsCorruption()) {
rdb_log_status_error(status, "data corruption detected!");
rdb_persist_corruption_marker(); /* NO_LINT_DEBUG */
sql_print_error("MyRocks: aborting because of data corruption.");
abort();
} elseif (!status.ok()) { switch (err_type) { case RDB_IO_ERROR_DICT_COMMIT: {
rdb_log_status_error(status, "Failed to write to WAL (dictionary)"); /* NO_LINT_DEBUG */
sql_print_error("MyRocks: aborting on WAL write error.");
abort(); break;
} default:
rdb_log_status_error(status, "Failed to read/write in RocksDB"); break;
}
}
} #ifdef __GNUC__ #endif
Rdb_dict_manager *rdb_get_dict_manager(void) { return &dict_manager; }
if (new_val != rocksdb_table_stats_sampling_pct) {
rocksdb_table_stats_sampling_pct = new_val;
if (properties_collector_factory) {
properties_collector_factory->SetTableStatsSamplingPct(
rocksdb_table_stats_sampling_pct);
}
}
RDB_MUTEX_UNLOCK_CHECK(rdb_sysvars_mutex);
}
/* Thisfunctionallowssettingtheratelimiter'sbytespersecondvalue butonlyiftheratelimiteristurnedonwhichhastobedoneatstartup. Iftherateisalready0(turnedoff)orwearechangingitto0(trying toturnitoff)thisfunctionwillpushawarningtotheclientanddo nothing. Thisissimilartothecodeininnodb_doublewrite_update(foundin storage/innobase/handler/ha_innodb.cc).
*/ void rocksdb_set_rate_limiter_bytes_per_sec(
my_core::THD *const thd,
my_core::st_mysql_sys_var *const var MY_ATTRIBUTE((__unused__)), void *const var_ptr MY_ATTRIBUTE((__unused__)), constvoid *const save) { const uint64_t new_val = *static_cast<const uint64_t *>(save); if (new_val == 0 || rocksdb_rate_limiter_bytes_per_sec == 0) { /* Ifarate_limiterwasnotenabledatstartupwecan'tchangeitnor canwedisableitifonewascreatedatstartup
*/
push_warning_printf(thd, Sql_condition::WARN_LEVEL_WARN, ER_WRONG_ARGUMENTS, "RocksDB: rocksdb_rate_limiter_bytes_per_sec cannot " "be dynamically changed to or from 0. Do a clean " "shutdown if you want to change it from or to 0.");
} elseif (new_val != rocksdb_rate_limiter_bytes_per_sec) { /* Apply the new value to the rate limiter and store it locally */
DBUG_ASSERT(rocksdb_rate_limiter != nullptr);
rocksdb_rate_limiter_bytes_per_sec = new_val;
rocksdb_rate_limiter->SetBytesPerSecond(new_val);
}
}
//psergey-todo: what is the purpose of the below?? constchar *val_copy= val? my_strdup(PSI_INSTRUMENT_ME, val, MYF(0)): nullptr;
my_free(*static_cast<char**>(var_ptr));
*static_cast<constchar**>(var_ptr) = val_copy;
}
staticint rocksdb_validate_update_cf_options(
THD * /* unused */, struct st_mysql_sys_var * /*unused*/, void *save, struct st_mysql_value *value) { char buff[STRING_BUFFER_USUAL_SIZE]; constchar *str; int length;
length = sizeof(buff);
str = value->val_str(value, buff, &length); // In some cases, str can point to buff in the stack. // This can cause invalid memory access after validation is finished. // To avoid this kind case, let's alway duplicate the str if str is not // nullptr
*(constchar **)save = (str == nullptr) ? nullptr : my_strdup(PSI_INSTRUMENT_ME, str, MYF(0));
if (str == nullptr) { return HA_EXIT_SUCCESS;
}
Rdb_cf_options::Name_to_config_t option_map;
// Basic sanity checking and parsing the options into a map. If this fails // then there's no point to proceed. if (!Rdb_cf_options::parse_cf_options(str, &option_map)) {
my_error(ER_WRONG_VALUE_FOR_VAR, MYF(0), "rocksdb_update_cf_options", str); // Free what we've copied with my_strdup above.
my_free((void*)(*(constchar **)save)); return HA_EXIT_FAILURE;
} // Loop through option_map and create missing column families for (Rdb_cf_options::Name_to_config_t::iterator it = option_map.begin();
it != option_map.end(); ++it) {
cf_manager.get_or_create_cf(rdb, it->first);
} return HA_EXIT_SUCCESS;
}
if (!val) {
*reinterpret_cast<char **>(var_ptr) = nullptr;
RDB_MUTEX_UNLOCK_CHECK(rdb_sysvars_mutex); return;
}
DBUG_ASSERT(val != nullptr);
// Reset the pointers regardless of how much success we had with updating // the CF options. This will results in consistent behavior and avoids // dealing with cases when only a subset of CF-s was successfully updated.
*reinterpret_cast<constchar **>(var_ptr) = val;
// Do the real work of applying the changes.
Rdb_cf_options::Name_to_config_t option_map;
// This should never fail, because of rocksdb_validate_update_cf_options if (!Rdb_cf_options::parse_cf_options(val, &option_map)) {
my_free(*reinterpret_cast<char**>(var_ptr));
RDB_MUTEX_UNLOCK_CHECK(rdb_sysvars_mutex); return;
}
// For each CF we have, see if we need to update any settings. for (constauto &cf_name : cf_manager.get_cf_names()) {
DBUG_ASSERT(!cf_name.empty());
if (!per_cf_options.empty()) {
Rdb_cf_options::Name_to_config_t opt_map;
rocksdb::Status s = rocksdb::StringToMap(per_cf_options, &opt_map);
if (s != rocksdb::Status::OK()) { // NO_LINT_DEBUG
sql_print_warning( "MyRocks: failed to convert the options for column " "family '%s' to a map. %s",
cf_name.c_str(), s.ToString().c_str());
} else {
DBUG_ASSERT(rdb != nullptr);
// Finally we can apply the options.
s = rdb->SetOptions(cfh, opt_map);
if (s != rocksdb::Status::OK()) { // NO_LINT_DEBUG
sql_print_warning( "MyRocks: failed to apply the options for column " "family '%s'. %s",
cf_name.c_str(), s.ToString().c_str());
} else { // NO_LINT_DEBUG
sql_print_information( "MyRocks: options for column family '%s' " "have been successfully updated.",
cf_name.c_str());
// Make sure that data is internally consistent as well and update // the CF options. This is necessary also to make sure that the CF // options will be correctly reflected in the relevant table: // ROCKSDB_CF_OPTIONS in INFORMATION_SCHEMA.
rocksdb::ColumnFamilyOptions cf_options = rdb->GetOptions(cfh);
std::string updated_options;
s = rocksdb::GetStringFromColumnFamilyOptions(&updated_options,
cf_options);
auto s = dict_manager.get_value(rocksdb::Slice(lookup_key), &value); if (s.IsNotFound()) {
res = 0;
} elseif (s.ok()) { // decode the value if (value.length() == sizeof(res)) {
memcpy(&res, value.data(), sizeof(res));
res = be64toh(res);
} else
res = ulonglong(-1);
} else {
res = ulonglong(-1);
} return res;
}
maria_declare_plugin(rocksdb_se){
MYSQL_STORAGE_ENGINE_PLUGIN, /* Plugin Type */
&rocksdb_storage_engine, /* Plugin Descriptor */ "ROCKSDB", /* Plugin Name */ "Monty Program Ab", /* Plugin Author */ "RocksDB storage engine", /* Plugin Description */
PLUGIN_LICENSE_GPL, /* Plugin Licence */
myrocks::rocksdb_init_func, /* Plugin Entry Point */
myrocks::rocksdb_done_func, /* Plugin Deinitializer */ 0x0001, /* version number (0.1) */
myrocks::rocksdb_status_vars, /* status variables */
myrocks::rocksdb_system_variables, /* system variables */ "1.0", /* string version */
myrocks::MYROCKS_MARIADB_PLUGIN_MATURITY_LEVEL
},
myrocks::rdb_i_s_cfstats, myrocks::rdb_i_s_dbstats,
myrocks::rdb_i_s_perf_context, myrocks::rdb_i_s_perf_context_global,
myrocks::rdb_i_s_cfoptions, myrocks::rdb_i_s_compact_stats,
myrocks::rdb_i_s_global_info, myrocks::rdb_i_s_ddl,
myrocks::rdb_i_s_sst_props, myrocks::rdb_i_s_index_file_map,
myrocks::rdb_i_s_lock_info, myrocks::rdb_i_s_trx_info,
myrocks::rdb_i_s_deadlock_info
maria_declare_plugin_end;
Messung V0.5 in Prozent
¤ 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.1.968Bemerkung:
(vorverarbeitet am 2026-10-08)
¤
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.