FTS_MSG_ADD_TABLE, /*!< Add table to the optimize thread's
work queue */
FTS_MSG_DEL_TABLE, /*!< Remove a table from the optimize
threads work queue */
FTS_MSG_SYNC_TABLE /*!< Sync fts cache of a table */
};
/** Compressed list of words that have been read from FTS INDEX
that needs to be optimized. */ struct fts_zip_t {
lint status; /*!< Status of (un)/zip operation */
ulint n_words; /*!< Number of words compressed */
ulint block_sz; /*!< Size of a block in bytes */
ib_vector_t* blocks; /*!< Vector of compressed blocks */
ib_alloc_t* heap_alloc; /*!< Heap to use for allocations */
ulint pos; /*!< Offset into blocks */
ulint last_big_block; /*!< Offset of last block in the blocksarraythatisofsize block_sz.Blocksbeyondthisoffset
are of size FTS_MAX_WORD_LEN */
z_streamp zp; /*!< ZLib state */
/*!< The value of the last word read fromtheFTSINDEXtable.Thisis
used to discard duplicates */
fts_string_t word; /*!< UTF-8 string */
ulint max_words; /*!< maximum number of words to read
in one pass */
};
/** Prepared statements used during optimize */ struct fts_optimize_graph_t { /*!< Delete a word from FTS INDEX */
que_t* delete_nodes_graph; /*!< Insert a word into FTS INDEX */
que_t* write_nodes_graph; /*!< COMMIT a transaction */
que_t* commit_graph; /*!< Read the nodes from FTS_INDEX */
que_t* read_nodes_graph;
};
/** Used by fts_optimize() to store state. */ struct fts_optimize_t {
trx_t* trx; /*!< The transaction used for all SQL */
ib_alloc_t* self_heap; /*!< Heap to use for allocations */
char* name_prefix; /*!< FTS table name prefix */
dict_table_t* table; /*!< Table that has to be queried */
dict_index_t* index; /*!< The FTS index to be optimized */
fts_doc_ids_t* to_delete; /*!< doc ids to delete, we check against thisvectorandpurgethematching entriesduringtheoptimizing process.Thevectorentriesare
sorted on doc id */
ulint del_pos; /*!< Offset within to_delete vector, thisisusedtokeeptrackofwhere
we are up to in the vector */
ibool done; /*!< TRUE when optimize finishes */
ib_vector_t* words; /*!< Word + Nodes read from FTS_INDEX,
it contains instances of fts_word_t */
fts_zip_t* zip; /*!< Words read from the FTS_INDEX */
fts_optimize_graph_t /*!< Prepared statements used during */
graph; /*optimize */
ulint n_completed; /*!< Number of FTS indexes that have
been optimized */
ibool del_list_regenerated; /*!< BEING_DELETED list regenerated */
};
/** Used by the optimize, to keep state during compacting nodes. */ struct fts_encode_t {
doc_id_t src_last_doc_id;/*!< Last doc id read from src node */
byte* src_ilist_ptr; /*!< Current ptr within src ilist */
};
/** We use this information to determine when to start the optimize
cycle for a table. */ struct fts_slot_t { /** table, or NULL if the slot is unused */
dict_table_t* table;
/** whether this slot is being processed */ bool running;
ulint added; /*!< Number of doc ids added since the
last time this table was optimized */
ulint deleted; /*!< Number of doc ids deleted since the
last time this table was optimized */
/** time(NULL) of completing fts_optimize_table_bk() */
time_t last_run;
/** time(NULL) of latest successful fts_optimize_table() */
time_t completed;
};
/** A table remove message for the FTS optimize thread. */ struct fts_msg_del_t
{ /** the table to remove */
dict_table_t *table; /** condition variable to signal message consumption */
pthread_cond_t *cond;
};
/** The FTS optimize message work queue message type. */ struct fts_msg_t {
fts_msg_type_t type; /*!< Message type */
void* ptr; /*!< The message contents */
mem_heap_t* heap; /*!< The heap used to allocate this message,themessageconsumerwill
free the heap. */
};
/** The number of words to read and optimize in a single pass. */
ulong fts_num_word_optimize;
/**********************************************************************//**
Create an instance of fts_zip_t.
@return a new instance of fts_zip_t */ static
fts_zip_t*
fts_zip_create( /*===========*/
mem_heap_t* heap, /*!< in: heap */
ulint block_sz, /*!< in: size of a zip block.*/
ulint max_words) /*!< in: max words to read */
{
fts_zip_t* zip;
zip = static_cast<fts_zip_t*>(mem_heap_zalloc(heap, sizeof(*zip)));
/**********************************************************************//**
Read a word */ static
byte*
fts_zip_read_word( /*==============*/
fts_zip_t* zip, /*!< in: Zip state + data */
fts_string_t* word) /*!< out: uncompressed word */
{ short len = 0; void* null = NULL;
byte* ptr = word->f_str; int flush = Z_NO_FLUSH;
/* Either there was an error or we are at the Z_STREAM_END. */ if (zip->status != Z_OK) { return(NULL);
}
case Z_BUF_ERROR: /* No progress possible. */ case Z_STREAM_END:
inflateEnd(zip->zp); break;
case Z_STREAM_ERROR: default:
ut_error;
}
}
/* All blocks must be freed at end of inflate. */ if (zip->status != Z_OK) { for (ulint i = 0; i < ib_vector_size(zip->blocks); ++i) { if (ib_vector_getp(zip->blocks, i)) {
ut_free(ib_vector_getp(zip->blocks, i));
ib_vector_set(zip->blocks, i, &null);
}
}
}
if (ptr != NULL) {
ut_ad(word->f_len == strlen((char*) ptr));
}
/* Skip duplicate words */ if (zip->word.f_len == word_len &&
!memcmp(zip->word.f_str, word_data, word_len)) return DB_SUCCESS;
/* Initialize deflate if not done yet */ if (!compress_inited)
{ int err = deflateInit(zip->zp, 9); if (err != Z_OK)
{
sql_print_error("InnoDB: ZLib deflateInit() failed: %d", err); return DB_ERROR;
}
compress_inited = true;
}
/* Update current word */
memcpy(zip->word.f_str, word_data, word_len);
zip->word.f_len = word_len;
ut_a(zip->zp->avail_in == 0);
ut_a(zip->zp->next_in == NULL);
/* Compress the word with length prefix */
uint16_t len = static_cast<uint16_t>(word_len);
zip->zp->next_in = reinterpret_cast<byte*>(&len);
zip->zp->avail_in = sizeof(len);
/* Compress the word, create output blocks as necessary */ while (zip->zp->avail_in > 0)
{ /* No space left in output buffer, create a new one */ if (zip->zp->avail_out == 0)
{
byte* block= static_cast<byte*>(ut_malloc_nokey(zip->block_sz));
ib_vector_push(zip->blocks, &block);
zip->zp->next_out= block;
zip->zp->avail_out= static_cast<uInt>(zip->block_sz);
}
switch (zip->status = deflate(zip->zp, Z_NO_FLUSH))
{ case Z_OK: if (zip->zp->avail_in == 0)
{
zip->zp->next_in= const_cast<byte*>(word_data);
zip->zp->avail_in = static_cast<uInt>(len);
ut_a(len <= FTS_MAX_WORD_LEN);
len = 0;
} continue; case Z_STREAM_END: case Z_BUF_ERROR: case Z_STREAM_ERROR: default:
ut_error;
}
}
/* All data should have been compressed */
ut_a(zip->zp->avail_in == 0);
zip->zp->next_in = NULL;
++zip->n_words;
/* Continue until we reach max words */ return zip->n_words < zip->max_words ? DB_SUCCESS : DB_SUCCESS_LOCKED_REC;
};
for (uint8_t selected= fts_select_index(cs, word->f_str, word->f_len);
selected < FTS_NUM_AUX_INDEX; selected++)
{ for (;;)
{
AuxRecordReader aux_reader(optim->zip, compress_processor,
AuxCompareMode::GREATER);
fts_zip_t *zip = optim->zip; if (error == DB_SUCCESS && zip->status == Z_OK && zip->n_words > 0)
{ /* All data should have been read */
ut_a(zip->zp->avail_in == 0);
fts_zip_deflate_end(zip);
} else deflateEnd(zip->zp);
return error;
}
dberr_t fts_table_fetch_doc_ids(FTSQueryExecutor *executor, constchar *tbl_name,
fts_doc_ids_t *doc_ids) noexcept
{
ut_ad(executor != nullptr);
executor->trx()->op_info = "fetching FTS doc ids"; /* Append doc_ids straight into the caller's vector rather than
buffering the whole set in the reader first. */
CommonTableReader reader(doc_ids->doc_ids);
dberr_t err= executor->read_all_common(tbl_name, reader);
if (err == DB_SUCCESS)
fts_doc_ids_sort(doc_ids->doc_ids);
return err;
}
/**********************************************************************//** Do a binary search for a doc id in the array
@return +ve index if found -ve index where it should be inserted ifnot found */ int
fts_bsearch( /*========*/
doc_id_t* array, /*!< in: array to sort */ int lower, /*!< in: the array lower bound */ int upper, /*!< in: the array upper bound */
doc_id_t doc_id) /*!< in: the doc id to search for */
{ int orig_size = upper;
if (upper == 0) { /* Nothing to search */ return(-1);
} else { while (lower < upper) { int i = (lower + upper) >> 1;
/**********************************************************************//**
Search in the to delete array whether any of the doc ids within
the [first, last] range are to be deleted
@return +ve index if found -ve index where it should be inserted ifnot found */ static int
fts_optimize_lookup( /*================*/
ib_vector_t* doc_ids, /*!< in: array to search */
ulint lower, /*!< in: lower limit of array */
doc_id_t first_doc_id, /*!< in: doc id to lookup */
doc_id_t last_doc_id) /*!< in: doc id to lookup */
{ int pos; int upper = static_cast<int>(ib_vector_size(doc_ids));
doc_id_t* array = (doc_id_t*) doc_ids->data;
/* If i is 1, it could be first_doc_id is less than eitherthefirstorsecondarrayitem,doa
double check */ if (i == 1 && array[0] <= last_doc_id
&& first_doc_id < array[0]) {
pos = 0;
} elseif (i < upper && array[i] <= last_doc_id) {
/* Check if the "next" doc id is within the
first & last doc id of the node. */
pos = i;
}
}
return(pos);
}
/**********************************************************************//**
Encode the word pos list into the node
@return DB_SUCCESS or error code*/ static MY_ATTRIBUTE((nonnull))
dberr_t
fts_optimize_encode_node( /*=====================*/
fts_node_t* node, /*!< in: node to fill*/
doc_id_t doc_id, /*!< in: doc id to encode */
fts_encode_t* enc) /*!< in: encoding state.*/
{
byte* dst;
ulint enc_len;
ulint pos_enc_len;
doc_id_t doc_id_delta;
dberr_t error = DB_SUCCESS; const byte* src = enc->src_ilist_ptr;
if (node->first_doc_id == 0) {
ut_a(node->last_doc_id == 0);
node->first_doc_id = doc_id;
}
/* Calculate the space required to store the ilist. */
ut_ad(doc_id > node->last_doc_id);
doc_id_delta = doc_id - node->last_doc_id;
enc_len = fts_get_encoded_len(static_cast<ulint>(doc_id_delta));
/* Calculate the size of the encoded pos array. */ while (*src) {
fts_decode_vlc(&src);
}
/* Skip the 0x00 byte at the end of the word positions list. */
++src;
/* Number of encoded pos bytes to copy. */
pos_enc_len = ulint(src - enc->src_ilist_ptr);
/* Total number of bytes required for copy. */
enc_len += pos_enc_len;
/* Check we have enough space in the destination buffer for
copying the document word list. */ if (!node->ilist) {
ulint new_size;
/* While there is data in the source node and space to copy
into in the destination node. */ while (copied < src_node->ilist_size
&& dst_node->ilist_size < FTS_ILIST_MAX_SIZE) {
test_again: /* Check whether the doc id is in the delete list, if sothenweskiptheentriesbutweneedtotrackthe deltafordecodingtheentriesfollowingthisdocument's
entries. */ if (*del_pos >= 0 && *del_pos < (int) ib_vector_size(del_vec)) {
doc_id_t* update;
/* Skip the entries for this document. */ while (*enc->src_ilist_ptr) {
fts_decode_vlc((const byte**)&enc->src_ilist_ptr);
}
/* Skip the end of word position marker. */
++enc->src_ilist_ptr;
} else {
/* DOC ID already becomes larger than
del_doc_id, check the next del_doc_id */ if (del_doc_id > 0 && doc_id > del_doc_id) {
del_doc_id = 0;
++*del_pos;
delta = 0; goto test_again;
}
/* Decode and copy the word positions into
the dest node. */
fts_optimize_encode_node(dst_node, doc_id, enc);
++dst_node->doc_count;
ut_a(dst_node->last_doc_id == doc_id);
}
/* Bytes copied so for from source. */
copied = ulint(enc->src_ilist_ptr - src_node->ilist);
}
if (copied >= src_node->ilist_size) {
ut_a(doc_id == src_node->last_doc_id);
}
enc->src_last_doc_id = doc_id;
return(error);
}
/**********************************************************************//**
Determine the starting pos within the deleted doc id vector for a word.
@returndelete position */ static MY_ATTRIBUTE((nonnull, warn_unused_result)) int
fts_optimize_deleted_pos( /*=====================*/
fts_optimize_t* optim, /*!< in: optimize state data */
fts_word_t* word) /*!< in: the word data to check */
{ int del_pos;
ib_vector_t* del_vec = optim->to_delete->doc_ids;
/* Get the first and last dict ids for the word, we will use thesevaluestodeterminewhichdocidsneedtoberemoved whenwecoalescethenodes.Thiswaywecanreducethenumer ofelementsthatneedtobesearchedinthedeleteddocids vectorandsecondlywecanremovethedocidsduringthe
coalescing phase. */ if (ib_vector_size(del_vec) > 0) {
fts_node_t* node;
doc_id_t last_id;
doc_id_t first_id;
ulint size = ib_vector_size(word->nodes);
del_pos = -1; /* Note that there is nothing to delete. */
}
return(del_pos);
}
/**********************************************************************//**
Compact the nodes for a word, we also remove any doc ids during the
compaction pass.
@return DB_SUCCESS or error code.*/ static
ib_vector_t*
fts_optimize_word( /*==============*/
fts_optimize_t* optim, /*!< in: optimize state data */
fts_word_t* word) /*!< in: the word to optimize */
{
fts_encode_t enc;
ib_vector_t* nodes;
ulint i = 0; int del_pos;
fts_node_t* dst_node = NULL;
ib_vector_t* del_vec = optim->to_delete->doc_ids;
ulint size = ib_vector_size(word->nodes);
err = executor->insert_aux_record(selected, &insert_data); if (err != DB_SUCCESS)
{
sql_print_error("InnoDB: (%s) during optimize, when " "inserting a word to the FTS index.",
ut_strerr(err)); return err;
}
ut_free(node->ilist);
node->ilist= nullptr;
node->ilist_size= node->ilist_size_alloc= 0;
}
/**********************************************************************//**
Compact the nodes for a given word, the nodes passed in are
already optimized.
@return status one of RESTART, EXIT, ERROR */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t
fts_optimize_compact( /*=================*/
FTSQueryExecutor* executor, /*!< in: query executor */
fts_optimize_t* optim, /*!< in: optimize state data */
dict_index_t* index, /*!< in: current FTS being optimized */
time_t start_time) noexcept /*!< in: optimize start time */
{
ulint i;
dberr_t error = DB_SUCCESS;
ulint size = ib_vector_size(optim->words);
for (i = 0; i < size && error == DB_SUCCESS && !optim->done; ++i) {
fts_word_t* word;
ib_vector_t* nodes;
word = (fts_word_t*) ib_vector_get(optim->words, i);
/* nodes is allocated from the word heap and will be destroyed whenthewordisfreed.Wehoweverhavetobecarefulabout
the ilist, that needs to be freed explicitly. */
nodes = fts_optimize_word(optim, word);
/* Update the data on disk. */
error = fts_optimize_write_word(
executor, index, &word->text, nodes);
if (error == DB_SUCCESS) { /* Write the last word optimized to the config table,
we use this value for restarting optimize. */
error = fts_config_set_index_value(
executor, index,
FTS_LAST_OPTIMIZED_WORD, &word->text);
}
/* Free the word that was optimized. */
fts_word_free(word);
/* This will free the heap from which optim itself was allocated. */
mem_heap_free(heap);
}
/** Get the max time optimize should run in millisecs. @paramexecutorqueryexecutor @paramtableusertabletobeoptimized
@return max optimize time limit in millisecs. */ static
ulint fts_optimize_get_time_limit(FTSQueryExecutor *executor, const dict_table_t *table) noexcept
{
ulint time_limit= 0;
fts_string_t value;
value.f_len= FTS_MAX_CONFIG_VALUE_LEN;
value.f_str= static_cast<byte*>(ut_malloc_nokey(value.f_len + 1));
dberr_t error= fts_config_get_value(executor, table,
FTS_OPTIMIZE_LIMIT_IN_SECS, &value); if (error == DB_SUCCESS)
time_limit= strtoul(reinterpret_cast<char*>(value.f_str), nullptr, 10);
ut_free(value.f_str); /* FIXME: This is returning milliseconds, while the variable
is being stored and interpreted as seconds! */ return(time_limit * 1000);
}
/** Run OPTIMIZE on the given table. Note: this can take a very longtime(hours). @paramexecutorqueryexecutor @paramoptimoptimizeinstance @paramindexcurrentftsbeingoptimized
@param word starting word to optimize */ static void fts_optimize_words(FTSQueryExecutor *executor, fts_optimize_t *optim,
dict_index_t *index, fts_string_t *word) noexcept
{
ut_a(!optim->done); /* Get the time limit from the config table. */
fts_optimize_time_limit=
fts_optimize_get_time_limit(executor, index->table); const time_t start_time= time(NULL);
while (!optim->done)
{
trx_t *trx= optim->trx;
ut_a(ib_vector_size(optim->words) == 0); /* Read the index records to optimize. */
dberr_t error= fts_index_fetch_nodes(
executor, index, word, optim->words, nullptr, AuxCompareMode::EQUAL); if (error == DB_SUCCESS)
{ /* There must be some nodes to read. */
ut_a(ib_vector_size(optim->words) > 0); /* Optimize the nodes that were read and write back to DB. */
error = fts_optimize_compact(executor, optim, index, start_time); if (error == DB_SUCCESS) fts_sql_commit(optim->trx); else fts_sql_rollback(optim->trx);
}
ib_vector_reset(optim->words);
/** Optimize is complete. Set the completion time, and reset the optimizestartstringforthisFTSindexto"". @paramexecutorqueryexecutor @paramoptimoptimizeinstance @paramindextablewithoneFTSindex
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t
fts_optimize_index_completed(FTSQueryExecutor *executor,
fts_optimize_t *optim,
dict_index_t *index) noexcept
{
fts_string_t word;
dberr_t error= DB_SUCCESS;
byte buf[sizeof(ulint)]; /* If we've reached the end of the index then set the start
word to the empty string. */
word.f_len= 0;
word.f_str= buf;
*word.f_str= '\0';
if (UNIV_UNLIKELY(error != DB_SUCCESS))
sql_print_error("InnoDB: (%s) while updating last optimized word!",
ut_strerr(error)); return error;
}
/** Read the words that will be optimized in this pass. @paramexecutorqueryexecutor @paramoptimoptimizeinstance @paramindextablewithoneFTSindex @paramwordbuffertouse
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t
fts_optimize_index_read_words(
FTSQueryExecutor* executor,
fts_optimize_t* optim,
dict_index_t* index,
fts_string_t* word)
{
dberr_t error = DB_SUCCESS;
if (optim->del_list_regenerated) {
word->f_len = 0;
} else { /* Get the last word that was optimized from
the config table. */
error = fts_config_get_index_value(
executor, index, FTS_LAST_OPTIMIZED_WORD, word);
}
/* If record not found then we start from the top. */ if (error == DB_RECORD_NOT_FOUND) {
word->f_len = 0;
error = DB_SUCCESS;
}
optim->index = index; while (error == DB_SUCCESS) {
error = fts_index_fetch_words(
executor, optim, word, fts_num_word_optimize);
if (error == DB_SUCCESS) { /* Reset the last optimized word to '' if no
more words could be read from the FTS index. */ if (optim->zip->n_words == 0) {
word->f_len = 0;
*word->f_str = 0;
}
break;
}
}
return(error);
}
/**********************************************************************//**
Run OPTIMIZE on the given FTS index. Note: this can take a very long
time (hours).
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t
fts_optimize_index( /*===============*/
FTSQueryExecutor* executor, /*!< in: query executor */
fts_optimize_t* optim, /*!< in: optimize instance */
dict_index_t* index) /*!< in: table with one FTS index */
{
fts_string_t word;
dberr_t error;
byte str[FTS_MAX_WORD_LEN + 1];
optim->done = FALSE; /* Optimize until !done */
/* We need to read the last word optimized so that we start from
the next word. */
word.f_str = str;
/* We set the length of word to the size of str since we
need to pass the max len info to the fts_get_config_value() function. */
word.f_len = sizeof(str) - 1;
memset(word.f_str, 0x0, word.f_len);
/* Read the words that will be optimized in this pass. */
error = fts_optimize_index_read_words(executor, optim, index, &word);
/* If we couldn't read any records then optimize is complete.Incrementthenumberofindexesthathave beenoptimizedandsetFTSindexoptimizestateto
completed. */ if (error == DB_SUCCESS && optim->zip->n_words == 0) {
if (error == DB_SUCCESS) {
++optim->n_completed;
}
}
}
return(error);
}
/** Purge the doc ids that are in the snapshot from themasterdeletedtable. @paramexecutorqueryexecutor @paramoptimoptimizeinstance
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t fts_optimize_purge_deleted_doc_ids(FTSQueryExecutor *executor,
fts_optimize_t *optim) noexcept
{
dberr_t error= DB_SUCCESS;
ut_a(ib_vector_size(optim->to_delete->doc_ids) > 0); for (ulint i= 0;
i < ib_vector_size(optim->to_delete->doc_ids) && error != DB_SUCCESS;
++i)
{
doc_id_t *update= static_cast<doc_id_t*>(ib_vector_get(optim->to_delete->doc_ids, i));
error= executor->delete_common_record("DELETED", *update); if (error == DB_SUCCESS)
error= executor->delete_common_record("DELETED_CACHE", *update);
}
if (error != DB_SUCCESS)
fts_sql_rollback(optim->trx); return error;
}
/** Delete the document ids in the pending delete, and delete tables. @paramexecutorqueryexecutor @paramoptimoptimizeinstance
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t fts_optimize_purge_deleted_doc_id_snapshot(FTSQueryExecutor *executor,
fts_optimize_t *optim) noexcept
{
dberr_t error= executor->delete_all_common_records("BEING_DELETED"); if (error == DB_SUCCESS)
error= executor->delete_all_common_records("BEING_DELETED_CACHE"); return error;
}
/** Check if there are records in BEING_DELETED table @paramexecutorqueryexecutor @paramoptimoptimizeftsinstance @paramn_rowsnumberofrowsexistinbeing_deletedtable
@return DB_SUCCESS if all OK */ static
dberr_t fts_optimize_being_deleted_count(FTSQueryExecutor *executor,
fts_optimize_t *optim,
ulint *n_rows) noexcept
{
CommonTableReader reader;
dberr_t err= executor->read_all_common("BEING_DELETED", reader); if (err == DB_SUCCESS) *n_rows= reader.size(); return err;
}
/** Create a snapshot of deleted document IDs by moving them from DELETEDtoBEING_DELETEDandfromDELETED_CACHEto BEING_DELETED_CACHE. @paramexecutorqueryexecutor @paramoptimoptimizeftsinstance
@return DB_SUCCESS or error code */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t fts_optimize_create_deleted_doc_id_snapshot(FTSQueryExecutor *executor,
fts_optimize_t *optim) noexcept
{
dberr_t err= DB_SUCCESS;
CommonTableReader reader;
for (ulint i= 0, n= reader.size(); i < n; i++)
{
err= executor->insert_common_record("BEING_DELETED_CACHE", reader.get(i)); if (err != DB_SUCCESS) return err;
}
optim->del_list_regenerated= TRUE; return err;
}
/*********************************************************************//**
Read in the document ids that are to be purged during optimize. The
transaction is committed upon successfully read.
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t
fts_optimize_read_deleted_doc_id_snapshot( /*======================================*/
FTSQueryExecutor* executor, /*!< in: FTS query executor */
fts_optimize_t* optim) noexcept /*!< in: optimize instance */
{ /* Read the doc_ids to delete. */
dberr_t error = fts_table_fetch_doc_ids(
executor, "BEING_DELETED", optim->to_delete);
/*********************************************************************//**
Cleanup the snapshot tables and the master deleted table.
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull, warn_unused_result))
dberr_t
fts_optimize_purge_snapshot( /*========================*/
FTSQueryExecutor* executor, /*!< in: query executor */
fts_optimize_t* optim) noexcept /*!< in: optimize instance */
{
dberr_t error;
/* Delete the doc ids from the master deleted tables, that were
in the snapshot that was taken at the start of optimize. */
error = fts_optimize_purge_deleted_doc_ids(executor, optim);
if (error == DB_SUCCESS) { /* Destroy the deleted doc id snapshot. */
error = fts_optimize_purge_deleted_doc_id_snapshot(
executor, optim);
}
/*********************************************************************//**
Run OPTIMIZE on the given table by a background thread.
@return DB_SUCCESS if all OK */ static MY_ATTRIBUTE((nonnull))
dberr_t
fts_optimize_table_bk( /*==================*/
fts_slot_t* slot) /*!< in: table to optimiza */
{ const time_t now = time(NULL); const ulint interval = ulint(now - slot->last_run);
/* Avoid optimizing tables that were optimized recently. */ if (slot->last_run > 0
&& lint(interval) >= 0
&& interval < FTS_OPTIMIZE_INTERVAL_IN_SECS) {
if (error == DB_SUCCESS) {
slot->running = false;
slot->completed = slot->last_run;
}
} else { /* Note time this run completed. */
slot->last_run = now;
error = DB_SUCCESS;
}
return(error);
} /** Run OPTIMIZE on the given table. @paramtabletabletobeoptimized @paramthdthreadwhichexecutesoptimizetable
@return DB_SUCCESS if all OK */
dberr_t
fts_optimize_table(dict_table_t *table, THD *thd)
{
ut_ad(!srv_read_only_mode || recv_sys.rpo);
if (recv_sys.rpo) { return DB_READ_ONLY;
}
/* Serialize concurrent fts_optimize_table() on the same table: acquireMDL_EXCLUSIVEfirstsoasecondcallerblockshere,then downgradetoMDL_SHARED_UPGRADABLEsootheroperationscanproceed
while optimization is in progress. */
MDL_ticket* mdl_ticket = nullptr;
dict_sys.freeze(SRW_LOCK_CALL); if (dict_acquire_mdl<false, true>(table, thd, &mdl_ticket)
!= table) {
dict_sys.unfreeze(); if (mdl_ticket) {
thd->mdl_context.release_lock(mdl_ticket);
} return DB_TABLE_NOT_FOUND;
}
dict_sys.unfreeze();
if (mdl_ticket) {
mdl_ticket->downgrade_lock(MDL_SHARED_UPGRADABLE);
}
/* Create FTSQueryExecutor and open common tables */
FTSQueryExecutor executor(optim->trx, table);
error = executor.open_all_deletion_tables(); if (error != DB_SUCCESS) {
err_exit:
fts_optimize_free(optim); if (mdl_ticket) {
thd->mdl_context.release_lock(mdl_ticket);
} return error;
}
error = executor.open_config_table(); if (error) { goto err_exit; }
// FIXME: Call this only at the start of optimize, currently we // rely on DB_DUPLICATE_KEY to handle corrupting the snapshot.
/* Check whether there are still records in BEING_DELETED table */
ulint n_rows = 0;
error= fts_optimize_being_deleted_count(&executor, optim, &n_rows);
if (error == DB_SUCCESS && n_rows == 0) { /* Take a snapshot of the deleted document ids, they are copied
to the BEING_ tables. */
error = fts_optimize_create_deleted_doc_id_snapshot(
&executor, optim);
}
/* A duplicate error is OK, since we don't erase the docidsfromthebeingdeletedstateuntilallFTS
indexes have been optimized. */ if (error == DB_DUPLICATE_KEY) {
error = DB_SUCCESS;
}
if (error == DB_SUCCESS) {
/* These document ids will be filtered out during the indexoptimizationphase.Theyareinthesnapshotthatwe
took above, at the start of the optimize. */
error = fts_optimize_read_deleted_doc_id_snapshot(&executor, optim);
if (error == DB_SUCCESS) {
/* Commit the read of being deleted
doc ids transaction. */
fts_sql_commit(optim->trx);
/* We would do optimization only if there
are deleted records to be cleaned up */ if (ib_vector_size(optim->to_delete->doc_ids) > 0) {
error = fts_optimize_indexes(&executor, optim);
}
} else {
ut_a(optim->to_delete == NULL);
}
/* Only after all indexes have been optimized can we deletethe(snapshot)docidsinthependingdelete,
and master deleted tables. */ if (error == DB_SUCCESS
&& optim->n_completed == ib_vector_size(fts->indexes)) {
if (ib_vector_size(optim->to_delete->doc_ids) > 0) {
/* Purge the doc ids that were in the snapshotfromthesnapshottablesand
the master deleted table. */
error = fts_optimize_purge_snapshot(
&executor, optim);
}
}
}
fts_optimize_free(optim);
if (mdl_ticket) {
thd->mdl_context.release_lock(mdl_ticket);
}
return(error);
}
/********************************************************************//**
Add the table to add to the OPTIMIZER's list.
@returnnew message instance */ static
fts_msg_t*
fts_optimize_create_msg( /*====================*/
fts_msg_type_t type, /*!< in: type of message */ void* ptr) /*!< in: message payload */
{
mem_heap_t* heap;
fts_msg_t* msg;
/**********************************************************************//**
Remove the table from the OPTIMIZER's list. We do wait for
acknowledgement from the consumer of the message. */ void
fts_optimize_remove_table( /*======================*/
dict_table_t* table) /*!< in: table to remove */
{ if (!fts_optimize_wq) return;
if (fts_opt_start_shutdown)
{
sql_print_information("InnoDB: Try to remove table %s after FTS optimize " "thread exiting.", table->name.m_name); while (fts_optimize_wq)
std::this_thread::sleep_for(std::chrono::milliseconds(10)); return;
}
/** Send sync fts cache for the table.
@param[in] table table to sync */ void
fts_optimize_request_sync_table(
dict_table_t* table)
{ /* if the optimize system not yet initialized, return */ if (!fts_optimize_wq) { return;
}
mysql_mutex_lock(&fts_optimize_wq->mutex);
/* FTS optimizer thread is already exited */ if (fts_opt_start_shutdown) {
sql_print_information("InnoDB: Try to sync table %s " "after FTS optimize thread exiting.",
table->name.m_name);
} elseif (table->fts->sync_message) { /* If the table already has SYNC message in
fts_optimize_wq queue then ignore it */
} else {
add_msg(fts_optimize_create_msg(FTS_MSG_SYNC_TABLE, table));
table->fts->sync_message = true;
DBUG_EXECUTE_IF("fts_optimize_wq_count_check",
DBUG_ASSERT(fts_optimize_wq->length <= 1000););
}
mysql_mutex_unlock(&fts_optimize_wq->mutex);
}
/** Add a table to fts_slots if it doesn't already exist. */ staticbool fts_optimize_new_table(dict_table_t* table)
{
ut_ad(table);
/** Remove a table from fts_slots if it exists.
@param remove table to be removed from fts_slots */ staticbool fts_optimize_del_table(fts_msg_del_t *remove)
{ const dict_table_t* table = remove->table;
ut_ad(table); for (ulint i = 0; i < ib_vector_size(fts_slots); ++i) {
fts_slot_t* slot;
/**********************************************************************//**
Calculate how many tables in fts_slots need to be optimized.
@return no. of tables to optimize */ static ulint fts_optimize_how_many()
{
ulint n_tables = 0; const time_t current_time = time(NULL);
for (ulint i = 0; i < ib_vector_size(fts_slots); ++i) { const fts_slot_t* slot = static_cast<const fts_slot_t*>(
ib_vector_get_const(fts_slots, i)); if (!slot->table) { continue;
}
/**********************************************************************//**
Check if the total memory used by all FTS table exceeds the maximum limit.
@returntrueif a sync is needed, false otherwise */ staticbool fts_is_sync_needed()
{
ulint total_memory = 0; const time_t now = time(NULL); double time_diff = difftime(now, last_check_sync_time);
while (!done && srv_shutdown_state <= SRV_SHUTDOWN_INITIATED) { #ifdef WITH_WSREP
ut_d(extern Atomic_relaxed<bool> wsrep_sst_disable_writes);
ut_ad(!wsrep_sst_disable_writes); #endif /* If there is no message in the queue and we have tables
to optimize then optimize the tables. */
/* Server is being shutdown, sync the data from FTS cache to disk
if needed */ if (n_tables > 0) { for (ulint i = 0; i < ib_vector_size(fts_slots); i++) {
fts_slot_t* slot = static_cast<fts_slot_t*>(
ib_vector_get(fts_slots, i));
if (slot->table) {
fts_optimize_sync_table(slot->table);
}
}
}
fts_opt_thd = innobase_create_background_thd("InnoDB FTS optimizer"); /* Add fts tables to fts_slots which could be skipped duringdict_load_table_one()becausefts_optimize_thread
wasn't even started. */
dict_sys.freeze(SRW_LOCK_CALL); for (dict_table_t* table = UT_LIST_GET_FIRST(dict_sys.table_LRU);
table != NULL;
table = UT_LIST_GET_NEXT(table_LRU, table)) { if (!table->fts || !dict_table_has_fts_index(table)) { continue;
}
/* fts_optimize_thread is not started yet. So there is no needtoacquirefts_optimize_wq->mutexforaddingthefts
table to the fts slots. */
ut_ad(!table->can_be_evicted);
fts_optimize_new_table(table);
table->fts->in_queue = true;
}
dict_sys.unfreeze();
/** Shut down the fts optimize thread. */ void fts_optimize_shutdown()
{
ut_ad(!srv_read_only_mode);
ut_ad(!recv_sys.rpo);
/* If there is an ongoing activity on dictionary, such as
srv_master_evict_from_table_cache(), wait for it */
dict_sys.freeze(SRW_LOCK_CALL);
mysql_mutex_lock(&fts_optimize_wq->mutex); /* Tells FTS optimizer system that we are exiting from optimizerthread,messagessentthereafterwillnotbe
processed */
fts_opt_start_shutdown = true;
dict_sys.unfreeze();
/* We tell the OPTIMIZE thread to switch to state done, we can'tdeletetheworkqueueherebecausetheaddthreadneeds
deregister the FTS tables. */
timer->disarm();
task_group.cancel_pending(&task);
#ifdef WITH_WSREP /** Pause the optimize subsystem. */ void fts_optimize_pause()
{
ut_ad(!srv_read_only_mode); /* Prevent fts_optimize_callback() from being scheduled. */
timer->disarm(); /* Wait for any current fts_optimize_callback() to finish. */
task.wait();
}
/** Resume after fts_optimize_stop() */ void fts_optimize_resume()
{ /* Schedule fts_optimize_callback() immediately.
It will reschedule itself via the timer when needed. */
srv_thread_pool->submit_task(&task);
} #endif
/** Sync the table during commit phase
@param[in] table table to be synced */ void fts_sync_during_ddl(dict_table_t* table, THD* thd)
{ if (!fts_optimize_wq) return;
mysql_mutex_lock(&fts_optimize_wq->mutex); constauto sync_message= table->fts->sync_message;
mysql_mutex_unlock(&fts_optimize_wq->mutex); if (!sync_message) return;
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.