yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4121872371
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/rules/rewrite/NormalizeOlapTableStreamScan.java:
##########
@@ -262,6 +263,15 @@ private Plan makeSnapshotScan(LogicalOlapTableStreamScan
scan, CascadesContext c
OlapTable baseTable = streamWrapper.getBaseTable();
List<Slot> originSlots = scan.getOutput();
selectedPartitionIds =
streamWrapper.filterConsumedPartitionIds(selectedPartitionIds);
+ // What this read is about to answer with is recorded here, where the
plan is final and the read states
+ // are in place, and before the read is built: the refresh that reads
it back may not record a
+ // partition it replaced through an image of the table as of the
offset as caught up -- what it wrote
+ // is that image, and the delta that would bring the table up to date
does not apply to it. See
+ // MTMVTask#executePartitionBasedRefresh.
+ if (streamWrapper.readsSnapshotOfAnOlderImage(selectedPartitionIds)) {
Review Comment:
Right, and the reason it slipped is in the question it asked. It asked
whether the read answered with the image at the offset, and took everything
else for the table as it is now. A read that drops a partition with no
consumption baseline answers with neither: there is no offset to read that
partition from, and rows it holds by now are simply missing from the answer. So
the partition was recorded as caught up while the rows it should have joined in
were absent, and a current snapshot was published over it.
`81fedd86e27` asks the question the answer deserves:
`OlapTableStreamWrapper#answersWithTheCurrentTable`. It says no for a partition
read behind its offset, and no for a partition the read dropped **that holds
rows** -- the read is missing them either way. It says yes for a partition the
read dropped that holds nothing, which is what the round before this one asked
for: withholding there withholds a partition that was read as the table is, and
costs a rebuild it does not owe.
The two ways are pinned per key type in `OlapTableStreamSnapshotReadTest`,
including the dropped-partition-with-rows and dropped-partition-without-rows
pair, and the suites around these reads are green:
`test_ivm_partial_rebuild_is_not_recorded`,
`test_ivm_fallback_keeps_the_whole_scope`,
`test_ivm_fallback_stream_multi_batch` and `..._dup`, `test_ivm_snapshot`,
`test_ivm_partition_sync_limit_with_window`. The judgement stays per batch, so
a batch that answers incompletely keeps every partition it replaced unrecorded;
that errs towards one more rebuild and never towards rows.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -663,54 +746,299 @@ public void alterIvmInfo(IvmInfo ivmInfo) {
* into the journal, and a replay that replaces the field would leave the
caller's reference
* pointing at state that is no longer the MV's. Changing the states is
the MV's own job, under its
* write lock.
- *
- * <p>A missing map -- an image written before the field existed, or a
non-IVM MV -- reads as empty.
*/
public Map<String, MTMVPartitionState> getPartitionStates() {
readMvLock();
try {
- if (partitionStates == null) {
- return Collections.emptyMap();
- }
return
Collections.unmodifiableMap(MTMVPartitionState.copyOf(partitionStates));
} finally {
readMvUnlock();
}
}
+ /**
+ * The partitions whose requirement has been raised and not met, which a
refresh has to rebuild rather
+ * than catch up.
+ *
+ * <p>Detached names rather than the states themselves: a caller that only
routes by them has no
+ * business holding the map the MV journals, and it needs nothing else
from an entry.
+ */
+ public Set<String> getPartitionsNeedingRebuild() {
+ // Built before the lock, like the map getLatestEpochs returns: which
entries go in is what needs
+ // the lock, not having somewhere to put them.
+ Set<String> res = Sets.newLinkedHashSet();
+ readMvLock();
+ try {
+ for (Entry<String, MTMVPartitionState> entry :
partitionStates.entrySet()) {
+ if (entry.getValue().isDirty()) {
+ res.add(entry.getKey());
+ }
+ }
+ return res;
+ } finally {
+ readMvUnlock();
+ }
+ }
+
+ /**
+ * Whether every partition the MV holds needs a rebuild, which is when a
whole-MV refresh does nothing
+ * the per-partition routing would not.
+ *
+ * <p>A partition that holds data and does not need one makes this false:
a whole-MV refresh would
+ * recompute it for nothing, which is the waste the per-partition routing
exists to avoid. A partition
+ * that was never refreshed does not count against it -- a whole-MV
refresh fills it, which its
+ * per-partition branch would do as well -- and it needs no clause of its
own: an aligned entry is
+ * {@code {0, 1}}, so it is behind its requirement already. An MV with no
partitions is not an
+ * escalation either.
+ *
+ * <p>Read in place rather than through {@link #getPartitionStates()}: the
caller asks a yes/no
+ * question, and copying the map to answer it would allocate a state
object per partition, under this
+ * lock, on every refresh -- including the ones that escalate nothing.
+ */
+ public boolean allPartitionsNeedRebuild() {
+ readMvLock();
+ try {
+ return !partitionStates.isEmpty()
+ &&
partitionStates.values().stream().allMatch(MTMVPartitionState::isDirty);
+ } finally {
+ readMvUnlock();
+ }
+ }
+
// ALTER_PARTITION_STATES replay applies a detached snapshot here,
mirroring alterIvmInfo(). Live
// invalidation changes submit their journal from the mutating method
instead.
//
// A payload without the member carries no state at all, which is not the
same as an empty map that
// says the states are now empty: leaving them alone is the only answer
that cannot lose state.
public void alterPartitionStates(Map<String, MTMVPartitionState>
partitionStates) {
- if (partitionStates == null) {
+ replayAlterPartitionStates(partitionStates, null, false);
+ }
+
+ /**
+ * ALTER_PARTITION_STATES replay: applies the states the payload carries,
and drops the snapshots it
+ * names. Both in one lock acquisition, because a reader that saw the new
requirement while the
+ * snapshot was still there could let a transparent rewrite serve rows the
rebuild has to replace.
+ *
+ * <p>A payload without the states carries none, which is not the same as
an empty map that says the
+ * states are now empty: leaving them alone is the only answer that cannot
lose state.
+ *
+ * <p>{@code merge} says what the payload's states are. False, which is
what a payload written before the
+ * member existed means and what the invalidation channel still writes,
carries the map itself and
+ * replaces. True carries only the partitions a change touched -- see
submitPartitionStatesDelta -- so the
+ * entries it does not name belong to other records (an invalidation that
ran during the refresh, an entry
+ * alignment added) and are merged over rather than dropped.
+ */
+ public void replayAlterPartitionStates(Map<String, MTMVPartitionState>
partitionStates,
+ Set<String> removedSnapshotPartitions, boolean merge) {
+ writeMvLock();
+ try {
+ if (partitionStates != null) {
+ if (merge) {
+ for (Entry<String, MTMVPartitionState> entry :
partitionStates.entrySet()) {
+ this.partitionStates.put(entry.getKey(), new
MTMVPartitionState(entry.getValue()));
+ }
+ } else {
+ this.partitionStates =
MTMVPartitionState.copyOf(partitionStates);
+ }
+ }
+ refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
+ } finally {
+ writeMvUnlock();
+ }
+ }
+
+ /**
+ * The {@code latestEpoch} of the given MV partitions, taken under the MV
read lock.
+ *
+ * <p>This is the value a refresh has to remember: what it read from the
base tables is described by
+ * the requirement in force when it started reading, so writing that value
back as the new
+ * {@code refreshEpoch} is what keeps an invalidation arriving mid-refresh
from being swallowed. A
+ * partition without an entry is left out -- a caller writes an epoch only
for what it captured.
+ */
+ public Map<String, Long> getLatestEpochs(Set<String> partitionNames) {
+ if (CollectionUtils.isEmpty(partitionNames)) {
+ return Collections.emptyMap();
+ }
+ // Sized before the lock: the state map is what needs it, and building
the map is not part of that.
+ Map<String, Long> res =
Maps.newHashMapWithExpectedSize(partitionNames.size());
+ readMvLock();
+ try {
+ for (String partitionName : partitionNames) {
+ MTMVPartitionState state = partitionStates.get(partitionName);
+ if (state != null) {
+ res.put(partitionName, state.getLatestEpoch());
+ }
+ }
+ return res;
+ } finally {
+ readMvUnlock();
+ }
+ }
+
+ /**
+ * Brings the partition states in line with the MV's partitions: every
partition gets an entry, and
+ * every entry whose partition is gone is dropped.
+ *
+ * <p>Alignment is what makes "the partition exists" and "the entry
exists" the same thing, and it is
+ * why an invalidation cannot miss: rows are only written by a refresh,
and every refresh aligns
+ * before it reads a base table, so a partition that holds rows always has
an entry for the mark to
+ * land on. The other direction is what makes the criterion safe -- an
entry created here describes a
+ * partition with no rows yet, so requiring one generation of it discards
no requirement that was
+ * made earlier.
+ *
+ * <p>What it changes is journaled, because the entry has to be on disk
before the rows it describes
+ * can be: a crash between this and the task result would otherwise leave
a partition that holds rows
+ * with no entry at all, and every later invalidation of it would find
nothing to land on. That is the
+ * one shape in which the criterion cannot be read -- "no entry" is
supposed to mean "no rows" -- so
+ * the entry is made durable before any base table is read rather than
derived again on the next run.
+ *
+ * <p>It is deliberately not a hook on every path that creates or drops a
partition. An entry is
+ * derived state, and rebuilding it from the live partition set also
repairs whatever a crash left
+ * behind: the drop of a partition and the removal of its entry are two
journal records, and only
+ * their order -- partition first -- is safe, which leaves at most a stale
entry that the next
+ * alignment drops.
+ *
+ * <p>Only an IVM MV is aligned. For a non-IVM MV the map stays as it is,
and every reader treats
+ * "empty" and "no state" the same.
+ */
+ public void alignPartitionStates() {
+ if (!isIvm()) {
return;
}
+ EditLogItem editLogItem = null;
writeMvLock();
try {
- this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+ // Read here rather than handed in by the caller: a caller has to
read the names before it takes
+ // this lock, and a partition created in between -- by a
concurrent refresh's partition sync --
+ // would then be dropped by the retainAll below, taking with it
the state a following
+ // invalidation has to land on. The read is cheap and takes no
lock of its own, so doing it here
+ // does not add an edge to the lock order.
+ Set<String> livePartitions = Sets.newHashSet(getPartitionNames());
+ boolean changed =
partitionStates.keySet().retainAll(livePartitions);
+ for (String partitionName : livePartitions) {
+ if (!partitionStates.containsKey(partitionName)) {
+ partitionStates.put(partitionName,
MTMVPartitionState.initial());
+ changed = true;
+ }
+ }
+ if (changed) {
+ editLogItem =
submitPartitionStatesChange(Collections.emptySet());
+ }
} finally {
writeMvUnlock();
}
+ if (editLogItem != null) {
+ editLogItem.await();
+ }
}
- public void invalidateIvmBaseline() {
- EditLogItem editLogItem;
+ /**
+ * The snapshots of the partitions that are clean after this result's
epochs were applied.
+ *
+ * <p>An invalidation that reached a partition while the task ran leaves
it dirty, and its snapshot
+ * must stay gone: dropping the entry is what keeps transparent rewrite
away from rows the rebuild has
+ * to replace, and a result written back afterwards would undo exactly
that. Removing only the entry
+ * keeps the rest of the map, which the removal on the invalidation side
cannot express.
+ *
+ * <p>The caller holds the MV write lock and has already applied the
epochs, so {@code isDirty} here
+ * reads the state the data is actually described by.
+ *
+ * <p>Only an IVM MV has partition states, so only its write-back is
narrowed here: every entry of a
+ * non-IVM MV has no state to be dirty in and is written back as it always
was.
+ */
+ private Map<String, MTMVRefreshPartitionSnapshot>
snapshotsOfCleanPartitions(
+ Map<String, MTMVRefreshPartitionSnapshot> snapshots) {
+ if (MapUtils.isEmpty(snapshots)) {
+ return snapshots;
+ }
+ Map<String, MTMVRefreshPartitionSnapshot> res =
Maps.newHashMapWithExpectedSize(snapshots.size());
+ for (Entry<String, MTMVRefreshPartitionSnapshot> entry :
snapshots.entrySet()) {
+ MTMVPartitionState state = partitionStates.get(entry.getKey());
+ // No entry means the partition was created after the alignment,
so it can only hold rows this
+ // task wrote; a dirty one needs its rebuild before anything may
read it through the MV.
+ if (state == null || !state.isDirty()) {
+ res.put(entry.getKey(), entry.getValue());
+ }
+ }
+ return res;
+ }
+
+ /**
+ * Records the epochs the given partitions were read at, which is how a
refresh turns a requirement
+ * into the state of the data.
+ *
+ * <p>Only {@code refreshEpoch} is written: a refresh writes back the
requirement it captured, and the
+ * requirement may have been raised again since that capture. A payload
built from the captured map
+ * would overwrite the newer value and lose the rebuild it asks for, so
{@code latestEpoch} is left
+ * alone here.
+ *
+ * <p>The caller holds the MV write lock (it is applied together with the
rest of a task result).
+ */
+ private void applyRefreshedEpochs(Map<String, Long> capturedEpochs) {
+ if (MapUtils.isEmpty(capturedEpochs)) {
+ return;
+ }
+ for (Entry<String, Long> entry : capturedEpochs.entrySet()) {
+ MTMVPartitionState state = partitionStates.get(entry.getKey());
+ if (state == null) {
+ // The partition was dropped while the task ran, so its state
went with it.
+ continue;
+ }
+ state.setRefreshEpoch(entry.getValue());
Review Comment:
Right. The captures and snapshots a task holds are keyed by partition name,
and a retry's partition sync can drop a partition and add one of the same name
back -- which the alignment gives `{0, 1}`, empty. Writing the old captures
back credits the new partition with what the old one held, and the failure mode
is not only the wrong number: a partition clean at an epoch that a later change
only raises to is one no refresh rebuilds, so the rows the recreation removed
are published as current.
`81fedd86e27` drops both, in the hook that already re-reads the partitions
after a retry (`adoptPartitionsCreatedByTheRetry`): a live partition whose
state is the one an entry starts with -- `MTMVPartitionState.initial()` -- is
one the alignment just created, so the capture and the snapshot this task holds
for that name go. The snapshot has to go with the capture: leaving it would let
the sync judgement call the recreated, empty partition caught up even with its
epoch no longer written back.
`MTMVTaskTest` pins it: with `p1` recreated and `p2` the partition the
rebuild replaced, `p1`'s capture and snapshot are dropped, `p2`'s stand, and
both are dirty.
What I did not build is the end-to-end case you describe -- two retries with
a same-name re-add. It needs repeatedly inducing a "missing MV partition" delta
failure by dropping and re-adding populated base partitions, and the mechanism
is small enough that I pinned it at the unit level instead; if you want it as a
suite, say so and I will add one on its own rather than growing a suite that
already drives several shapes. The reporting side is left as it was on purpose:
`CompletedPartitions` still names the partition this task refreshed under that
name, which is what it did; the epochs are what must not be inherited.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]