diff --git a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/AbstractCdcEventsApplier.java b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/AbstractCdcEventsApplier.java index ffef91114..479a1db4a 100644 --- a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/AbstractCdcEventsApplier.java +++ b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/AbstractCdcEventsApplier.java @@ -87,7 +87,7 @@ public int apply(Iterable evts) throws IgniteCheckedException { if (evt.value() != null) { evtsApplied += applyIf(currCacheId, () -> isApplyBatch(updBatch, key), hasRemoves); - updBatch.put(key, toValue(currCacheId, evt.value(), ver)); + updBatch.put(key, toValue(currCacheId, evt, ver)); } else { evtsApplied += applyIf(currCacheId, hasUpdates, () -> isApplyBatch(rmvBatch, key)); @@ -152,7 +152,7 @@ private boolean isApplyBatch(Map map, K key) { protected abstract K toKey(CdcEvent evt); /** @return Value. */ - protected abstract V toValue(int cacheId, Object val, GridCacheVersion ver); + protected abstract V toValue(int cacheId, CdcEvent evt, GridCacheVersion ver); /** Stores DR data. */ protected abstract void putAllConflict(int cacheId, Map drMap); diff --git a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/CdcEventsIgniteApplier.java b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/CdcEventsIgniteApplier.java index bce50d38d..9afda275b 100644 --- a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/CdcEventsIgniteApplier.java +++ b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/CdcEventsIgniteApplier.java @@ -35,8 +35,8 @@ import org.apache.ignite.internal.util.typedef.internal.CU; import org.apache.ignite.internal.util.typedef.internal.U; -import static org.apache.ignite.internal.processors.cache.GridCacheUtils.EXPIRE_TIME_CALCULATE; -import static org.apache.ignite.internal.processors.cache.GridCacheUtils.TTL_NOT_CHANGED; +import static org.apache.ignite.internal.processors.cache.GridCacheUtils.EXPIRE_TIME_ETERNAL; +import static org.apache.ignite.internal.processors.cache.GridCacheUtils.TTL_ETERNAL; /** * Contains logic to process {@link CdcEvent} and apply them to the destination cluster. @@ -93,16 +93,18 @@ public CdcEventsIgniteApplier(IgniteEx ignite, int maxBatchSize, IgniteLogger lo } /** {@inheritDoc} */ - @Override protected GridCacheDrInfo toValue(int cacheId, Object val, GridCacheVersion ver) { + @Override protected GridCacheDrInfo toValue(int cacheId, CdcEvent evt, GridCacheVersion ver) { CacheObject cacheObj; + Object val = evt.value(); + if (val instanceof CacheObject) cacheObj = (CacheObject)val; else cacheObj = new CacheObjectImpl(val, null); - return cache(cacheId).configuration().getExpiryPolicyFactory() != null ? - new GridCacheDrExpirationInfo(cacheObj, ver, TTL_NOT_CHANGED, EXPIRE_TIME_CALCULATE) : + return evt.expireTime() != EXPIRE_TIME_ETERNAL ? + new GridCacheDrExpirationInfo(cacheObj, ver, TTL_ETERNAL, evt.expireTime()) : new GridCacheDrInfo(cacheObj, ver); } diff --git a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/CacheVersionConflictResolverImpl.java b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/CacheVersionConflictResolverImpl.java index 56e764fd5..ce1a17fce 100644 --- a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/CacheVersionConflictResolverImpl.java +++ b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/CacheVersionConflictResolverImpl.java @@ -24,6 +24,7 @@ import org.apache.ignite.internal.processors.cache.version.GridCacheVersionConflictContext; import org.apache.ignite.internal.processors.cache.version.GridCacheVersionedEntryEx; import org.apache.ignite.internal.util.tostring.GridToStringInclude; +import org.apache.ignite.internal.util.typedef.internal.CU; import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; @@ -88,7 +89,30 @@ public CacheVersionConflictResolverImpl(byte clusterId, String conflictResolveFi ) { GridCacheVersionConflictContext res = new GridCacheVersionConflictContext<>(ctx, oldEntry, newEntry); - if (isUseNew(ctx, oldEntry, newEntry)) + boolean expireExists = oldEntry.ttl() != CU.TTL_ETERNAL + || newEntry.ttl() != CU.TTL_ETERNAL + || oldEntry.expireTime() != CU.EXPIRE_TIME_ETERNAL + || newEntry.expireTime() != CU.EXPIRE_TIME_ETERNAL; + + boolean useNew = isUseNew(ctx, oldEntry, newEntry); + + if (expireExists) { + if (newEntry.expireTime() > oldEntry.expireTime()) { + res.merge( + useNew ? newEntry.value(ctx) : oldEntry.value(ctx), + newEntry.ttl(), + newEntry.expireTime() + ); + } + else { + res.merge( + useNew ? newEntry.value(ctx) : oldEntry.value(ctx), + oldEntry.ttl(), + oldEntry.expireTime() + ); + } + } + else if (useNew) res.useNew(); else res.useOld(); diff --git a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/DebugCacheVersionConflictResolverImpl.java b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/DebugCacheVersionConflictResolverImpl.java index cb3ce3210..000cbf378 100644 --- a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/DebugCacheVersionConflictResolverImpl.java +++ b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/conflictresolve/DebugCacheVersionConflictResolverImpl.java @@ -54,6 +54,8 @@ public DebugCacheVersionConflictResolverImpl(byte clusterId, String conflictReso "start=" + oldEntry.isStartVersion() + ", oldVer=" + oldEntry.version() + ", newVer=" + newEntry.version() + + ", oldExpire=[" + oldEntry.ttl() + "," + oldEntry.expireTime() + ']' + + ", newExpire=[" + newEntry.ttl() + "," + newEntry.expireTime() + ']' + ", old=" + oldVal + ", new=" + newVal + ", res=" + res + ']'); diff --git a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/thin/CdcEventsIgniteClientApplier.java b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/thin/CdcEventsIgniteClientApplier.java index 5484277a7..2a9912131 100644 --- a/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/thin/CdcEventsIgniteClientApplier.java +++ b/modules/cdc-ext/src/main/java/org/apache/ignite/cdc/thin/CdcEventsIgniteClientApplier.java @@ -26,7 +26,7 @@ import org.apache.ignite.internal.processors.cache.version.GridCacheVersion; import org.apache.ignite.internal.util.collection.IntHashMap; import org.apache.ignite.internal.util.collection.IntMap; -import org.apache.ignite.internal.util.typedef.T2; +import org.apache.ignite.internal.util.typedef.T3; import org.apache.ignite.internal.util.typedef.internal.CU; /** @@ -35,7 +35,7 @@ * @see TcpClientCache#putAllConflict(Map) * @see TcpClientCache#removeAllConflict(Map) */ -public class CdcEventsIgniteClientApplier extends AbstractCdcEventsApplier> { +public class CdcEventsIgniteClientApplier extends AbstractCdcEventsApplier> { /** Client connected to the destination cluster. */ private final IgniteClient client; @@ -59,12 +59,12 @@ public CdcEventsIgniteClientApplier(IgniteClient client, int maxBatchSize, Ignit } /** {@inheritDoc} */ - @Override protected T2 toValue(int cacheId, Object val, GridCacheVersion ver) { - return new T2<>(val, ver); + @Override protected T3 toValue(int cacheId, CdcEvent evt, GridCacheVersion ver) { + return new T3<>(evt.value(), ver, evt.expireTime()); } /** {@inheritDoc} */ - @Override protected void putAllConflict(int cacheId, Map> drMap) { + @Override protected void putAllConflict(int cacheId, Map> drMap) { cache(cacheId).putAllConflict(drMap); } diff --git a/modules/cdc-ext/src/test/java/org/apache/ignite/cdc/AbstractReplicationTest.java b/modules/cdc-ext/src/test/java/org/apache/ignite/cdc/AbstractReplicationTest.java index 144b8b872..02d394d84 100644 --- a/modules/cdc-ext/src/test/java/org/apache/ignite/cdc/AbstractReplicationTest.java +++ b/modules/cdc-ext/src/test/java/org/apache/ignite/cdc/AbstractReplicationTest.java @@ -446,7 +446,7 @@ public void testActiveActiveReplication() throws Exception { /** Test that destination cluster applies expiration policy on received entries. */ @Test public void testWithExpiryPolicy() throws Exception { - Factory factory = () -> new CreatedExpiryPolicy(new Duration(TimeUnit.SECONDS, 10)); + Factory factory = () -> new CreatedExpiryPolicy(new Duration(TimeUnit.SECONDS, 30)); IgniteCache srcCache = createCache(srcCluster[0], ACTIVE_PASSIVE_CACHE, factory); IgniteCache destCache = createCache(destCluster[0], ACTIVE_PASSIVE_CACHE, factory); @@ -465,7 +465,13 @@ public void testWithExpiryPolicy() throws Exception { assertTrue(waitForCondition(() -> !srcCache.containsKey(0), getTestTimeout())); log.warning(">>>>>> Waiting for removing in destination cache"); - assertTrue(waitForCondition(() -> !destCache.containsKey(0), 20_000)); + + Duration ttl = factory.create().getExpiryForCreation(); + + assertTrue(waitForCondition( + () -> !destCache.containsKey(0), + 2 * TimeUnit.MILLISECONDS.convert(ttl.getDurationAmount(), ttl.getTimeUnit()) + )); } finally { for (IgniteInternalFuture fut : futs)