Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,19 @@ public interface AbstractInsertOverwriteManager {

void taskGroupSuccess(long groupId, OlapTable targetTable, boolean forceDropPartition) throws DdlException;

/**
* The replacement half of {@link #taskGroupSuccess}: makes every replacement the group recorded visible.
* A caller that publishes the group itself calls this and then {@link #finishTaskGroup}, so that the task
* bookkeeping -- an edit-log write per task, each waiting for its journal -- runs outside whatever lock the
* caller holds the replacement under. A manager whose owning frontend runs the whole sequence behind one
* call replaces there and finishes there.
*/
void replacePartitionsOfTaskGroup(long groupId, OlapTable targetTable, boolean forceDropPartition)
throws DdlException;

/** The bookkeeping half of {@link #taskGroupSuccess}, after the replacement it follows. */
void finishTaskGroup(long groupId);

void taskSuccess(long taskId) throws Exception;

void taskGroupFail(long groupId) throws Exception;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,21 @@ public void taskGroupFail(long groupId) {
// here we will make all raplacement of this group visiable. if someone fails, nothing happen.
@Override
public void taskGroupSuccess(long groupId, OlapTable targetTable, boolean forceDropPartition) throws DdlException {
replacePartitionsOfTaskGroup(groupId, targetTable, forceDropPartition);
finishTaskGroup(groupId);
}

/**
* The replacement half of {@link #taskGroupSuccess}: makes every replacement this group recorded visible.
*
* <p>Separate from the bookkeeping that follows it because the caller holds the target table's write lock
* for the last cancellation check before it publishes: the bookkeeping writes an edit-log entry per task,
* each of which waits for its journal to be acknowledged, and readers and writers of the target must not
* stay blocked through those waits -- the utility's own replacement, which the lock is about, does not.
*/
@Override
public void replacePartitionsOfTaskGroup(long groupId, OlapTable targetTable, boolean forceDropPartition)
throws DdlException {
try {
Map<Long, Long> relations = partitionPairs.get(groupId);
ArrayList<String> oldNames = new ArrayList<>();
Expand All @@ -205,6 +220,16 @@ public void taskGroupSuccess(long groupId, OlapTable targetTable, boolean forceD
+ "all new partition will not be visible and will be recycled by partition GC.");
throw e;
}
}

/**
* The bookkeeping half of {@link #taskGroupSuccess}: makes the group's tasks final and forgets the group,
* after the replacement has been made. See
* {@link #replacePartitionsOfTaskGroup} for why it is not the replacement's caller to hold this under a
* table lock.
*/
@Override
public void finishTaskGroup(long groupId) {
LOG.info("insert overwrite auto detect partition task group [" + groupId + "] succeed");
for (Long taskId : taskGroups.get(groupId)) {
Env.getCurrentEnv().getEditLog()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,19 @@ public void taskGroupSuccess(long groupId, OlapTable targetTable, boolean forceD
targetTable.getDatabase().getFullName(), targetTable.getName(), groupId, forceDropPartition);
}

@Override
public void replacePartitionsOfTaskGroup(long groupId, OlapTable targetTable, boolean forceDropPartition)
throws DdlException {
// The owning frontend runs the replacement and the bookkeeping behind this one call, so the caller
// does not hold its lock across them; finishTaskGroup has nothing left to do here.
taskGroupSuccess(groupId, targetTable, forceDropPartition);
}

@Override
public void finishTaskGroup(long groupId) {
// Finished on the owning frontend, which the replacement call reached.
}

@Override
public void taskSuccess(long taskId) throws Exception {
catalog.getFeServiceClient().taskSuccess(taskId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -355,6 +355,17 @@ public boolean requiresTransaction() {
return !emptyInsert || !streamUpdateInfos.isEmpty();
}

/**
* Records on the insert context that what this insert wrote is durable: its transaction has committed,
* whether or not its publication has finished. Each executor calls this where its own commit happens,
* because the response the client gets is a different question -- a publication timeout that follows a
* commit is reported as an error while the rows are committed. See
* {@link InsertCommandContext#setCommitted}.
*/
protected void markCommitted() {
insertCtx.ifPresent(insertCommandContext -> insertCommandContext.setCommitted(true));
}

public void setStreamUpdateInfos(List<TableStreamUpdateInfo> streamUpdateInfos) {
this.streamUpdateInfos = streamUpdateInfos;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,7 @@ protected void onComplete() throws UserException {
transactionManager.commit(txnId);
}
txnStatus = TransactionStatus.COMMITTED;
markCommitted();
long t2 = System.currentTimeMillis();

try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,25 @@
*/
public abstract class InsertCommandContext {

/**
* Set when the insert's transaction has committed: what it wrote -- its rows, and any table stream offset
* update it carried -- is durable. That is not the same as the insert having succeeded. A load whose
* publication times out after its commit is committed, while the session's visibility-timeout mode reports
* it to the client as an error; and a load whose plan folded to an empty relation commits nothing at all,
* since it runs no transaction unless it has an offset to commit.
*
* <p>Whoever owns a transaction records it where that commit happens; see
* {@code AbstractInsertExecutor#markCommitted()}. A caller that owns what happens next -- an overwrite
* publishing its temporary partitions -- decides by it: what is durable has to be published, and where
* nothing is durable a failure or a cancellation still has everything to take back.
*/
private boolean committed = false;

public boolean hasCommitted() {
return committed;
}

public void setCommitted(boolean committed) {
this.committed = committed;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -756,6 +756,10 @@ private void runInternal(ConnectContext ctx, StmtExecutor executor) throws Excep
AbstractInsertExecutor insertExecutor = initPlan(ctx, executor);
// An empty Table Stream read still needs to commit its offset update atomically.
if (!insertExecutor.requiresTransaction()) {
// Nothing is committed on this path, so a caller that asks whether the rows are durable is told
// they are not -- no row and no offset was written, which is what it needs before treating a
// cancellation or a later failure here as too late to take back.
// See InsertCommandContext#setCommitted.
return;
}
if (insertExecutorListener != null) {
Expand Down
Loading
Loading