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
6 changes: 6 additions & 0 deletions docs/_docs/persistence/change-data-capture.adoc
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,10 @@ The CDC can be disabled manually or by configured directory maximum size. In thi

WARNING: All changes in skipped segments will be lost!

Caches of the data regions without persistence do not write changes to the WAL while the CDC is disabled manually.
On the nodes with such regions the segment that is current at the disabling is skipped as well, so the gap appears
even if no segment is archived while the CDC is disabled.

So when enabled there will be gap between segments: `0000000000000002.wal`, `0000000000000010.wal`, `0000000000000011.wal`, for example.
In this case `ignite-cdc.sh` will fail with the something like "Found missed segments. Some events are missed. Exiting! [lastSegment=2, nextSegment=10]".

Expand Down Expand Up @@ -191,6 +195,8 @@ NOTE: The command requires an active cluster and will be canceled if the cluster

NOTE: The command will be canceled if cluster was not rebalanced or topology changed (node left/joined, baseline changed).

NOTE: The command fails if the CDC is disabled by the `cdc.disabled` distributed property.

To forcefully resend all cache data to CDC you can run the following link:tools/control-script[Control Script] command:

[source,shell]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;
import java.util.stream.LongStream;
import org.apache.ignite.Ignite;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.cache.CachePeekMode;
Expand Down Expand Up @@ -63,7 +65,6 @@
import static org.apache.ignite.cdc.AbstractCdcTest.ChangeEventType.UPDATE;
import static org.apache.ignite.cdc.AbstractCdcTest.KEYS_CNT;
import static org.apache.ignite.cdc.CdcSelfTest.addData;
import static org.apache.ignite.events.EventType.EVT_WAL_SEGMENT_ARCHIVED;
import static org.apache.ignite.internal.commandline.CommandHandler.EXIT_CODE_INVALID_ARGUMENTS;
import static org.apache.ignite.internal.commandline.CommandHandler.EXIT_CODE_OK;
import static org.apache.ignite.internal.commandline.CommandHandler.EXIT_CODE_UNEXPECTED_ERROR;
Expand Down Expand Up @@ -116,8 +117,6 @@ public class CdcCommandTest extends GridCommandHandlerAbstractTest {
.setDefaultDataRegionConfiguration(new DataRegionConfiguration()
.setCdcEnabled(true)));

cfg.setIncludeEventTypes(EVT_WAL_SEGMENT_ARCHIVED);

