YoushouldhavereceivedacopyoftheGNUGeneralPublicLicense alongwiththisprogram;ifnot,writetotheFreeSoftware
Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA 02110-1301 USA */
ORDER *st_select_lex::find_common_window_func_partition_fields(THD *thd)
{
ORDER *ord;
Item *item;
DBUG_ASSERT(window_funcs.elements);
List_iterator_fast<Item_window_func> it(window_funcs);
Item_window_func *first_wf= it++; if (!first_wf->window_spec->partition_list) return0;
List<Item> common_fields;
uint first_partition_elements= 0; for (ord= first_wf->window_spec->partition_list->first; ord; ord= ord->next)
{ if ((*ord->item)->real_item()->type() == Item::FIELD_ITEM)
common_fields.push_back(*ord->item, thd->mem_root);
first_partition_elements++;
} if (window_specs.elements == 1 &&
common_fields.elements == first_partition_elements) return first_wf->window_spec->partition_list->first;
List_iterator<Item> li(common_fields);
Item_window_func *wf; while (common_fields.elements && (wf= it++))
{ if (!wf->window_spec->partition_list) return0; while ((item= li++))
{ for (ord= wf->window_spec->partition_list->first; ord; ord= ord->next)
{ if (item->eq(*ord->item, false)) break;
} if (!ord)
li.remove();
}
li.rewind();
} if (!common_fields.elements) return0; if (common_fields.elements == first_partition_elements) return first_wf->window_spec->partition_list->first;
SQL_I_List<ORDER> res_list; for (ord= first_wf->window_spec->partition_list->first, item= li++;
ord; ord= ord->next)
{ if (item != *ord->item) continue; if (add_to_list(thd, res_list, item, ord->direction)) return0;
item= li++;
} return res_list.first;
}
///////////////////////////////////////////////////////////////////////////// // Sorting window functions to minimize the number of table scans // performed during the computation of these functions /////////////////////////////////////////////////////////////////////////////
#define CMP_LT -2// Less than #define CMP_LT_C -1// Less than and compatible #define CMP_EQ 0// Equal to #define CMP_GT_C 1// Greater than and compatible #define CMP_GT 2// Greater then
static int compare_order_lists(SQL_I_List<ORDER> *part_list1, int spec_number1,
SQL_I_List<ORDER> *part_list2, int spec_number2)
{ if (part_list1 == part_list2) return CMP_EQ;
ORDER *elem1= part_list1->first;
ORDER *elem2= part_list2->first; for ( ; elem1 && elem2; elem1= elem1->next, elem2= elem2->next)
{ int cmp; // remove all constants as we don't need them for comparision while(elem1 && ((*elem1->item)->real_item())->const_item())
{
elem1= elem1->next; continue;
}
/* Window frames are equal. Let's use only one of them. */ if (!win_spec1->name().str && win_spec2->name().str)
win_spec1->window_frame= win_spec2->window_frame; else
win_spec2->window_frame= win_spec1->window_frame;
return CMP_EQ;
}
if (cmp == CMP_GT || cmp == CMP_LT) return cmp;
/* one of the partitions lists is the proper beginning of the another */
cmp= compare_window_spec_joined_lists(win_spec1, win_spec2);
class Rowid_seq_cursor
{ public:
Rowid_seq_cursor() : io_cache(NULL), ref_buffer(0) {} virtual ~Rowid_seq_cursor()
{ if (ref_buffer)
my_free(ref_buffer); if (io_cache)
{
end_slave_io_cache(io_cache);
my_free(io_cache);
io_cache= NULL;
}
}
private: /* Length of one rowid element */
size_t ref_length;
/* If io_cache=!NULL, use it */
IO_CACHE *io_cache;
uchar *ref_buffer; /* Buffer for the last returned rowid */
ha_rows rownum; /* Number of the rowid that is about to be returned */
ha_rows current_ref_buffer_rownum; bool ref_buffer_valid;
/* The following are used when we are reading from an array of pointers */
uchar *cache_start;
uchar *cache_pos;
uchar *cache_end; public:
/* Informsthecursorthatweneedtomoveintothenextpartition. Thenextpartitionisprovidedintwoways: -intable->record[0].. -rownumparameterhastherownumber.
*/ void on_next_partition(ha_rows rownum)
{ /* Remember the sort key value from the new partition */
move_to(rownum);
bound_tracker.check_if_next_group();
end_of_partition= false;
}
/* Thisreturns-1whenendofpartitionwasreached.
*/ int next() override
{ int res; if (end_of_partition) return -1;
if ((res= Table_read_cursor::next()) ||
(res= fetch()))
{ /* TODO(cvicentiu) This does not consider table read failures.
Perhaps assuming end of table like this is fine in that case. */
/* This row is the final row in the table. To maintain semantics thatcursorsalwayspointtothelastvalidrow,movebackonestep,
but mark end_of_partition as true. */
Table_read_cursor::prev();
end_of_partition= true; return res;
}
if (bound_tracker.compare_with_cache())
{ /* This row is part of a new partition, don't move
forward any more untill we get informed of a new partition. */
Table_read_cursor::prev();
end_of_partition= true; return -1;
} return0;
}
/* Clear all sum functions handled by this cursor. */ void clear_sum_functions()
{
List_iterator_fast<Item_sum> iter_sum_func(sum_functions);
Item_sum *sum_func; while ((sum_func= iter_sum_func++))
{
sum_func->clear();
}
}
/* Sum functions that this cursor handles. */
List<Item_sum> sum_functions;
bool use_minus= is_preceding; if (order_direction == -1)
use_minus= !use_minus;
if (use_minus)
item_add= new (thd->mem_root) Item_func_minus(thd, src_expr, n_val); else
item_add= new (thd->mem_root) Item_func_plus(thd, src_expr, n_val);
bool is_outside_computation_bounds() const override
{ if (end_of_partition) returntrue; returnfalse;
}
private: void walk_till_non_peer()
{ if (cursor.fetch()) // ERROR return; // Current row is not a peer. if (order_direction * range_expr->cmp_read_only() <= 0) return;
remove_value_from_items();
int res; while (!(res= cursor.next()))
{ /* Note, no need to fetch the value explicitly here. The partition readcursorwillfetchittocheckifthepartitionhaschanged. TODO(cvicentiu)makethispieceofinformationnotnecessaryby reimplementingPartition_read_cursor.
*/ if (order_direction * range_expr->cmp_read_only() <= 0) break;
remove_value_from_items();
} if (res)
end_of_partition= true;
}
bool use_minus= is_preceding; if (order_direction == -1)
use_minus= !use_minus;
if (use_minus)
item_add= new (thd->mem_root) Item_func_minus(thd, src_expr, n_val); else
item_add= new (thd->mem_root) Item_func_plus(thd, src_expr, n_val);
bool is_outside_computation_bounds() const override
{ if (!added_values) returntrue; returnfalse;
}
ha_rows get_curr_rownum() const override
{ if (end_of_partition) return cursor.get_rownum(); // Cursor does not pass over partition bound. else return cursor.get_rownum() - 1; // Cursor is placed on first non peer.
}
private: bool added_values;
void walk_till_non_peer()
{
cursor.fetch(); // Current row is not a peer. if (order_direction * range_expr->cmp_read_only() < 0) return;
add_value_to_items(); // Add current row.
added_values= true; int res; while (!(res= cursor.next()))
{ if (order_direction * range_expr->cmp_read_only() < 0) break;
add_value_to_items();
} if (res)
end_of_partition= true;
}
};
void pre_next_partition(ha_rows rownum) override
{ // Save the value of the current_row
peer_tracker.check_if_next_group();
cursor.on_next_partition(rownum); // Add the current row now because our cursor has already seen it
add_value_to_items();
}
void next_row() override
{ // Check if our cursor is pointing at a peer of the current row. // If not, move forward until that becomes true if (dont_move)
{ /* Ourcurrentisnotapeerofthecurrentrow. Noneedtomovethebound.
*/ return;
}
walk_till_non_peer();
}
void pre_next_partition(ha_rows rownum) override
{ // Fetch the value from the first row
peer_tracker.check_if_next_group();
cursor.move_to(rownum);
}
void next_partition(ha_rows rownum) override {}
void pre_next_row() override
{ // Check if the new current_row is a peer of the row that our cursor is // pointing to.
move= peer_tracker.check_if_next_group();
}
void next_row() override
{ if (move)
{ /* Ourcursorispointingatthefirstrowthatwasapeeroftheprevious currentrow.Or,itwasthefirstrowinthepartition.
*/ if (cursor.fetch()) return;
// todo: need the following check ? if (!peer_tracker.compare_with_cache()) return;
remove_value_from_items();
do
{ if (cursor.next() || cursor.fetch()) return; if (!peer_tracker.compare_with_cache()) return;
remove_value_from_items();
} while (1);
}
}
///////////////////////////////////////////////////////////////////////////// // UNBOUNDED frame bounds (shared between RANGE and ROWS) /////////////////////////////////////////////////////////////////////////////
/* UNBOUNDEDPRECEDINGframebound
*/ class Frame_unbounded_preceding : public Frame_cursor
{ public:
Frame_unbounded_preceding(THD *thd,
SQL_I_List<ORDER> *partition_list,
SQL_I_List<ORDER> *order_list)
{}
/* Walk to the end of the partition, find how many rows there are. */ while (!cursor.next())
num_rows_in_partition++;
set_win_funcs_row_count(num_rows_in_partition);
}
/* Walk to the end of the partition, find how many rows there are. */ do
{ if (!order_item->is_null())
num_rows_in_partition++;
} while (!cursor.next());
void next_partition(ha_rows rownum) override
{ /* Positionourcursortopointatthefirstrowinthenewpartition (forrownum=0,itisalreadythere,otherwise,itlagsbehind)
*/
cursor.move_to(rownum); /* Cursor is in the same spot as current row. */
n_rows_behind= 0;
bool is_outside_computation_bounds() const override
{ /* As a bottom boundary, rows have not yet been added. */ if (!is_top_bound && n_rows - n_rows_behind) returntrue; returnfalse;
}
private: void move_cursor_if_possible()
{
longlong rows_difference= n_rows - n_rows_behind; if (rows_difference > 0) /* We still have to wait. */ return;
/* The cursor points to the first row in the frame. */ if (rows_difference == 0)
{ if (!is_top_bound)
{
cursor.fetch();
add_value_to_items();
} /* For top bound we don't have to remove anything as nothing was added. */ return;
}
/* We need to catch up by one row. */
DBUG_ASSERT(rows_difference == -1);
if (is_top_bound)
{
cursor.fetch();
remove_value_from_items();
cursor.next();
} else
{
cursor.next();
cursor.fetch();
add_value_to_items();
} /* We've advanced one row. We are no longer behind. */
n_rows_behind--;
}
};
private: void next_part_top(ha_rows rownum)
{ for (ha_rows i= 0; i < n_rows; i++)
{ if (cursor.fetch()) break;
remove_value_from_items(); if (cursor.next())
at_partition_end= true;
}
}
void next_part_bottom(ha_rows rownum)
{ if (cursor.fetch()) return;
add_value_to_items();
for (ha_rows i= 0; i < n_rows; i++)
{ if (cursor.next())
{
at_partition_end= true; break;
}
add_value_to_items();
} return;
}
void next_row_top()
{ if (cursor.fetch()) // PART END OR FAILURE
{
at_partition_end= true; return;
}
remove_value_from_items(); if (cursor.next())
{
at_partition_end= true; return;
}
}
void next_row_bottom()
{ if (at_partition_end) return;
if (cursor.next())
{
at_partition_end= true; return;
}
void pre_next_partition(ha_rows rownum) override
{ /* TODO(cvicentiu) Sum functions get cleared on next partition anyway during thewindowfunctioncomputationalgorithm.Eitherperformthisonlyin cursors,orremoveitfrompre_next_partition.
*/
curr_rownum= rownum;
clear_sum_functions();
}
/* Scan the rows between the top bound and bottom bound. Add all the values
between them, top bound row and bottom bound row inclusive. */ void compute_values_for_current_row()
{
THD *thd= current_thd; if (top_bound.is_outside_computation_bounds() ||
bottom_bound.is_outside_computation_bounds()) return;
for (ha_rows idx= start_rownum; idx <= bottom_rownum
&& ((idx & 0xFF) || !thd->check_killed(true)); idx++)
{ if (cursor.fetch()) //EOF break;
add_value_to_items(); if (cursor.next()) // EOF break;
}
}
};
/* A cursor that follows a target cursor. Each time a new row is added, thewindowfunctionsareclearedandonlyhavetherowatwhichthetarget ispointataddedtothem.
void pre_next_partition(ha_rows rownum) override
{ /* The offset is dependant on the current row values. We can only get
* it here accurately. When fetching other rows, it changes. */
save_offset_value();
}
void pre_next_row() override
{ /* The offset is dependant on the current row values. We can only get
* it here accurately. When fetching other rows, it changes. */
save_offset_value();
}
private: /* Check if a our position is within bounds.
* The position is passed as a parameter to avoid recalculating it. */ bool position_is_within_bounds()
{ if (!offset) return !position_cursor.is_outside_computation_bounds();
if (overflowed) returnfalse;
/* No valid bound to compare to. */ if (position_cursor.is_outside_computation_bounds() ||
top_bound->is_outside_computation_bounds() ||
bottom_bound->is_outside_computation_bounds()) returnfalse;
/* We are over the bound. */ if (position < top_bound->get_curr_rownum()) returnfalse; if (position > bottom_bound->get_curr_rownum()) returnfalse;
returntrue;
}
/* Get the current position, accounting for the offset value, if present. NOTE:Thisfunctiondoesnotcheckover/underflow.
*/ void get_current_position()
{
position = position_cursor.get_curr_rownum();
overflowed= false; if (offset)
{ if (offset_value < 0 &&
position + offset_value > position)
{
overflowed= true;
} if (offset_value > 0 &&
position + offset_value < position)
{
overflowed= true;
}
position += offset_value;
}
}
if (bound->offset == NULL) /* this is UNBOUNDED */
{ /* The following serve both RANGE and ROWS: */ if (is_preceding) returnnew Frame_unbounded_preceding(thd,
spec->partition_list,
spec->order_list);
if (frame->units == Window_frame::UNITS_ROWS)
{
ha_rows n_rows= bound->offset->val_int(); /* These should be handled in the parser */
DBUG_ASSERT(!bound->offset->null_value);
DBUG_ASSERT((longlong) n_rows >= 0); if (is_preceding) returnnew Frame_n_rows_preceding(is_top_bound, n_rows);
if (bound->precedence_type == Window_frame_bound::CURRENT)
{ if (frame->units == Window_frame::UNITS_ROWS)
{ if (is_top_bound) returnnew Frame_rows_current_row_top;
if (init_read_record(&info, current_thd, tbl, NULL/*select*/, filesort_result, 0, 1, FALSE)) returntrue;
Cursor_manager *cursor_manager; while ((cursor_manager= iter_cursor_managers++))
cursor_manager->initialize_cursors(&info);
/* One partition tracker for each window function. */
List<Group_bound_tracker> partition_trackers;
Item_window_func *win_func; while ((win_func= iter_win_funcs++))
{
Group_bound_tracker *tracker= new Group_bound_tracker(thd,
win_func->window_spec->partition_list); // TODO(cvicentiu) This should be removed and placed in constructor.
tracker->init();
partition_trackers.push_back(tracker);
}
while (true)
{ if ((err= info.read_record())) break; // End of file.
/* Remember current row so that we can restore it before computing
each window function. */
tbl->file->position(tbl->record[0]);
memcpy(rowid_buf, tbl->file->ref, tbl->file->ref_length);
Group_bound_tracker *tracker; while ((win_func= iter_win_funcs++) &&
(tracker= iter_part_trackers++) &&
(cursor_manager= iter_cursor_managers++))
{ if (tracker->check_if_next_group() || (rownum == 0))
{ /* TODO(cvicentiu)
Clearing window functions should happen through cursors. */
win_func->window_func()->clear();
cursor_manager->notify_cursors_partition_changed(rownum);
} else
{
cursor_manager->notify_cursors_next_row();
}
/* Check if we found any error in the window function while adding values
through cursors. */ if (unlikely(thd->is_error() || thd->is_killed()))
{
ret= true; break;
}
/* Return to current row after notifying cursors for each window
function. */ if (tbl->file->ha_rnd_pos(tbl->record[0], rowid_buf))
{
ret= true; break;
}
}
/* We now have computed values for each window function. They can now
be saved in the current row. */ if (save_window_function_values(window_functions, tbl, rowid_buf))
{
ret= true; break;
}
rownum++;
}
/* Make a list that is a concation of two lists of ORDER elements */
static ORDER* concat_order_lists(MEM_ROOT *mem_root, ORDER *list1, ORDER *list2)
{ if (!list1)
{
list1= list2;
list2= NULL;
}
ORDER *res= NULL; // first element in the new list
ORDER *prev= NULL; // last element in the new list
ORDER *cur_list= list1; // this goes through list1, list2 while (cur_list)
{ for (ORDER *cur= cur_list; cur; cur= cur->next)
{
ORDER *copy= (ORDER*)alloc_root(mem_root, sizeof(ORDER));
memcpy((void *) copy, (void *) cur, sizeof(ORDER)); if (prev)
prev->next= copy;
prev= copy; if (!res)
res= copy;
}
switch (type)
{ /* Distinct is not yet supported. */ case Item_sum::GROUP_CONCAT_FUNC:
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "GROUP_CONCAT() aggregate as window function"); returntrue; case Item_sum::SUM_DISTINCT_FUNC:
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "SUM(DISTINCT) aggregate as window function"); returntrue; case Item_sum::AVG_DISTINCT_FUNC:
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "AVG(DISTINCT) aggregate as window function"); returntrue; case Item_sum::COUNT_DISTINCT_FUNC:
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "COUNT(DISTINCT) aggregate as window function"); returntrue; case Item_sum::JSON_ARRAYAGG_FUNC:
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "JSON_ARRAYAGG() aggregate as window function"); returntrue; case Item_sum::JSON_OBJECTAGG_FUNC:
my_error(ER_NOT_SUPPORTED_YET, MYF(0), "JSON_OBJECTAGG() aggregate as window function"); returntrue; default: break;
}
return window_functions.push_back(win_func);
}
/* Computethevalueofwindowfunctionforallrows.
*/ bool Window_func_runner::exec(THD *thd, TABLE *tbl, SORT_INFO *filesort_result)
{
List_iterator_fast<Item_window_func> it(window_functions);
Item_window_func *win_func; while ((win_func= it++))
{
win_func->set_phase_to_computation(); // TODO(cvicentiu) Setting the aggregator should probably be done during // setup of Window_funcs_sort.
win_func->window_func()->set_aggregator(thd,
Aggregator::SIMPLE_AGGREGATOR);
}
it.rewind();
List<Cursor_manager> cursor_managers; if (get_window_functions_required_cursors(thd, window_functions,
&cursor_managers)) returntrue;
/* Go through the sorted array and compute the window function */ bool is_error= compute_window_func(thd,
window_functions,
cursor_managers,
tbl, filesort_result); while ((win_func= it++))
{
win_func->set_phase_to_retrieval();
}
/* Sort the table based on the most specific sorting criteria of
the window functions. */ if (create_sort_index(thd, join, join_tab, filesort)) returntrue;
/* The iterator should point to a valid function at the start of execution. */
DBUG_ASSERT(win_func); do
{
spec= win_func->window_spec; int win_func_order_elements= spec->partition_list->elements +
spec->order_list->elements; if (win_func_order_elements >= longest_order_elements)
{
win_func_with_longest_order= win_func;
longest_order_elements= win_func_order_elements;
} if (runner.add_function_to_run(win_func)) returntrue;
it++;
win_func= it.peek();
} while (win_func && !(win_func->marker & MARKER_SORTORDER_CHANGE));
ORDER* sort_order= concat_order_lists(thd->mem_root,
spec->partition_list->first,
spec->order_list->first); if (sort_order == NULL) // No partition or order by clause.
{ /* TODO(cvicentiu) This is used as a way to allow an empty OVER () clauseforwindowfunctions.However,abetterapproachis tonotcallFilesortatallinthiscaseandjustreadwhateverorder thetemporarytablehas. Duetocursorsnotworkingforout_of_memorycases(yet!),wehavetorun filesorttogenerateasortbufferoftheresults. Inthiscasewesortbythefirstfieldofthetemporarytable. Weshouldhavethisfieldavailable,evenifitisawindow_function field.Wedon'tcareoftheparticularsortingresultinthiscase.
*/
ORDER *order= (ORDER *)alloc_root(thd->mem_root, sizeof(ORDER));
memset((void *) order, 0, sizeof(*order));
Item_field *item= new (thd->mem_root) Item_field(thd, join_tab->table->field[0]); if (item)
item->set_refers_to_temp_table();
order->item= (Item **)alloc_root(thd->mem_root, 2 * sizeof(Item *));
order->item[1]= NULL;
order->item[0]= item;
order->field= join_tab->table->field[0];
sort_order= order;
}
filesort= new (thd->mem_root) Filesort(sort_order, HA_POS_ERROR, true, NULL);
/* Apply the same condition that the subsequent sort has. */
filesort->select= sel;
Explain_aggr_window_funcs*
Window_funcs_computation::save_explain_plan(MEM_ROOT *mem_root, bool is_analyze)
{
Explain_aggr_window_funcs *xpl= new Explain_aggr_window_funcs;
List_iterator<Window_funcs_sort> it(win_func_sorts);
Window_funcs_sort *srt; if (!xpl) return0; while ((srt = it++))
{
Explain_aggr_filesort *eaf= new Explain_aggr_filesort(mem_root, is_analyze, srt->filesort); if (!eaf) return0;
xpl->sorts.push_back(eaf, mem_root);
} return xpl;
}
bool st_select_lex::add_window_func(THD *thd, Item_window_func *win_func)
{ if (parsing_place != SELECT_LIST)
fields_in_window_functions+= win_func->window_func()->argument_count(); /* We may use it later for other clauses, now just ORDER_CLAUSE */ if (thd->where == THD_WHERE::ORDER_CLAUSE)
parent_lex->clause_winfuncs.push_back(win_func, thd->mem_root); return window_funcs.push_back(win_func);
}
///////////////////////////////////////////////////////////////////////////// // Unneeded comments (will be removed when we develop a replacement for // the feature that was attempted here ///////////////////////////////////////////////////////////////////////////// /* TODOGetthiscodetosetcan_compute_window_functionduringpreparation, notduringexecution.
Theruleforhavingcompatiblesortingisthus: Eachpartitionordermustcontaintheotherwindowfunctionspartitions prefixes,orbeaprefixitself.Thismustholdtrueforallpartitions. Analogfortheorderbyclause.
*/ #if0
List<Item_window_func> window_functions;
SQL_I_List<ORDER> largest_partition;
SQL_I_List<ORDER> largest_order_by; bool can_compute_window_live = !need_tmp; // Construct the window_functions item list and check if they can be // computed using only one sorting. // // TODO: Perhaps group functions into compatible sorting bins // to minimize the number of sorting passes required to compute all of them. while ((item= it++))
{ if (item->type() == Item::WINDOW_FUNC_ITEM)
{
Item_window_func *item_win = (Item_window_func *) item;
window_functions.push_back(item_win); if (!can_compute_window_live) continue; // No point checking since we have to perform multiple sorts.
Window_spec *spec = item_win->window_spec; // Having an empty partition list on one window function and a // not empty list on a separate window function causes the sorting // to be incompatible. // // Example: // over (partition by a, order by x) && over (order by x). // // The first function requires an ordering by a first and then by x, // while the second function requires an ordering by x first. // The same restriction is not required for the order by clause. if (largest_partition.elements && !spec->partition_list.elements)
{
can_compute_window_live= FALSE; continue;
}
can_compute_window_live= test_if_order_compatible(largest_partition,
spec->partition_list); if (!can_compute_window_live) continue;
can_compute_window_live= test_if_order_compatible(largest_order_by,
spec->order_list); if (!can_compute_window_live) continue;
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.