diff --git a/src/dsql/Nodes.h b/src/dsql/Nodes.h index d3853c4090d..f0c258ef0f7 100644 --- a/src/dsql/Nodes.h +++ b/src/dsql/Nodes.h @@ -1006,6 +1006,8 @@ class AggNode : public TypedNode static constexpr unsigned CAP_WANTS_AGG_CALLS = 0x04; // wants winPass call in a window static constexpr unsigned CAP_WANTS_WIN_PASS_CALL = 0x08; + // depends on the window frame boundaries, even when not respecting a user-defined frame + static constexpr unsigned CAP_USES_WINDOW_FRAME = 0x10; protected: struct AggInfo diff --git a/src/dsql/WinNodes.h b/src/dsql/WinNodes.h index bc4cfa8e3e5..caadc2318d6 100644 --- a/src/dsql/WinNodes.h +++ b/src/dsql/WinNodes.h @@ -135,7 +135,7 @@ class CumeDistWinNode final : public WinFuncNode unsigned getCapabilities() const override { - return CAP_SUPPORTS_WINDOW_FRAME | CAP_WANTS_WIN_PASS_CALL; + return CAP_SUPPORTS_WINDOW_FRAME | CAP_USES_WINDOW_FRAME | CAP_WANTS_WIN_PASS_CALL; } void getChildren(NodeRefsHolder& holder, bool dsql) const override diff --git a/src/jrd/recsrc/AggregatedStream.cpp b/src/jrd/recsrc/AggregatedStream.cpp index 6d0248c102f..6ca074ca5f2 100644 --- a/src/jrd/recsrc/AggregatedStream.cpp +++ b/src/jrd/recsrc/AggregatedStream.cpp @@ -178,7 +178,10 @@ bool BaseAggWinStream::evaluateGroup(thread_db* tdbb) const } } else + { cacheValues(tdbb, request, m_group, impure->groupValues, DummyAdjustFunctor()); + static_cast(this)->groupRecord(tdbb, request, true); + } // Loop thru records until either a value change or EOF @@ -193,6 +196,8 @@ bool BaseAggWinStream::evaluateGroup(thread_db* tdbb) const if (lookForChange(tdbb, request, m_group, NULL, impure->groupValues)) impure->state = STATE_FETCHED; + else + static_cast(this)->groupRecord(tdbb, request, false); } else impure->state = STATE_EOF; diff --git a/src/jrd/recsrc/RecordSource.h b/src/jrd/recsrc/RecordSource.h index a07e808fa43..b7e88c8c2bf 100644 --- a/src/jrd/recsrc/RecordSource.h +++ b/src/jrd/recsrc/RecordSource.h @@ -1090,6 +1090,12 @@ namespace Jrd protected: void internalGetPlan(thread_db* tdbb, PlanEntry& planEntry, unsigned level, bool recurse) const override; bool internalGetRecord(thread_db* tdbb) const override; + + public: + // Called by evaluateGroup for each record that belongs to the group being evaluated. + void groupRecord(thread_db* /*tdbb*/, Request* /*request*/, bool /*first*/) const + { + } }; class WindowedStream : public RecordSource @@ -1143,6 +1149,14 @@ namespace Jrd SINT64 partitionPending, rangePending; Block partitionBlock, windowBlock; impure_value_ex startOffset, endOffset; + + // Last positions of the peer groups of the partition, found when it is evaluated. + // Bit n is set when the record at position peerBase + n ends its peer group. + FB_UINT64* peerBits; + ULONG peerCapacity; // allocated words + ULONG peerUsed; // words that may have bits set + SINT64 peerBase; + bool peerOverflow; // the partition is too big for the bitmap }; public: @@ -1161,6 +1175,8 @@ namespace Jrd bool isDependent(const StreamList& streams) const override; void nullRecords(thread_db* tdbb) const override; + void groupRecord(thread_db* tdbb, Request* request, bool first) const; + protected: void internalGetPlan(thread_db* tdbb, PlanEntry& planEntry, unsigned level, bool recurse) const override; void internalOpen(thread_db* tdbb) const override; @@ -1186,6 +1202,8 @@ namespace Jrd SINT64 locateFrameGroups(thread_db* tdbb, Request* request, Impure* impure, const Frame* frame, const impure_value_ex* offsetValue, SINT64 position, bool startFrame) const; + SINT64 locatePeerGroupEnd(thread_db* tdbb, Request* request, Impure* impure, + SINT64 position) const; private: NestConst m_order; @@ -1196,6 +1214,8 @@ namespace Jrd NestValueArray m_winPassSources, m_winPassTargets; Exclusion m_exclusion; UCHAR m_invariantOffsets; // 0x1 | 0x2 bitmask + bool m_needsFrame = false; // some function depends on the frame + bool m_usesPeerGroupEnd = false; // peer group ends are found when evaluating the partition }; public: diff --git a/src/jrd/recsrc/WindowedStream.cpp b/src/jrd/recsrc/WindowedStream.cpp index 1094976a070..06a46eff828 100644 --- a/src/jrd/recsrc/WindowedStream.cpp +++ b/src/jrd/recsrc/WindowedStream.cpp @@ -29,6 +29,7 @@ #include "../jrd/vio_proto.h" #include "../jrd/optimizer/Optimizer.h" #include "RecordSource.h" +#include #include using namespace Firebird; @@ -483,6 +484,7 @@ WindowedStream::WindowStream::WindowStream(thread_db* tdbb, CompilerScratch* csb m_invariantOffsets(0) { // Separate nodes that requires the winPass call. + // While here, verify if some node depends on the frame. If none, it should not be evaluated. const NestConst* const sourceEnd = m_windowMap->sourceList.end(); @@ -508,9 +510,34 @@ WindowedStream::WindowStream::WindowStream(thread_db* tdbb, CompilerScratch* csb m_winPassSources.add(*source); m_winPassTargets.add(*target); } + + if ((capabilities & AggNode::CAP_WANTS_AGG_CALLS) || + (capabilities & AggNode::CAP_USES_WINDOW_FRAME) || + (capabilities & AggNode::CAP_RESPECTS_WINDOW_FRAME) == AggNode::CAP_RESPECTS_WINDOW_FRAME) + { + m_needsFrame = true; + } } } + // Peer group ends are needed by RANGE ... AND CURRENT ROW frames and by the pending range of frames + // without exclusion. They are found during the partition evaluation, which already reads all + // its records, avoiding a look-ahead (and returning back) for each record. + if (m_order && m_needsFrame) + { + const bool rangeOrGroups = m_frameExtent->unit == FrameExtent::Unit::RANGE || + m_frameExtent->unit == FrameExtent::Unit::GROUPS; + const bool unboundedFrame = + m_frameExtent->frame1->bound == Frame::Bound::PRECEDING && !m_frameExtent->frame1->value && + m_frameExtent->frame2->bound == Frame::Bound::FOLLOWING && !m_frameExtent->frame2->value; + + m_usesPeerGroupEnd = + (m_frameExtent->unit == FrameExtent::Unit::RANGE && + m_frameExtent->frame2->bound == Frame::Bound::CURRENT_ROW) || + (rangeOrGroups && m_exclusion == Exclusion::NO_OTHERS && !unboundedFrame && + m_frameExtent->frame2->bound != Frame::Bound::CURRENT_ROW); + } + m_arithNodes.resize(2); if (m_order) @@ -594,6 +621,69 @@ void WindowedStream::WindowStream::close(thread_db* tdbb) const BaseAggWinStream::close(tdbb); } +// Maximum number of records of a partition for which peer group ends are kept in a bitmap. +static constexpr SINT64 MAX_PEER_BITS = SINT64(1) << 26; + +// Called for each record of the partition being evaluated, with the record already fetched. +// Finds the peer group ends. +void WindowedStream::WindowStream::groupRecord(thread_db* tdbb, Request* request, bool first) const +{ + if (!m_usesPeerGroupEnd) + return; + + Impure* const impure = getImpure(request); + const SINT64 position = (SINT64) m_next->getPosition(request) - 1; + + if (first) + { + if (impure->peerBits && impure->peerUsed) + memset(impure->peerBits, 0, impure->peerUsed * sizeof(FB_UINT64)); + + impure->peerUsed = 0; + impure->peerBase = position; + impure->peerOverflow = false; + + cacheValues(tdbb, request, &m_order->expressions, impure->orderValues, DummyAdjustFunctor()); + return; + } + + if (!lookForChange(tdbb, request, &m_order->expressions, m_order, impure->orderValues)) + return; + + // The previous record ended a peer group. + + const SINT64 bit = position - 1 - impure->peerBase; + + if (bit >= MAX_PEER_BITS) + impure->peerOverflow = true; + else if (!impure->peerOverflow) + { + const ULONG word = (ULONG) (bit / 64); + + if (word >= impure->peerCapacity) + { + const ULONG capacity = MAX(word + 1, impure->peerCapacity * 2); + const auto bits = FB_NEW_POOL(*tdbb->getDefaultPool()) FB_UINT64[capacity]; + + memset(bits, 0, capacity * sizeof(FB_UINT64)); + + if (impure->peerBits) + { + memcpy(bits, impure->peerBits, impure->peerUsed * sizeof(FB_UINT64)); + delete [] impure->peerBits; + } + + impure->peerBits = bits; + impure->peerCapacity = capacity; + } + + impure->peerBits[word] |= FB_UINT64(1) << (bit % 64); + impure->peerUsed = MAX(impure->peerUsed, word + 1); + } + + cacheValues(tdbb, request, &m_order->expressions, impure->orderValues, DummyAdjustFunctor()); +} + bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const { JRD_reschedule(tdbb); @@ -630,7 +720,21 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const return false; } - m_next->locate(tdbb, count); + if (m_usesPeerGroupEnd) + { + m_next->locate(tdbb, 0); + + for (FB_UINT64 i = 0; i < count; ++i) + { + if (!m_next->getRecord(tdbb)) + fb_assert(false); + + groupRecord(tdbb, request, i == 0); + } + } + else + m_next->locate(tdbb, count); + impure->state = STATE_EOF; } @@ -653,7 +757,27 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const exclusion1.invalidate(); exclusion2.invalidate(); - if (impure->rangePending > 0 && m_exclusion == Exclusion::NO_OTHERS) + // Reposition the stream in the current record if it was moved. + // Use force when the stream was located without fetching, as then its position does not + // reflect the record buffer contents. + const auto syncCurrentRecord = [&](bool force = false) + { + if (force || (SINT64) m_next->getPosition(request) != position + 1) + { + m_next->locate(tdbb, position); + + if (!m_next->getRecord(tdbb)) + fb_assert(false); + } + }; + + if (!m_needsFrame) + { + // No function depends on the frame, so there is no need to evaluate it. + if (position == impure->partitionBlock.startPosition) + aggInit(tdbb, request, m_windowMap); + } + else if (impure->rangePending > 0 && m_exclusion == Exclusion::NO_OTHERS) --impure->rangePending; else { @@ -758,28 +882,7 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const else if (m_frameExtent->unit == FrameExtent::Unit::RANGE && m_frameExtent->frame2->bound == Frame::Bound::CURRENT_ROW) { - SINT64 rangePos = position; - cacheValues(tdbb, request, &m_order->expressions, impure->orderValues, - DummyAdjustFunctor()); - - while (++rangePos <= impure->partitionBlock.endPosition) - { - if (!m_next->getRecord(tdbb)) - fb_assert(false); - - if (lookForChange(tdbb, request, &m_order->expressions, m_order, - impure->orderValues)) - { - break; - } - } - - impure->windowBlock.endPosition = rangePos - 1; - - m_next->locate(tdbb, position); - - if (!m_next->getRecord(tdbb)) - fb_assert(false); + impure->windowBlock.endPosition = locatePeerGroupEnd(tdbb, request, impure, position); } // range between ... and {preceding | following} else if (m_frameExtent->unit == FrameExtent::Unit::RANGE && @@ -806,31 +909,15 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const (m_frameExtent->unit == FrameExtent::Unit::RANGE || m_frameExtent->unit == FrameExtent::Unit::GROUPS)) { - SINT64 rangePos = position; - cacheValues(tdbb, request, &m_order->expressions, impure->orderValues, - DummyAdjustFunctor()); - - while (++rangePos <= impure->partitionBlock.endPosition) - { - if (!m_next->getRecord(tdbb)) - fb_assert(false); - - if (lookForChange(tdbb, request, &m_order->expressions, m_order, - impure->orderValues)) - { - break; - } - } + // {range | groups} between ... and current row: the window end is the peer group end. + const SINT64 peerGroupEnd = m_frameExtent->frame2->bound == Frame::Bound::CURRENT_ROW ? + impure->windowBlock.endPosition : + locatePeerGroupEnd(tdbb, request, impure, position); - impure->rangePending = rangePos - position - 1; + impure->rangePending = peerGroupEnd - position; } - m_next->locate(tdbb, position); - - if (!m_next->getRecord(tdbb)) - fb_assert(false); - - //// TODO: There is no need to pass record by record when m_aggSources.isEmpty() + syncCurrentRecord(); const bool invalidFrame = !impure->windowBlock.isValid() || impure->windowBlock.endPosition < impure->windowBlock.startPosition || @@ -875,10 +962,16 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const aggExecute(tdbb, request, m_aggSources, m_aggTargets); - m_next->locate(tdbb, position); - - if (!m_next->getRecord(tdbb)) - fb_assert(false); + syncCurrentRecord(); + } + else if (lastWindow.isValid() && + impure->windowBlock.startPosition == lastWindow.startPosition && + impure->windowBlock.endPosition == position && + lastWindow.endPosition == position - 1) + { + // The window grows only by the current record, which is already fetched. + aggPass(tdbb, request, m_aggSources, m_aggTargets); + aggExecute(tdbb, request, m_aggSources, m_aggTargets); } else { @@ -887,6 +980,8 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const // This may be incompatible with some function like LIST, but currently LIST cannot // be used in ordered windows anyway. + bool recordMoved = false; + if (!lastWindow.isValid() || impure->windowBlock.startPosition > lastWindow.startPosition || impure->windowBlock.endPosition < lastWindow.endPosition) @@ -898,6 +993,7 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const { if (impure->windowBlock.startPosition < lastWindow.startPosition) { + recordMoved = true; m_next->locate(tdbb, impure->windowBlock.startPosition); SINT64 pending = lastWindow.startPosition - impure->windowBlock.startPosition; @@ -925,10 +1021,7 @@ bool WindowedStream::WindowStream::internalGetRecord(thread_db* tdbb) const aggExecute(tdbb, request, m_aggSources, m_aggTargets); - m_next->locate(tdbb, position); - - if (!m_next->getRecord(tdbb)) - fb_assert(false); + syncCurrentRecord(recordMoved); } } } @@ -1245,6 +1338,46 @@ SINT64 WindowedStream::WindowStream::locateFrameGroups(thread_db* tdbb, Request* return restoreAndReturn(startFrame ? groupBlock.startPosition : groupBlock.endPosition); } +// Returns the last position of the peer group of the current record. The stream is not repositioned +// when the peer group ends were found during the partition evaluation. +SINT64 WindowedStream::WindowStream::locatePeerGroupEnd(thread_db* tdbb, Request* request, + Impure* impure, SINT64 position) const +{ + if (m_usesPeerGroupEnd && !impure->peerOverflow) + { + const SINT64 first = position - impure->peerBase; + const SINT64 words = impure->peerUsed; + + for (SINT64 word = first / 64; word < words; ++word) + { + FB_UINT64 bits = impure->peerBits[word]; + + if (word == first / 64) + bits &= ~FB_UINT64(0) << (first % 64); + + if (bits) + return impure->peerBase + word * 64 + (SINT64) std::countr_zero(bits); + } + + // The last peer group of the partition. + return impure->partitionBlock.endPosition; + } + + SINT64 rangePos = position; + cacheValues(tdbb, request, &m_order->expressions, impure->orderValues, DummyAdjustFunctor()); + + while (++rangePos <= impure->partitionBlock.endPosition) + { + if (!m_next->getRecord(tdbb)) + fb_assert(false); + + if (lookForChange(tdbb, request, &m_order->expressions, m_order, impure->orderValues)) + break; + } + + return rangePos - 1; +} + SINT64 WindowedStream::WindowStream::locateFrameRange(thread_db* tdbb, Request* request, Impure* impure, const Frame* frame, const dsc* offsetDesc, SINT64 position) const {