cfg.setPluginProviders(new AbstractTestPluginProvider() {
@Override public String name() {
return "Test WAL provider";
Expand Down Expand Up @@ -217,24 +216,28 @@ public void testDeleteLostSegmentLinksApplicationNotClosed() throws Exception {
/** */
@Test
public void testDeleteLostSegmentLinks() throws Exception {
checkDeleteLostSegmentLinks(F.asList(0L, 2L), F.asList(2L), true);
checkDeleteLostSegmentLinks(1, true);
}

/** */
@Test
public void testDeleteLostSegmentLinksOneNode() throws Exception {
checkDeleteLostSegmentLinks(F.asList(0L, 2L), F.asList(2L), false);
checkDeleteLostSegmentLinks(1, false);
}

/** */
@Test
public void testDeleteLostSegmentLinksMultipleGaps() throws Exception {
checkDeleteLostSegmentLinks(F.asList(0L, 3L, 5L), F.asList(5L), true);
checkDeleteLostSegmentLinks(2, true);
}

/** */
private void checkDeleteLostSegmentLinks(List<Long> expBefore, List<Long> expAfter, boolean allNodes) throws Exception {
archiveSegmentLinks(expBefore);
private void checkDeleteLostSegmentLinks(int disableCnt, boolean allNodes) throws Exception {
archiveSegmentLinks(disableCnt);

// The segment that is current at a disabling is never linked: the even segments are linked.
List<Long> expBefore = LongStream.rangeClosed(0, disableCnt).map(i -> 2 * i).boxed().collect(Collectors.toList());
List<Long> expAfter = F.asList(2L * disableCnt);

checkLinks(srv0, expBefore);
checkLinks(srv1, expBefore);
Expand All @@ -255,34 +258,32 @@ private void checkLinks(IgniteEx srv, List<Long> expLinks) {
File[] links = ft.walCdcSegments();

assertEquals(expLinks.size(), links.length);
Arrays.stream(links).map(File::toPath).map(ft::walSegmentIndex)
.allMatch(expLinks::contains);
assertTrue(Arrays.stream(links).map(File::toPath).map(ft::walSegmentIndex)
.allMatch(expLinks::contains));
}

/** Archive given segments links with possible gaps. */
private void archiveSegmentLinks(List<Long> idxs) throws Exception {
for (long idx = 0; idx <= idxs.stream().mapToLong(v -> v).max().getAsLong(); idx++) {
cdcDisabled.propagate(!idxs.contains(idx));
/** Archives segments, disabling and enabling CDC the given number of times. */
private void archiveSegmentLinks(int disableCnt) throws Exception {
archiveSegment(0);

for (int i = 1; i <= disableCnt; i++) {
cdcDisabled.propagate(true);
cdcDisabled.propagate(false);

archiveSegment();
// The first record closes the segment that was current at the disabling, the data goes to the next one.
archiveSegment(2 * i);
}
}

/** */
private void archiveSegment() throws Exception {
CountDownLatch latch = new CountDownLatch(G.allGrids().size());
/** Writes data and waits for the given segment to be archived on all nodes. */
private void archiveSegment(long idx) throws Exception {
addData(srv1.cache(DEFAULT_CACHE_NAME), 0, 1);

for (Ignite srv : G.allGrids()) {
srv.events().localListen(evt -> {
latch.countDown();
IgniteWriteAheadLogManager wal = ((IgniteEx)srv).context().cache().context().wal(true);

return false;
}, EVT_WAL_SEGMENT_ARCHIVED);
assertTrue(waitForCondition(() -> wal.lastArchivedSegment() >= idx, getTestTimeout()));
}

addData(srv1.cache(DEFAULT_CACHE_NAME), 0, KEYS_CNT);

latch.await(getTestTimeout(), TimeUnit.MILLISECONDS);
}

/** */
Expand All @@ -303,6 +304,56 @@ public void testParseResend() {
"Please specify a value for argument: --caches");
}

/** */
@Test
public void testResendCdcDisabled() throws Exception {
injectTestSystemOut();

cdcDisabled.propagate(true);

String out = executeCommand(EXIT_CODE_UNEXPECTED_ERROR, CDC, RESEND, CACHES, DEFAULT_CACHE_NAME);

if (cliCommandHandler())
assertContains(log, out, "CDC is disabled");
}

/** */
@Test
public void testResendCancelOnCdcDisabled() throws Exception {
injectTestSystemOut();

addData(srv0.cache(DEFAULT_CACHE_NAME), 0, KEYS_CNT);

CountDownLatch preload = new CountDownLatch(1);
CountDownLatch disabled = new CountDownLatch(1);

AtomicInteger cnt = new AtomicInteger();

onLogLsnr = rec -> {
if (cnt.incrementAndGet() < KEYS_CNT / 2)
return;

preload.countDown();

U.await(disabled);
};

IgniteInternalFuture<Object> fut = GridTestUtils.runAsync(() -> {
String out = executeCommand(EXIT_CODE_UNEXPECTED_ERROR, CDC, RESEND, CACHES, DEFAULT_CACHE_NAME);

if (cliCommandHandler())
assertContains(log, out, "CDC is disabled");
});

preload.await();

cdcDisabled.propagate(true);

disabled.countDown();

fut.get();
}

/** */
@Test
public void testResendCacheData() throws Exception {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
import org.jetbrains.annotations.Nullable;

import static org.apache.ignite.cluster.ClusterState.INACTIVE;
import static org.apache.ignite.internal.processors.cache.persistence.wal.FileWriteAheadLogManager.CDC_DISABLED;
import static org.apache.ignite.internal.util.lang.ClusterNodeFunc.nodeIds;

/**
Expand Down Expand Up @@ -154,12 +155,14 @@ protected CdcCacheDataResendJob(CdcResendCommandArg arg, AffinityTopologyVersion
caches.add(cache);
}

if (log.isInfoEnabled())
log.info("CDC cache data resend started [caches=" + String.join(", ", arg.caches()) + ']');

wal = ignite.context().cache().context().wal(true);
exchange = ignite.context().cache().context().exchange();

ensureCdcEnabled();

if (log.isInfoEnabled())
log.info("CDC cache data resend started [caches=" + String.join(", ", arg.caches()) + ']');

try {
Iterator<IgniteInternalCache<?, ?>> iter = caches.iterator();

Expand Down Expand Up @@ -198,6 +201,7 @@ private void resendCacheData(IgniteInternalCache<?, ?> cache) throws IgniteCheck
break;

ensureTopologyNotChanged();
ensureCdcEnabled();

KeyCacheObject key = row.key();

Expand Down Expand Up @@ -256,5 +260,11 @@ private void ensureTopologyNotChanged() {
lastFut = fut;
}
}

/** */
private void ensureCdcEnabled() {
if (wal.cdcForceDisabled())
throw new IgniteException("CDC is disabled by the '" + CDC_DISABLED + "' distributed property.");
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,11 @@ public interface IgniteWriteAheadLogManager extends GridCacheSharedManager, Igni
*/
public boolean isFullSync();

/**
* @return {@code True} if CDC is disabled by the {@code cdc.disabled} distributed property.
*/
public boolean cdcForceDisabled();

/**
* @return Current serializer version.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -980,7 +980,7 @@ public boolean persistenceEnabled() {
* @return {@code True} if {@link DataRecord} should be loged in the WAL.
*/
public boolean logDataRecords() {
return walEnabled() && (persistenceEnabled || cdcEnabled());
return walEnabled() && (persistenceEnabled || (cdcEnabled() && !wal().cdcForceDisabled()));
}

/** @return {@code True} if CDC enabled. */
Expand Down
Loading
Loading