- private void onCreateTransactionComplete(Throwable failure, Object response) {
- if(failure instanceof NoShardLeaderException) {
- // There's no leader for the shard yet - schedule and try again, unless we're out
- // of retries. Note: createTxTries is volatile as it may be written by different
- // threads however not concurrently, therefore decrementing it non-atomically here
- // is ok.
- if(--createTxTries > 0) {
- LOG.debug("Tx {} Shard {} has no leader yet - scheduling create Tx retry",
- getIdentifier(), shardName);
-
- getActorContext().getActorSystem().scheduler().scheduleOnce(CREATE_TX_TRY_INTERVAL,
- new Runnable() {
- @Override
- public void run() {
- tryCreateTransaction();
- }
- }, getActorContext().getClientDispatcher());
- return;
+ private void tryFindPrimaryShard() {
+ LOG.debug("Tx {} Retrying findPrimaryShardAsync for shard {}", getIdentifier(), shardName);
+
+ this.primaryShardInfo = null;
+ Future<PrimaryShardInfo> findPrimaryFuture = getActorContext().findPrimaryShardAsync(shardName);
+ findPrimaryFuture.onComplete(new OnComplete<PrimaryShardInfo>() {
+ @Override
+ public void onComplete(final Throwable failure, final PrimaryShardInfo newPrimaryShardInfo) {
+ onFindPrimaryShardComplete(failure, newPrimaryShardInfo);
+ }
+ }, getActorContext().getClientDispatcher());
+ }
+
+ private void onFindPrimaryShardComplete(final Throwable failure, final PrimaryShardInfo newPrimaryShardInfo) {
+ if (failure == null) {
+ this.primaryShardInfo = newPrimaryShardInfo;
+ tryCreateTransaction();
+ } else {
+ LOG.debug("Tx {}: Find primary for shard {} failed", getIdentifier(), shardName, failure);
+
+ onCreateTransactionComplete(failure, null);
+ }
+ }
+
+ private void onCreateTransactionComplete(final Throwable failure, final Object response) {
+ // An AskTimeoutException will occur if the local shard forwards to an unavailable remote leader or
+ // the cached remote leader actor is no longer available.
+ boolean retryCreateTransaction = primaryShardInfo != null
+ && (failure instanceof NoShardLeaderException || failure instanceof AskTimeoutException);
+
+ // Schedule a retry unless we're out of retries. Note: totalCreateTxTimeout is volatile as it may
+ // be written by different threads however not concurrently, therefore decrementing it
+ // non-atomically here is ok.
+ if (retryCreateTransaction && totalCreateTxTimeout > 0) {
+ long scheduleInterval = CREATE_TX_TRY_INTERVAL_IN_MS;
+ if (failure instanceof AskTimeoutException) {
+ // Since we use the createTxMessageTimeout for the CreateTransaction request and it timed
+ // out, subtract it from the total timeout. Also since the createTxMessageTimeout period
+ // has already elapsed, we can immediately schedule the retry (10 ms is virtually immediate).
+ totalCreateTxTimeout -= createTxMessageTimeout.duration().toMillis();
+ scheduleInterval = 10;