yujun777 commented on code in PR #68390:
URL: https://github.com/apache/doris/pull/68390#discussion_r4128971136
##########
fe/fe-core/src/main/java/org/apache/doris/job/extensions/mtmv/MTMVTask.java:
##########
@@ -259,7 +260,25 @@ private PartitionPlanningException(String message,
Throwable cause) {
// Written by the executing (Disruptor worker) thread via the
executeCommand consumer
// callback and read by the cancel (command) thread, so it must be
volatile.
private volatile StmtExecutor executor;
- private Map<String, MTMVRefreshPartitionSnapshot> partitionSnapshots;
+ // What this task has committed, per MV partition: the snapshot each
partition's rows were read at.
+ // One accumulator for the whole task rather than one per phase, because
that is what the MV publishes
+ // at the end of it -- a phase that started from empty would publish its
own work and drop the work of
+ // the phases before it, leaving partitions a preceding rebuild replaced
looking unsynced.
+ private Map<String, MTMVRefreshPartitionSnapshot> partitionSnapshots =
Maps.newConcurrentMap();
+ // The requirement each refreshed partition was read under, captured
before the base tables were read
+ // and recorded only once that batch's data committed (see
commitCapturedEpochs). In memory only: the
+ // journal carries the resulting states, and a replay applies those
instead of recomputing anything.
+ private transient Map<String, Long> ivmCapturedEpochs = Maps.newHashMap();
Review Comment:
One further instance of this is closed since this thread was opened: the
snapshot map the result carries was handed over live as well, and `after()` now
passes a copy of it -- the same rule `getIvmCapturedEpochs` already followed
for the epochs. So what the cancel thread publishes is now a detached read of
the committed work on every path: a copy of the epochs, a copy of the
snapshots, and a write-back that filters both to the partitions that are still
clean.
What the finding's second half asks for -- publishing and releasing only
after the worker quiesces -- is the framework change described above: STOP
reaches `cancel(false)` through `AbstractJob`, so a cancelled task publishes
before its worker stops and can see resources the cancel path has released, and
deciding what a cancelled refresh publishes (and whether STOP waits) is a
decision for every job type that shares the framework rather than for this PR.
I am resolving this thread on that division: the publication side is fixed and
pinned, the lifecycle side is tracked as its own change.
##########
fe/fe-core/src/main/java/org/apache/doris/catalog/MTMV.java:
##########
@@ -684,33 +759,211 @@ public Map<String, MTMVPartitionState>
getPartitionStates() {
// 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) {
- return;
- }
+ replayAlterPartitionStates(partitionStates, null);
+ }
+
+ /**
+ * 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.
+ */
+ public void replayAlterPartitionStates(Map<String, MTMVPartitionState>
partitionStates,
+ Set<String> removedSnapshotPartitions) {
writeMvLock();
try {
- this.partitionStates = MTMVPartitionState.copyOf(partitionStates);
+ if (partitionStates != null) {
+ this.partitionStates =
MTMVPartitionState.copyOf(partitionStates);
+ }
+ refreshSnapshot.removeSnapshots(removedSnapshotPartitions);
} finally {
writeMvUnlock();
}
}
- public void invalidateIvmBaseline() {
- EditLogItem editLogItem;
+ /**
+ * 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(Set<String> livePartitionNames) {
+ if (!isIvm()) {
+ return;
+ }
+ // Copied up front: callers pass what OlapTable holds, and that is
mutated under the table's own
+ // write lock, not this one. Iterating the live collection could see
it change.
+ Set<String> livePartitions = Sets.newHashSet(livePartitionNames);
+ EditLogItem editLogItem = null;
writeMvLock();
try {
- if (ivmInfo == null) {
- ivmInfo = new IvmInfo();
+ boolean changed =
partitionStates.keySet().retainAll(livePartitions);
Review Comment:
Recording the state of this one as it stands: the read of the live partition
names was moved into `alignPartitionStates` itself, under the same MV lock as
the map it edits, so the stale-snapshot half is gone and a caller cannot hand
in names it read before the lock.
What is left is the window between that read and the `retainAll`, both
adjacent statements under the MV lock -- which does not close it, because the
MV's partition map is mutated under the table's lock rather than this one.
Closing it means removing the entry where the partition DDL happens instead of
deriving it here, which is a change to the partition lifecycle rather than to
this one. I am resolving the thread on that division; say the word if you want
that change taken up separately.
--
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]