/** Read the next record to buffer N.
@param N index into array of merge info structure */ #define ROW_MERGE_READ_GET_NEXT(N) \ do { \
b[N] = row_merge_read_rec( \
block[N], buf[N], b[N], index, \
fd[N], &foffs[N], &mrec[N], offsets[N], \
crypt_block[N], space); \ if (UNIV_UNLIKELY(!b[N])) { \ if (mrec[N]) { \ gotoexit; \
} \
} \
} while (0)
/*********************************************************************//**
Create a temporary "fts sort index" used to merge sort the
tokenized doc string. The index has three "fields":
1) Tokenized word, 2) Doc ID (depend on number of records to sort, it can be a 4 bytes or8 bytes
integer value) 3) Word's position in original doc.
@see fts_create_one_index_table()
@return dict_index_t structure for the fts sort index */
dict_index_t*
row_merge_create_fts_sort_index( /*============================*/
dict_index_t* index, /*!< in: Original FTS index basedonwhichthissortindex
is created */
dict_table_t* table, /*!< in,out: table that FTS index
is being created on */ bool opt_doc_id_size) /*!< in: whether to use 4 bytes insteadof8bytesintegerto
store Doc ID during sort */
{
dict_index_t* new_index;
dict_field_t* field;
dict_field_t* idx_field;
CHARSET_INFO* charset;
// FIXME: This name shouldn't be hard coded here.
new_index = dict_mem_index_create(table, "tmp_fts_idx", DICT_FTS, 3);
/* The third field is on the word's position in the original doc */
field = dict_index_get_nth_field(new_index, 2);
field->name = NULL;
field->prefix_len = 0;
field->descending = false;
field->col = static_cast<dict_col_t*>(
mem_heap_zalloc(new_index->heap, sizeof(dict_col_t)));
field->col->mtype = DATA_INT;
field->col->len = 4 ;
field->fixed_len = 4;
field->col->prtype = DATA_NOT_NULL;
ut_ad(trx->mysql_thd != NULL); constchar* path = thd_innodb_tmpdir(trx->mysql_thd); /* There will be FTS_NUM_AUX_INDEX number of "sort buckets" for eachparallelsortthread.Each"sortbucket"holdsrecordsfor
a particular "FTS index partition" */ for (j = 0; j < fts_sort_pll_degree; j++) {
if (row_merge_file_create(psort_info[j].merge_file[i],
path) == OS_FILE_CLOSED) { goto func_exit;
}
/* Need to align memory for O_DIRECT write */
psort_info[j].merge_block[i] = static_cast<row_merge_block_t*>(
aligned_malloc(block_size, 1024));
if (!psort_info[j].merge_block[i]) {
ret = FALSE; goto func_exit;
}
/* If tablespace is encrypted, allocate additional buffer for
encryption/decryption. */ if (srv_encrypt_log) { /* Need to align memory for O_DIRECT write */
psort_info[j].crypt_block[i] = static_cast<row_merge_block_t*>(
aligned_malloc(block_size, 1024));
if (!psort_info[j].crypt_block[i]) {
ret = FALSE; goto func_exit;
}
} else {
psort_info[j].crypt_block[i] = NULL;
}
}
func_exit: if (!ret) {
row_fts_psort_info_destroy(psort_info, merge_info);
}
return(ret);
} /*********************************************************************//**
Clean up and deallocate FTS parallel sort structures, and close the
merge sort files */ void
row_fts_psort_info_destroy( /*=======================*/
fts_psort_t* psort_info, /*!< parallel sort info */
fts_psort_t* merge_info) /*!< parallel merge info */
{
ulint i;
ulint j;
if (psort_info) { for (j = 0; j < fts_sort_pll_degree; j++) { for (i = 0; i < FTS_NUM_AUX_INDEX; i++) { if (psort_info[j].merge_file[i]) {
row_merge_file_destroy(
psort_info[j].merge_file[i]);
}
/* Tokenize the data and add each word string, its corresponding
doc id and position to sort buffer */ while (parser
? (!t_ctx->processed_len
|| UT_LIST_GET_LEN(t_ctx->fts_token_list))
: t_ctx->processed_len < doc->text.f_len) {
ulint idx = 0;
ulint cur_len;
doc_id_t write_doc_id;
row_fts_token_t* fts_token = NULL;
if (parser != NULL) { if (t_ctx->processed_len == 0) {
UT_LIST_INIT(t_ctx->fts_token_list, &row_fts_token_t::token_list);
/* Parse the whole doc and cache tokens */
row_merge_fts_doc_tokenize_by_parser(doc,
parser, t_ctx);
/* Just indicate that we have parsed all words */
t_ctx->processed_len += 1;
}
/* Then get a token */
fts_token = UT_LIST_GET_FIRST(t_ctx->fts_token_list); if (fts_token) {
str.f_len = fts_token->text->f_len;
str.f_n_char = fts_token->text->f_n_char;
str.f_str = fts_token->text->f_str;
} else {
ut_ad(UT_LIST_GET_LEN(t_ctx->fts_token_list) == 0); /* Reach the end of the list */
t_ctx->processed_len = doc->text.f_len; break;
}
} else {
inc = innobase_mysql_fts_get_token(
doc->charset,
doc->text.f_str + t_ctx->processed_len,
doc->text.f_str + doc->text.f_len, &str);
ut_a(inc > 0);
}
/* Ignore string whose character number is less than
"fts_min_token_size" or more than "fts_max_token_size" */ if (!fts_check_token(&str, NULL, NULL)) { if (parser != NULL) {
UT_LIST_REMOVE(t_ctx->fts_token_list, fts_token);
ut_free(fts_token);
} else {
t_ctx->processed_len += inc;
}
/* if "cached_stopword" is defined, ignore words in the
stopword list */ if (!fts_check_token(&str, t_ctx->cached_stopword,
doc->charset)) { if (parser != NULL) {
UT_LIST_REMOVE(t_ctx->fts_token_list, fts_token);
ut_free(fts_token);
} else {
t_ctx->processed_len += inc;
}
continue;
}
/* There are FTS_NUM_AUX_INDEX auxiliary tables, find
out which sort buffer to put this word record in */
t_ctx->buf_used = fts_select_index(doc->charset,
(const byte*) str_buf.ptr(), str_buf.length());
/* For the temporary file, row_merge_buf_encode() uses 1byteforrepresentingthenumberofextra_sizebytes. Thisnumberwillalwaysbe1,becauseforthis3-fieldindex consistingofonevariable-sizecolumn,extra_sizewillalways be1or2,whichcanbeencodedinonebyte.
Theextra_sizeis1byteifthelengthofthe variable-lengthcolumnislessthan128bytesorthe
maximum length is less than 256 bytes. */
/* One variable length column, word with its length less than fts_max_token_size,addoneextrasizeandoneextrabyte.
SincethemaxlengthforFTStokennowislargerthan255, sowewillneedtosignifylengthbyteitself,soonly1to128
bytes can be used for 1 bytes, larger than that 2 bytes. */ if (len < 128 || field->type.len < 256) { /* Extra size is one byte. */
cur_len = 2 + len;
} else { /* Extra size is two bytes. */
cur_len = 3 + len;
}
dfield_dup(field, buf->heap);
field++;
/* The second field is the Doc ID */
ib_uint32_t doc_id_32_bit;
if (!opt_doc_id_size) {
fts_write_doc_id((byte*) &write_doc_id, doc_id);
/* The third field is the position. MySQL5.7changedthefulltextparserplugininterface byaddingMYSQL_FTPARSER_BOOLEAN_INFO::position.
Below we assume that the field is always 0. */
ulint pos = t_ctx->init_pos;
byte position[4]; if (parser == NULL) {
pos += t_ctx->processed_len + inc - str.f_len;
}
len = 4;
mach_write_to_4(position, pos);
dfield_set_data(field, &position, len);
/* Reserve one byte for the end marker of row_merge_block_t */ if (buf->total_size + data_size[idx] + cur_len
>= srv_sort_buf_size - 1) {
buf_full = TRUE; break;
}
/* Increment the number of tuples */
n_tuple[idx]++; if (parser != NULL) {
UT_LIST_REMOVE(t_ctx->fts_token_list, fts_token);
ut_free(fts_token);
} else {
t_ctx->processed_len += inc;
}
data_size[idx] += cur_len;
}
/* Update the data length and the number of new word tuples
added in this round of tokenization */ for (ulint i = 0; i < FTS_NUM_AUX_INDEX; i++) { /* The computation of total_size below assumes that no delete-markflagswillbestoredandthatallfields
are NOT NULL and fixed-length. */
/* If finish processing the last item, update "doc" with stringsinthedoc_item,otherwisecontinueprocessinglast
item */ if (processed) {
byte* data;
ulint data_len;
dfield = doc_item->field;
data = static_cast<byte*>(dfield_get_data(dfield));
data_len = dfield_get_len(dfield);
/* Current sort buffer full, need to recycle */ if (!processed) {
ut_ad(buf[0]->index->parser
|| t_ctx.processed_len < doc.text.f_len);
ut_ad(t_ctx.rows_added[t_ctx.buf_used]); break;
}
/* If we run out of current sort buffer, need to sort
and flush the sort buffer to disk */ if (t_ctx.rows_added[t_ctx.buf_used] && !processed) {
row_merge_buf_sort(buf[t_ctx.buf_used], NULL);
row_merge_buf_write(buf[t_ctx.buf_used], #ifndef DBUG_OFF
merge_file[t_ctx.buf_used], #endif
block[t_ctx.buf_used]);
/* Parent done scanning, and if finish processing all the docs, exit */ if (psort_info->state == FTS_PARENT_COMPLETE) { if (UT_LIST_GET_LEN(psort_info->fts_doc_list) == 0) { gotoexit;
} elseif (retried > 10000) {
ut_ad(!doc_item); /* retried too many times and cannot get new record */
ib::error() << "FTS parallel sort processed "
<< num_doc_processed
<< " records, the sort queue has "
<< UT_LIST_GET_LEN(psort_info->fts_doc_list)
<< " records. But sort cannot get the next" " records during alter table " << table->name; gotoexit;
}
} elseif (psort_info->state == FTS_PARENT_EXITING) { /* Parent abort */ goto func_exit;
}
if (doc_item == NULL) {
std::this_thread::yield();
}
exit: /* Do a final sort of the last (or latest) batch of records inblockmemory.Flushthemtotempfileifrecordscannot
be hold in one block memory */ for (i = 0; i < FTS_NUM_AUX_INDEX; i++) { if (t_ctx.rows_added[i]) {
row_merge_buf_sort(buf[i], NULL);
row_merge_buf_write(buf[i], #ifndef DBUG_OFF
merge_file[i], #endif
block[i]);
/* Write to temp file, only if records have beenflushedtotempfilebefore(offset>0): Thepseudocodeforsortisfollowing:
if (UT_LIST_GET_LEN(psort_info->fts_doc_list) > 0) { /* child can exit either with error or told by parent. */
ut_ad(error != DB_SUCCESS
|| psort_info->state == FTS_PARENT_EXITING);
}
/* Free fts doc list in case of error. */ do {
row_merge_fts_get_next_doc_item(psort_info, &doc_item);
} while (doc_item != NULL);
/*********************************************************************//**
Start the parallel tokenization and parallel merge sort */ void
row_fts_start_psort( /*================*/
fts_psort_t* psort_info) /*!< parallel sort structure */
{
ulint i = 0;
for (i = 0; i < fts_sort_pll_degree; i++) {
psort_info[i].psort_id = i;
psort_info[i].task = new tpool::waitable_task(fts_parallel_tokenization,&psort_info[i]);
srv_thread_pool->submit_task(psort_info[i].task);
}
}
/*********************************************************************//**
Function performs the merge and insertion of the sorted records. */ static void
fts_parallel_merge( /*===============*/ void* arg) /*!< in: parallel merge info */
{
fts_psort_t* psort_info = (fts_psort_t*) arg;
ulint id;
/* The first field is the tokenized word */
field = dtuple_get_nth_field(tuple, 0);
dfield_set_data(field, word->f_str, word->f_len);
/* The second field is first_doc_id */
field = dtuple_get_nth_field(tuple, 1);
fts_write_doc_id((byte*)&write_first_doc_id, node->first_doc_id);
dfield_set_data(field, &write_first_doc_id, sizeof(doc_id_t));
/* The third and fourth fileds(TRX_ID, ROLL_PTR) are filled already.*/ /* The fifth field is last_doc_id */
field = dtuple_get_nth_field(tuple, 4);
fts_write_doc_id((byte*)&write_last_doc_id, node->last_doc_id);
dfield_set_data(field, &write_last_doc_id, sizeof(doc_id_t));
/* The sixth field is doc_count */
field = dtuple_get_nth_field(tuple, 5);
mach_write_to_4((byte*)&write_doc_count, (ib_uint32_t)node->doc_count);
dfield_set_data(field, &write_doc_count, sizeof(ib_uint32_t));
/* The seventh field is ilist */
field = dtuple_get_nth_field(tuple, 6);
dfield_set_data(field, node->ilist, node->ilist_size);
ret = ins_ctx->btr_bulk->insert(tuple);
return(ret);
}
/********************************************************************//**
Insert processed FTS data to auxillary index tables.
@return DB_SUCCESS if insertion runs fine */ static MY_ATTRIBUTE((nonnull))
dberr_t
row_merge_write_fts_word( /*=====================*/
fts_psort_insert_t* ins_ctx, /*!< in: insert context */
fts_tokenizer_word_t* word) /*!< in: sorted and tokenized
word */
{
dberr_t ret = DB_SUCCESS;
/* Pop out each fts_node in word->nodes write them to auxiliary table */ for (ulint i = 0; i < ib_vector_size(word->nodes); i++) {
dberr_t error;
fts_node_t* fts_node;
if (UNIV_UNLIKELY(error != DB_SUCCESS)) {
ib::error() << "Failed to write word to FTS auxiliary" " index table "
<< ins_ctx->btr_bulk->table_name()
<< ", error " << error;
ret = error;
}
/*********************************************************************//**
Read sorted FTS data files and insert data tuples to auxillary tables.
@return DB_SUCCESS or error number */ static void
row_fts_insert_tuple( /*=================*/
fts_psort_insert_t*
ins_ctx, /*!< in: insert context */
fts_tokenizer_word_t* word, /*!< in: last processed
tokenized word */
ib_vector_t* positions, /*!< in: word position */
doc_id_t* in_doc_id, /*!< in: last item doc id */
dtuple_t* dtuple) /*!< in: entry to insert */
{
fts_node_t* fts_node = NULL;
dfield_t* dfield;
doc_id_t doc_id;
ulint position;
fts_string_t token_word;
ulint i;
/* Get fts_node for the FTS auxillary INDEX table */ if (ib_vector_size(word->nodes) > 0) {
fts_node = static_cast<fts_node_t*>(
ib_vector_last(word->nodes));
}
if (fts_node == NULL
|| fts_node->ilist_size > FTS_ILIST_MAX_SIZE) {
/* If dtuple == NULL, this is the last word to be processed */ if (!dtuple) { if (fts_node && ib_vector_size(positions) > 0) {
fts_cache_node_add_positions(
NULL, fts_node, *in_doc_id,
positions);
/* Write out the current word */
row_merge_write_fts_word(ins_ctx, word);
}
return;
}
/* Get the first field for the tokenized word */
dfield = dtuple_get_nth_field(dtuple, 0);
if (!word->text.f_str) {
fts_string_dup(&word->text, &token_word, ins_ctx->heap);
}
/* compare to the last word, to see if they are the same
word */ if (innobase_fts_text_cmp(ins_ctx->charset,
&word->text, &token_word) != 0) {
ulint num_item;
/* Getting a new word, flush the last position info
for the current word in fts_node */ if (ib_vector_size(positions) > 0) {
fts_cache_node_add_positions(
NULL, fts_node, *in_doc_id, positions);
}
/* Write out the current word */
row_merge_write_fts_word(ins_ctx, word);
/* Copy the new word */
fts_string_dup(&word->text, &token_word, ins_ctx->heap);
num_item = ib_vector_size(positions);
/* Clean up position queue */ for (i = 0; i < num_item; i++) {
ib_vector_pop(positions);
}
/* Get the word's position info */
dfield = dtuple_get_nth_field(dtuple, 2);
position = mach_read_from_4(static_cast<byte*>(dfield_get_data(dfield)));
/* If this is the same word as the last word, and they havethesameDocID,wejustneedtoadditsposition info.Otherwise,wewillflushpositioninfotothe
fts_node and initiate a new position vector */ if (!(*in_doc_id) || *in_doc_id == doc_id) {
ib_vector_push(positions, &position);
} else {
ulint num_pos = ib_vector_size(positions);
fts_cache_node_add_positions(NULL, fts_node,
*in_doc_id, positions); for (i = 0; i < num_pos; i++) {
ib_vector_pop(positions);
}
ib_vector_push(positions, &position);
}
/* record the current Doc ID */
*in_doc_id = doc_id;
}
/*********************************************************************//**
Propagate a newly added record up one level in the selection tree
@return parent where this value propagated to */ static
ulint
row_fts_sel_tree_propagate( /*=======================*/
ulint propogated, /*<! in: tree node propagated */ int* sel_tree, /*<! in: selection tree */ const mrec_t** mrec, /*<! in: sort record */
rec_offs** offsets, /*<! in: record offsets */
dict_index_t* index) /*<! in/out: FTS index */
{
ulint parent; int child_left; int child_right; int selected;
/* Find which parent this value will be propagated to */
parent = (propogated - 1) / 2;
/* Find out which value is smaller, and to propagate */
child_left = sel_tree[parent * 2 + 1];
child_right = sel_tree[parent * 2 + 2];
/*********************************************************************//**
Readjust selection tree after popping the root and read a new value
@return the new root */ static int
row_fts_sel_tree_update( /*====================*/ int* sel_tree, /*<! in/out: selection tree */
ulint propagated, /*<! in: node to propagate up */
ulint height, /*<! in: tree height */ const mrec_t** mrec, /*<! in: sort record */
rec_offs** offsets, /*<! in: record offsets */
dict_index_t* index) /*<! in: index dictionary */
{
ulint i;
for (i = 1; i <= height; i++) {
propagated = row_fts_sel_tree_propagate(
propagated, sel_tree, mrec, offsets, index);
}
return(sel_tree[0]);
}
/*********************************************************************//**
Build selection tree at a specified level */ static void
row_fts_build_sel_tree_level( /*=========================*/ int* sel_tree, /*<! in/out: selection tree */
ulint level, /*<! in: selection tree level */ const mrec_t** mrec, /*<! in: sort record */
rec_offs** offsets, /*<! in: record offsets */
dict_index_t* index) /*<! in: index dictionary */
{
ulint start; int child_left; int child_right;
ulint i;
ulint num_item = ulint(1) << level;
start = num_item - 1;
for (i = 0; i < num_item; i++) {
child_left = sel_tree[(start + i) * 2 + 1];
child_right = sel_tree[(start + i) * 2 + 2];
/* Select the smaller one to set parent pointer */ int cmp = cmp_rec_rec_simple(
mrec[child_left], mrec[child_right],
offsets[child_left], offsets[child_right],
index, NULL);
/*********************************************************************//**
Build a selection tree for merge. The selection tree is a binary tree and should have fts_sort_pll_degree / 2 levels. With root as level 0
@return number of tree levels */ static
ulint
row_fts_build_sel_tree( /*===================*/ int* sel_tree, /*<! in/out: selection tree */ const mrec_t** mrec, /*<! in: sort record */
rec_offs** offsets, /*<! in: record offsets */
dict_index_t* index) /*<! in: index dictionary */
{
ulint treelevel = 1;
ulint num = 2;
ulint i = 0;
ulint start;
/* No need to build selection tree if we only have two merge threads */ if (fts_sort_pll_degree <= 2) { return(0);
}
while (num < fts_sort_pll_degree) {
num = num << 1;
treelevel++;
}
start = (ulint(1) << treelevel) - 1;
for (i = 0; i < fts_sort_pll_degree; i++) {
sel_tree[i + start] = int(i);
}
i = treelevel; do {
row_fts_build_sel_tree_level(
sel_tree, --i, mrec, offsets, index);
} while (i > 0);
return(treelevel);
}
/*********************************************************************//**
Read sorted file containing index data tuples and insert these data
tuples to the index
@return DB_SUCCESS or error number */
dberr_t
row_fts_merge_insert( /*=================*/
dict_index_t* index, /*!< in: index */
dict_table_t* table, /*!< in: new table */
fts_psort_t* psort_info, /*!< parallel sort info */
ulint id) /* !< in: which auxiliary table's data
to insert to */
{ const byte** b;
mem_heap_t* tuple_heap;
mem_heap_t* heap;
dberr_t error = DB_SUCCESS;
ulint* foffs;
rec_offs** offsets;
fts_tokenizer_word_t new_word;
ib_vector_t* positions;
doc_id_t last_doc_id;
ib_alloc_t* heap_alloc;
ulint i;
mrec_buf_t** buf;
pfs_os_file_t* fd;
byte** block;
byte** crypt_block; const mrec_t** mrec; int* sel_tree;
ulint height;
ulint start;
fts_psort_insert_t ins_ctx;
fts_table_t fts_table; char aux_table_name[MAX_FULL_NAME_LEN];
dict_table_t* aux_table;
dict_index_t* aux_index;
trx_t* trx;
/* We use the insert query graph as the dummy graph
needed in the row module call */
trx = trx_create(); /* The merge thread runs under the user's DDL trx; inherit its
THD so row-lock waits use the user's timeout policy. */
trx->mysql_thd = psort_info->psort_common->trx->mysql_thd;
trx_start_if_not_started(trx, true);
/* We should set the flags2 with aux_table_name here,
in order to get the correct aux table names. */
index->table->flags2 |= DICT_TF2_FTS_AUX_HEX_NAME;
fts_table.type = FTS_INDEX_TABLE;
fts_table.index_id = index->id;
fts_table.table_id = table->id;
fts_table.table = index->table;
fts_table.suffix = fts_get_suffix(id);
/* Get aux index */
fts_get_table_name(&fts_table, aux_table_name);
aux_table = dict_table_open_on_name(aux_table_name, false,
DICT_ERR_IGNORE_NONE);
ut_ad(aux_table != NULL);
aux_index = dict_table_get_first_index(aux_table);
ut_ad(!aux_index->is_instant()); /* row_merge_write_fts_node() depends on the correct value */
ut_ad(aux_index->n_core_null_bytes
== UT_BITS_IN_BYTES(aux_index->n_nullable));
/* Set TRX_ID and ROLL_PTR */
dfield_set_data(dtuple_get_nth_field(ins_ctx.tuple, 2),
&reset_trx_id, DATA_TRX_ID_LEN);
dfield_set_data(dtuple_get_nth_field(ins_ctx.tuple, 3),
&reset_trx_id[DATA_TRX_ID_LEN], DATA_ROLL_PTR_LEN);
ut_d(ins_ctx.aux_index_id = id);
const ulint space = table->space_id;
for (i = 0; i < fts_sort_pll_degree; i++) { if (psort_info[i].merge_file[id]->n_rec == 0) { /* No Rows to read */
mrec[i] = b[i] = NULL;
} else { /* Read from temp file only if it has been writtento.Otherwise,blockmemoryholds
all the sorted records */ if (psort_info[i].merge_file[id]->offset > 0
&& (!row_merge_read(
fd[i], foffs[i],
(row_merge_block_t*) block[i],
(row_merge_block_t*) crypt_block[i],
space))) {
error = DB_CORRUPTION; gotoexit;
}
/* Fetch sorted records from sort buffer and insert them into
corresponding FTS index auxiliary tables */ for (;;) {
dtuple_t* dtuple; int min_rec = 0;
if (fts_sort_pll_degree <= 2) { while (!mrec[min_rec]) {
min_rec++;
¤ 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.0.19Bemerkung:
(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.