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 @@ -946,7 +946,7 @@ private static boolean tryAcquireLocks(
batchWaitMs = -1L;
}

if (!batch.getKey().lockTxEntries(batch.getValue().values(), batchWaitMs)) {
if (batch.getKey().lockTxEntries(batch.getValue().values(), batchWaitMs).containsValue(false)) {
locked = false;
break;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import java.util.HashMap;
import java.util.HashSet;
import java.util.Iterator;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.NoSuchElementException;
Expand Down Expand Up @@ -124,6 +125,7 @@
import org.apache.ignite.internal.util.future.GridFinishedFuture;
import org.apache.ignite.internal.util.future.GridFutureAdapter;
import org.apache.ignite.internal.util.lang.GridCloseableIterator;
import org.apache.ignite.internal.util.lang.GridClosureException;
import org.apache.ignite.internal.util.lang.GridPlainCallable;
import org.apache.ignite.internal.util.lang.GridPlainRunnable;
import org.apache.ignite.internal.util.tostring.GridToStringExclude;
Expand Down Expand Up @@ -165,6 +167,7 @@
import static org.apache.ignite.internal.processors.dr.GridDrType.DR_LOAD;
import static org.apache.ignite.internal.processors.dr.GridDrType.DR_NONE;
import static org.apache.ignite.internal.processors.metric.impl.MetricUtils.cacheMetricsRegistryName;
import static org.apache.ignite.internal.processors.rollingupgrade.feature.CoreFeatureRegistry.VERSIONED_TX_LOCK_FEATURE;
import static org.apache.ignite.internal.processors.task.TaskExecutionOptions.options;
import static org.apache.ignite.internal.thread.pool.IgniteThreadPoolExecutor.newFixedThreadPool;
import static org.apache.ignite.transactions.TransactionConcurrency.OPTIMISTIC;
Expand Down Expand Up @@ -3040,22 +3043,29 @@ public CacheMetricsImpl metrics0() {
}

/** {@inheritDoc} */
@Override public boolean lockTxEntries(Collection<CacheEntry<K, V>> entries, long waitTimeout)
throws IgniteCheckedException {
A.notNull(entries, "entries");

@Override public Map<CacheEntry<K, V>, Boolean> lockTxEntries(
Collection<CacheEntry<K, V>> entries,
long waitTimeout
) throws IgniteCheckedException {
return lockTxEntriesAsync(entries, waitTimeout).get();
}

/** {@inheritDoc} */
@Override public IgniteInternalFuture<Boolean> lockTxEntryAsync(CacheEntry<K, V> entry, long waitTimeout) {
A.notNull(entry, "entry");

return lockTxEntriesAsync(Collections.singleton(entry), waitTimeout);
return lockTxEntriesAsync(Collections.singleton(entry), waitTimeout).chain(fut -> {
try {
return fut.get().get(entry);
}
catch (IgniteCheckedException e) {
throw new GridClosureException(e);
}
});
}

/** {@inheritDoc} */
@Override public IgniteInternalFuture<Boolean> lockTxEntriesAsync(
@Override public IgniteInternalFuture<Map<CacheEntry<K, V>, Boolean>> lockTxEntriesAsync(
Collection<CacheEntry<K, V>> entries,
long waitTimeout
) {
Expand All @@ -3071,193 +3081,151 @@ public CacheMetricsImpl metrics0() {
return new GridFinishedFuture<>(
new IgniteCheckedException("Failed to acquire transactional lock in optimistic transaction."));

// Wait for previous per-transaction async operations to finish.
if (!ctx.kernalContext().rollingUpgrade().features().isActive(VERSIONED_TX_LOCK_FEATURE))
return new GridFinishedFuture<>(new IgniteCheckedException(
"Primary-side versioned transactional locks require rolling upgrade to be finalized."));

tx.txState().awaitLastFuture();

if (!tx.init())
return new GridFinishedFuture<>(new IgniteTxRollbackCheckedException(
"Failed to acquire transactional lock because transaction has been completed: " + tx));

if (entries.isEmpty())
return new GridFinishedFuture<>(true);
List<CacheEntry<K, V>> inputs = new ArrayList<>(entries);
Map<CacheEntry<K, V>, Boolean> results = new HashMap<>();
Map<KeyCacheObject, IgniteTxEntry> enlisted = new LinkedHashMap<>();
Map<KeyCacheObject, IgniteTxEntry> snapshots = new HashMap<>();
CacheOperationContext opCtx = ctx.operationContextPerCall();

for (CacheEntry<K, V> entry : inputs)
A.notNull(entry, "entry");

try {
tx.addActiveCache(ctx, false);
}
catch (IgniteCheckedException e) {
return new GridFinishedFuture<>(e);
}

Collection<KeyCacheObject> keys = new ArrayList<>(entries.size());
List<IgniteTxEntry> txEntries = new ArrayList<>(entries.size());
List<GridCacheVersion> expVers = new ArrayList<>(entries.size());
Set<IgniteTxKey> txKeys = new HashSet<>(entries.size());
for (CacheEntry<K, V> entry : inputs) {
KeyCacheObject key = ctx.toCacheKeyObject(entry.getKey());

CacheOperationContext opCtx = ctx.operationContextPerCall();
if (enlisted.containsKey(key)) {
throw new IgniteCheckedException("Failed to acquire transactional lock because " +
"entry is duplicated [key=" + key + ", entry=" + entry + ']');
}

for (CacheEntry<K, V> entry : entries) {
A.notNull(entry, "entry");
IgniteTxEntry oldEntry = tx.entry(ctx.txKey(key));

KeyCacheObject key = ctx.toCacheKeyObject(entry.getKey());
IgniteTxKey txKey = ctx.txKey(key);
if (oldEntry != null && (oldEntry.op() != READ || oldEntry.locked())) {
results.put(entry, true);

if (!txKeys.add(txKey))
continue;
continue;
}

IgniteTxEntry lockedTxEntry = tx.entry(txKey);
if (!(entry.version() instanceof GridCacheVersion)) {
throw new IgniteCheckedException(
"Failed to acquire transactional lock for entry with unsupported version type: " + entry);
}

if (lockedTxEntry != null && (lockedTxEntry.op() != READ || lockedTxEntry.locked()))
continue;
snapshots.put(key, oldEntry == null ? null : oldEntry.copy());

if (!(entry.version() instanceof GridCacheVersion)) {
tx.removeAndUnlockTxEntries(txEntries);
GridCacheEntryEx cached = ctx.isColocated()
? ctx.colocated().entryExx(key, tx.topologyVersion(), true) : entryEx(key);
IgniteTxEntry txEntry = tx.addEntry(
READ,
ctx.toCacheObject(entry.getValue()),
null,
null,
cached,
null,
null,
true,
-1L,
-1L,
null,
opCtx != null && opCtx.skipStore(),
opCtx != null && opCtx.skipReadThrough(),
opCtx != null && opCtx.keepBinaryInInterceptor(),
opCtx != null && opCtx.isKeepBinary(),
CU.isNearEnabled(ctx)
);

return new GridFinishedFuture<>(new IgniteCheckedException("Failed to acquire transactional lock for entry with " +
"unsupported version type [entry=" + entry + ", version=" + entry.version() + ']'));
txEntry.expectedLockVersion((GridCacheVersion)entry.version());
txEntry.versionedLockResult(null);
enlisted.put(key, txEntry);
}
}
catch (IgniteCheckedException e) {
for (Map.Entry<KeyCacheObject, IgniteTxEntry> entry : enlisted.entrySet())
tx.removeFailedLockEntry(entry.getValue(), snapshots.get(entry.getKey()));

CacheObject val = ctx.toCacheObject(entry.getValue());
GridCacheEntryEx entryEx = ctx.isColocated() ? ctx.colocated().entryExx(key, tx.topologyVersion(), true) : entryEx(key);
return new GridFinishedFuture<>(e);
}

IgniteTxEntry txEntry = tx.addEntry(
READ,
val,
null,
null,
entryEx,
null,
null,
true,
-1L,
-1L,
null,
opCtx != null && opCtx.skipStore(),
opCtx != null && opCtx.skipReadThrough(),
opCtx != null && opCtx.keepBinaryInInterceptor(),
opCtx != null && opCtx.isKeepBinary(),
CU.isNearEnabled(ctx)
);
if (enlisted.isEmpty())
return new GridFinishedFuture<>(results);

keys.add(key);
txEntries.add(txEntry);
expVers.add((GridCacheVersion)entry.version());
}
GridFutureAdapter<Map<CacheEntry<K, V>, Boolean>> res = new GridFutureAdapter<>();
GridCacheAdapter.FutureHolder holder = tx.txState().lastAsyncFuture();

if (keys.isEmpty())
return new GridFinishedFuture<>(true);
if (holder != null) {
holder.lock();

// Acquire transactional lock future from concrete cache implementation. Use txLockAsync which
// delegates to cache-specific lockAllAsync implementations for distributed caches.
long lockWaitStartTime = U.currentTimeMillis();
long timeout = tx.remainingTime();
long effectiveWaitTimeout = CU.isWaitTimeoutExpiresFirst(waitTimeout, timeout) ? waitTimeout : timeout;
long lockWaitEndTime = effectiveWaitTimeout > 0
? lockWaitStartTime + effectiveWaitTimeout
: effectiveWaitTimeout < 0 ? lockWaitStartTime : 0L;
try {
holder.saveFuture(res);
}
finally {
holder.unlock();
}
}

IgniteInternalFuture<Boolean> lockFut = txLockAsync(keys,
timeout,
// Keep the existing distributed mapping and lock order: one request per consecutive primary group,
// rather than one request per entry. Individual outcomes travel in the primary's batch response.
IgniteInternalFuture<Boolean> lockFut = txLockAsync(
enlisted.keySet(),
tx.remainingTime(),
waitTimeout,
tx,
/*isRead*/true,
/*retval*/false,
true,
false,
tx.isolation(),
/*invalidate*/false,
/*createTtl*/0L,
/*accessTtl*/0L);

IgniteInternalFuture<Boolean> res = new GridEmbeddedFuture<>(
lockFut,
(locked, ex) -> {
if (ex != null)
return new GridFinishedFuture<>(ex);

if (!locked) {
tx.removeAndUnlockTxEntries(txEntries);

return new GridFinishedFuture<>(false);
}

try {
for (int i = 0; i < txEntries.size(); i++) {
IgniteTxEntry txEntry = txEntries.get(i);
EntryGetResult getRes;
int retryCnt = 0;

while (true) {
try {
GridCacheEntryEx cached = txEntry.cached();

getRes = cached.innerGetVersioned(
null,
tx,
/*update-metrics*/false,
/*event*/false,
null,
tx.resolveTaskName(),
null,
false,
null);

break;
}
catch (GridCacheEntryRemovedException ignored) {
// NOWAIT still permits one immediate retry because renewing an obsolete entry does not
// wait for a lock. Other modes stop when their effective timeout expires.
boolean allowNowaitRetry = waitTimeout < 0 && retryCnt == 0;

if (!allowNowaitRetry && lockWaitEndTime != 0
&& U.currentTimeMillis() >= lockWaitEndTime) {
getRes = null;
false,
0L,
0L
);

break;
}
lockFut.listen(fut -> {
try {
boolean completed = fut.get();

if (log.isDebugEnabled())
log.debug("Got removed exception in lockTxEntries postLock (will retry): "
+ txEntry.cached());
if (tx.remainingTime() == -1)
throw tx.timeoutException();

KeyCacheObject key = txEntry.key();
GridCacheEntryEx cached = ctx.isColocated()
? ctx.colocated().entryExx(key, tx.topologyVersion(), true)
: entryEx(key, tx.topologyVersion());
if (tx.isRollbackOnly())
throw tx.rollbackException();

txEntry.cached(cached);
retryCnt++;
}
}
for (CacheEntry<K, V> entry : inputs) {
if (!results.containsKey(entry)) {
IgniteTxEntry txEntry = enlisted.get(ctx.toCacheKeyObject(entry.getKey()));

if (getRes == null || !expVers.get(i).equals(getRes.version())) {
tx.removeAndUnlockTxEntries(txEntries);
assert !completed || txEntry.versionedLockResult() != null;

return new GridFinishedFuture<>(false);
}
results.put(entry, completed && Boolean.TRUE.equals(txEntry.versionedLockResult()));
}

return new GridFinishedFuture<>(true);
}
catch (IgniteCheckedException e) {
tx.removeAndUnlockTxEntries(txEntries);

return new GridFinishedFuture<>(e);
}
}
);

// Register this future in transaction's async-holder so that subsequent operations
// that call tx.txState().awaitLastFuture() will wait for it.
GridCacheAdapter.FutureHolder holder = tx.txState().lastAsyncFuture();
for (Map.Entry<KeyCacheObject, IgniteTxEntry> entry : enlisted.entrySet()) {
IgniteTxEntry txEntry = entry.getValue();

if (holder != null) {
holder.lock();
if (completed && Boolean.TRUE.equals(txEntry.versionedLockResult()))
txEntry.markLocked();
else
tx.removeFailedLockEntry(txEntry, snapshots.get(entry.getKey()));
}

try {
holder.saveFuture(res);
res.onDone(results);
}
finally {
holder.unlock();
catch (IgniteCheckedException | RuntimeException e) {
res.onDone(e);
}
}
});

return res;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1336,7 +1336,7 @@ public IgniteInternalCache<K, V> delegate() {
}

/** {@inheritDoc} */
@Override public boolean lockTxEntries(Collection<CacheEntry<K, V>> entries, long waitTimeout)
@Override public Map<CacheEntry<K, V>, Boolean> lockTxEntries(Collection<CacheEntry<K, V>> entries, long waitTimeout)
throws IgniteCheckedException {
CacheOperationContext prev = gate.enter(opCtx);

Expand All @@ -1349,7 +1349,7 @@ public IgniteInternalCache<K, V> delegate() {
}

/** {@inheritDoc} */
@Override public IgniteInternalFuture<Boolean> lockTxEntriesAsync(
@Override public IgniteInternalFuture<Map<CacheEntry<K, V>, Boolean>> lockTxEntriesAsync(
Collection<CacheEntry<K, V>> entries,
long waitTimeout
) {
Expand Down
Loading
Loading