YoushouldhavereceivedacopyoftheGNUGeneralPublicLicense alongwiththisprogram;ifnot,writetotheFreeSoftware
Foundation, Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02111-1301 USA */
#include <my_global.h>
/* This C++ file's header file */ #include"./rdb_index_merge.h"
/* MySQL header files */ #include"../sql/sql_class.h"
Thishelpsmitigatepotentialtrimstallsonflashwhenlargefilesare beingdeletedtooquickly.
*/ if (m_merge_tmp_file_removal_delay > 0) {
uint64 curr_size = m_merge_buf_size * m_merge_file.m_num_sort_buffers; for (uint i = 0; i < m_merge_file.m_num_sort_buffers; i++) { if (my_chsize(m_merge_file.m_fd, curr_size, 0, MYF(MY_WME)) > 0) { // NO_LINT_DEBUG
sql_print_error("Error truncating file during fast index creation.");
}
my_sleep(m_merge_tmp_file_removal_delay * 1000); // Not aborting on fsync error since the tmp file is not used anymore if (mysql_file_sync(m_merge_file.m_fd, MYF(MY_WME))) { // NO_LINT_DEBUG
sql_print_error("Error flushing truncated MyRocks merge buffer.");
}
curr_size -= m_merge_buf_size;
}
}
/** Createamergefileinthegivenlocation.
*/ int Rdb_index_merge::merge_file_create() {
DBUG_ASSERT(m_merge_file.m_fd == -1);
int fd; #ifdef MARIAROCKS_NOT_YET // mysql_tmpfile_path use /* If no path set for tmpfile, use mysql_tmpdir by default */ if (m_tmpfile_path == nullptr) {
fd = mysql_tmpfile("myrocks");
} else {
fd = mysql_tmpfile_path(m_tmpfile_path, "myrocks");
} #else
fd = mysql_tmpfile("myrocks"); #endif if (fd < 0) { // NO_LINT_DEBUG
sql_print_error("Failed to create temp file during fast index creation."); return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
Ifbufferinmemoryisfull,writethebufferouttodisksortedusingthe offsettree,andclearthetree.(Happensinmerge_buf_write)
*/ int Rdb_index_merge::add(const rocksdb::Slice &key, const rocksdb::Slice &val) { /* Adding a record after heap is already created results in error */
DBUG_ASSERT(m_merge_min_heap.empty());
/* Checkifsortbufferisgoingtobeoutofspace,ifsowriteit outtodiskinsortedorderusingoffsettree.
*/ const ulonglong total_offset = RDB_MERGE_CHUNK_LEN +
m_rec_buf_unsorted->m_curr_offset +
RDB_MERGE_KEY_DELIMITER + RDB_MERGE_VAL_DELIMITER +
key.size() + val.size(); if (total_offset >= m_rec_buf_unsorted->m_total_size) { /* Iftheoffsettreeisemptyhere,thatmeansthattheproposedkeyto addistoolargeforthebuffer.
*/ if (m_offset_tree.empty()) { // NO_LINT_DEBUG
sql_print_error( "Sort buffer size is too small to process merge. " "Please set merge buffer size to a higher value."); return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
if (merge_buf_write()) { // NO_LINT_DEBUG
sql_print_error("Error writing sort buffer to disk."); return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
}
/* Find sort order of the new record */ auto res =
m_offset_tree.emplace(m_rec_buf_unsorted->m_block.get() + rec_offset,
m_cf_handle->GetComparator()); if (!res.second) {
my_printf_error(ER_DUP_ENTRY, "Failed to insert the record: the key already exists",
MYF(0)); return ER_DUP_ENTRY;
}
/* Write actual chunk size to first 8 bytes of the merge buffer */
merge_store_uint64(m_output_buf->m_block.get(),
m_rec_buf_unsorted->m_curr_offset + RDB_MERGE_CHUNK_LEN);
m_output_buf->m_curr_offset += RDB_MERGE_CHUNK_LEN;
/* Allocate buffers for each chunk */ for (ulonglong i = 0; i < m_merge_file.m_num_sort_buffers; i++) { constauto entry =
std::make_shared<merge_heap_entry>(m_cf_handle->GetComparator());
if (total_size == (size_t)-1) { return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
/* Can reach this condition if an index was added on table w/ no rows */ if (total_size - RDB_MERGE_CHUNK_LEN == 0) { break;
}
/* Read the first record from each buffer to initially populate the heap */ if (entry->read_rec(&entry->m_key, &entry->m_val)) { // NO_LINT_DEBUG
sql_print_error("Chunk size is too small to process merge."); return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
/* Ifmerge_read_recfails,itmeanstheeitherthechunkwascutoff orwe'vereachedtheendoftherespectivechunk.
*/ if (entry->read_rec(&entry->m_key, &entry->m_val)) { if (entry->read_next_chunk_from_disk(m_merge_file.m_fd)) { return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
/* Try reading record again, should never fail. */ if (entry->read_rec(&entry->m_key, &entry->m_val)) { return HA_ERR_ROCKSDB_MERGE_FILE_ERR;
}
}
/* Push entry back on to the heap w/ updated buffer + offset ptr */
m_merge_min_heap.push(std::move(entry));
/* Return the current top record on heap */
merge_heap_top(key, val); return HA_EXIT_SUCCESS;
}
int Rdb_index_merge::merge_heap_entry::read_next_chunk_from_disk(File fd) { if (m_chunk_info->read_next_chunk_from_disk(fd)) { return HA_EXIT_FAILURE;
}
/* Store key and value w/ their respective delimiters at the given offset */ void Rdb_index_merge::merge_buf_info::store_key_value( const rocksdb::Slice &key, const rocksdb::Slice &val) {
store_slice(key);
store_slice(val);
}
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.