diff --git a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java index 66dc4cc881fd2a..7f686dec238087 100644 --- a/fe/fe-common/src/main/java/org/apache/doris/common/Config.java +++ b/fe/fe-common/src/main/java/org/apache/doris/common/Config.java @@ -1227,6 +1227,10 @@ public class Config extends ConfigBase { @ConfField(mutable = true, masterOnly = true) public static int streaming_task_min_timeout_sec = 300; + @ConfField(mutable = true, masterOnly = true, description = { + "Minimum interval in seconds between snapshot offset persistence operations"}) + public static int streaming_job_snapshot_offset_persist_interval_sec = 300; + @ConfField(mutable = true, masterOnly = true) public static int streaming_cdc_light_rpc_timeout_sec = 90; diff --git a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java index fb0f1825324674..95956cbeb4919c 100644 --- a/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java +++ b/fe/fe-common/src/main/java/org/apache/doris/job/cdc/DataSourceConfigKeys.java @@ -35,6 +35,7 @@ public class DataSourceConfigKeys { public static final String OFFSET_LATEST = "latest"; public static final String OFFSET_SNAPSHOT = "snapshot"; public static final String SNAPSHOT_SPLIT_SIZE = "snapshot_split_size"; + public static final String SNAPSHOT_SPLIT_SIZE_DEFAULT = "40960"; public static final String SNAPSHOT_SPLIT_KEY = "snapshot_split_key"; public static final String SNAPSHOT_PARALLELISM = "snapshot_parallelism"; public static final String SNAPSHOT_PARALLELISM_DEFAULT = "1"; diff --git a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java index cc7e7e34a18ef3..09082ad0a6e29f 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java +++ b/fe/fe-core/src/main/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJob.java @@ -152,6 +152,7 @@ public class StreamingInsertJob extends AbstractJob(binlogSplit.getStartingOffset()); binlogOffsetPersist.put(SPLIT_ID, BinlogSplit.BINLOG_SPLIT_ID); + clearSnapshotState(); currentOffset = newOffset; hasMoreData = true; } @@ -266,6 +270,31 @@ public void updateOffset(Offset offset) { this.currentOffset = newOffset; } + protected void clearSnapshotState() { + if (MapUtils.isNotEmpty(chunkHighWatermarkMap)) { + chunkHighWatermarkMap = new HashMap<>(); + } + remainingSplits.clear(); + finishedSplits.clear(); + if (committedSplitProgress != null) { + clearProgress(committedSplitProgress); + } + if (cdcSplitProgress != null) { + clearProgress(cdcSplitProgress); + } + } + + public boolean shouldPersistOffset(long lastPersistTimeMs, long currentTimeMs) { + synchronized (splitsLock) { + if (currentOffset == null || !currentOffset.snapshotSplit()) { + return true; + } + } + long intervalMs = Math.max(1L, + (long) Config.streaming_job_snapshot_offset_persist_interval_sec) * 1000L; + return lastPersistTimeMs == 0L || currentTimeMs - lastPersistTimeMs >= intervalMs; + } + @Override public void setBoundBackendId(long boundBackendId) { this.boundBackendId = boundBackendId; diff --git a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java index 0e5bb8fb75319a..8eab2a7978cf41 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java +++ b/fe/fe-core/src/main/java/org/apache/doris/job/offset/jdbc/JdbcTvfSourceOffsetProvider.java @@ -294,10 +294,11 @@ public void updateOffset(Offset offset) { synchronized (splitsLock) { // Mirror binlog offset into bop so it survives FE checkpoint BinlogSplit bs = (BinlogSplit) newOffset.getSplits().get(0); - if (MapUtils.isNotEmpty(bs.getStartingOffset())) { - binlogOffsetPersist = new HashMap<>(bs.getStartingOffset()); - binlogOffsetPersist.put(SPLIT_ID, BinlogSplit.BINLOG_SPLIT_ID); - } + Preconditions.checkArgument(MapUtils.isNotEmpty(bs.getStartingOffset()), + "Committed binlog offset must not be empty"); + binlogOffsetPersist = new HashMap<>(bs.getStartingOffset()); + binlogOffsetPersist.put(SPLIT_ID, BinlogSplit.BINLOG_SPLIT_ID); + clearSnapshotState(); currentOffset = newOffset; hasMoreData = true; } diff --git a/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobOffsetPersistenceTest.java b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobOffsetPersistenceTest.java new file mode 100644 index 00000000000000..80785a34fefb0c --- /dev/null +++ b/fe/fe-core/src/test/java/org/apache/doris/job/extensions/insert/streaming/StreamingInsertJobOffsetPersistenceTest.java @@ -0,0 +1,279 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +package org.apache.doris.job.extensions.insert.streaming; + +import org.apache.doris.catalog.Env; +import org.apache.doris.common.Config; +import org.apache.doris.common.jmockit.Deencapsulation; +import org.apache.doris.job.cdc.request.CommitOffsetRequest; +import org.apache.doris.job.cdc.split.SnapshotSplit; +import org.apache.doris.job.common.FailureReason; +import org.apache.doris.job.common.JobStatus; +import org.apache.doris.job.common.TaskStatus; +import org.apache.doris.job.exception.JobException; +import org.apache.doris.job.manager.JobManager; +import org.apache.doris.job.manager.StreamingTaskManager; +import org.apache.doris.job.offset.jdbc.JdbcSourceOffsetProvider; +import org.apache.doris.transaction.GlobalTransactionMgrIface; +import org.apache.doris.transaction.TxnStateCallbackFactory; + +import org.junit.Assert; +import org.junit.Test; +import org.mockito.MockedStatic; +import org.mockito.Mockito; + +import java.util.Collections; +import java.util.HashMap; +import java.util.concurrent.locks.ReentrantReadWriteLock; + +public class StreamingInsertJobOffsetPersistenceTest { + + @Test + public void testFirstSnapshotCommitPersistsImmediately() throws Exception { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.getRemainingSplits().add(snapshotSplit("source_table:0")); + TestStreamingInsertJob job = newJob(provider, 1001L); + + job.commitOffset(snapshotRequest(1001L, "source_table:0", null)); + + Assert.assertEquals(1, job.journalCount); + Assert.assertNotNull(job.getOffsetProviderPersist()); + } + + @Test + public void testSnapshotCommitWithinIntervalDoesNotPersistAgain() throws Exception { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.getRemainingSplits().add(snapshotSplit("source_table:0")); + TestStreamingInsertJob job = newJob(provider, 1003L); + job.commitOffset(snapshotRequest(1003L, "source_table:0", null)); + + provider.getRemainingSplits().add(snapshotSplit("source_table:1")); + job.commitOffset(snapshotRequest(1003L, "source_table:1", null)); + + Assert.assertEquals(1, job.journalCount); + Assert.assertNotNull(job.getOffsetProviderPersist()); + } + + @Test + public void testBinlogCommitPersistsImmediately() throws Exception { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + TestStreamingInsertJob job = newJob(provider, 1002L); + + job.commitOffset(binlogRequest(1002L, "100")); + job.commitOffset(binlogRequest(1002L, "200")); + + Assert.assertEquals(2, job.journalCount); + Assert.assertNotNull(job.getOffsetProviderPersist()); + } + + @Test + public void testSnapshotToBinlogTransitionPersistsCompactedState() throws Exception { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.getRemainingSplits().add(snapshotSplit("source_table:0")); + TestStreamingInsertJob job = newJob(provider, 1008L); + job.commitOffset(snapshotRequest(1008L, "source_table:0", null)); + Assert.assertEquals(1, job.journalCount); + + job.commitOffset(binlogRequest(1008L, "200")); + + Assert.assertEquals(2, job.journalCount); + Assert.assertFalse(job.getOffsetProviderPersist().contains("source_table:0")); + Assert.assertTrue(provider.getFinishedSplits().isEmpty()); + Assert.assertTrue(provider.getChunkHighWatermarkMap().isEmpty()); + } + + @Test + public void testSnapshotOffsetPersistsOnNextCommitAfterInterval() throws Exception { + int oldInterval = Config.streaming_job_snapshot_offset_persist_interval_sec; + Config.streaming_job_snapshot_offset_persist_interval_sec = 300; + try { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.getRemainingSplits().add(snapshotSplit("source_table:0")); + TestStreamingInsertJob job = newJob(provider, 1011L); + + job.commitOffset(snapshotRequest(1011L, "source_table:0", null)); + Assert.assertEquals(1, job.journalCount); + Deencapsulation.setField(job, "lastOffsetPersistTimeMs", + System.currentTimeMillis() - 300_000L); + provider.getRemainingSplits().add(snapshotSplit("source_table:1")); + job.commitOffset(snapshotRequest(1011L, "source_table:1", null)); + + Assert.assertEquals(2, job.journalCount); + Assert.assertTrue((long) Deencapsulation.getField(job, "lastOffsetPersistTimeMs") > 0L); + } finally { + Config.streaming_job_snapshot_offset_persist_interval_sec = oldInterval; + } + } + + @Test + public void testAlterOffsetReplacesSnapshotState() throws Exception { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.getRemainingSplits().add(snapshotSplit("source_table:0")); + TestStreamingInsertJob job = newJob(provider, 1009L); + job.commitOffset(snapshotRequest(1009L, "source_table:0", null)); + + HashMap properties = new HashMap<>(); + properties.put(StreamingJobProperties.OFFSET_PROPERTY, "{\"lsn\":\"300\"}"); + Deencapsulation.invoke(job, "modifyPropertiesInternal", properties); + + Assert.assertTrue(job.getOffsetProviderPersist().contains("300")); + Assert.assertTrue(provider.getFinishedSplits().isEmpty()); + Assert.assertTrue(provider.getChunkHighWatermarkMap().isEmpty()); + } + + @Test + public void testNaturalFinishPersistsFinalState() throws Exception { + TestStreamingInsertJob job = newJob(new EndJdbcSourceOffsetProvider(), 1012L); + NoopStreamingMultiTblTask task = + (NoopStreamingMultiTblTask) Deencapsulation.getField(job, "runningStreamTask"); + + try (MockedStatic envMockedStatic = Mockito.mockStatic(Env.class)) { + Env env = Mockito.mock(Env.class); + JobManager jobManager = Mockito.mock(JobManager.class); + StreamingTaskManager streamingTaskManager = Mockito.mock(StreamingTaskManager.class); + GlobalTransactionMgrIface transactionMgr = Mockito.mock(GlobalTransactionMgrIface.class); + TxnStateCallbackFactory callbackFactory = Mockito.mock(TxnStateCallbackFactory.class); + envMockedStatic.when(Env::getCurrentEnv).thenReturn(env); + envMockedStatic.when(Env::getCurrentGlobalTransactionMgr).thenReturn(transactionMgr); + Mockito.when(env.getJobManager()).thenReturn(jobManager); + Mockito.when(jobManager.getStreamingTaskManager()).thenReturn(streamingTaskManager); + Mockito.when(transactionMgr.getCallbackFactory()).thenReturn(callbackFactory); + + long beforeFinish = System.currentTimeMillis(); + job.onStreamTaskSuccess(task); + + Assert.assertEquals(JobStatus.FINISHED, job.getJobStatus()); + Assert.assertTrue(job.getFinishTimeMs() >= beforeFinish); + Assert.assertEquals(1, job.journalCount); + Mockito.verify(callbackFactory).removeCallback(9001L); + } + } + + @Test + public void testReplayUpdatedRestoresFinalStateAndRemovesCallback() { + TestStreamingInsertJob job = newJob(new JdbcSourceOffsetProvider(), 1013L); + TestStreamingInsertJob replayJob = newJob(new JdbcSourceOffsetProvider(), 1014L); + replayJob.setJobStatus(JobStatus.FINISHED); + replayJob.setFinishTimeMs(1234L); + + try (MockedStatic envMockedStatic = Mockito.mockStatic(Env.class)) { + GlobalTransactionMgrIface transactionMgr = Mockito.mock(GlobalTransactionMgrIface.class); + TxnStateCallbackFactory callbackFactory = Mockito.mock(TxnStateCallbackFactory.class); + envMockedStatic.when(Env::getCurrentGlobalTransactionMgr).thenReturn(transactionMgr); + Mockito.when(transactionMgr.getCallbackFactory()).thenReturn(callbackFactory); + + job.replayOnUpdated(replayJob); + + Assert.assertEquals(JobStatus.FINISHED, job.getJobStatus()); + Assert.assertEquals(1234L, job.getFinishTimeMs()); + Mockito.verify(callbackFactory).removeCallback(9001L); + } + } + + @Test + public void testReplayUpdatedRestoresStartTimeAndFailureReason() { + TestStreamingInsertJob job = newJob(new JdbcSourceOffsetProvider(), 1015L); + TestStreamingInsertJob replayJob = newJob(new JdbcSourceOffsetProvider(), 1016L); + FailureReason replayFailureReason = new FailureReason("replay failure"); + replayJob.setStartTimeMs(1234L); + replayJob.setFailureReason(replayFailureReason); + + job.replayOnUpdated(replayJob); + + Assert.assertEquals(1234L, job.getStartTimeMs()); + Assert.assertSame(replayFailureReason, job.getFailureReason()); + } + + @Test + public void testReplayUpdatedClearsFailureReason() { + TestStreamingInsertJob job = newJob(new JdbcSourceOffsetProvider(), 1017L); + TestStreamingInsertJob replayJob = newJob(new JdbcSourceOffsetProvider(), 1018L); + job.setFailureReason(new FailureReason("stale failure")); + replayJob.setFailureReason(null); + + job.replayOnUpdated(replayJob); + + Assert.assertNull(job.getFailureReason()); + } + + private static TestStreamingInsertJob newJob(JdbcSourceOffsetProvider provider, long taskId) { + TestStreamingInsertJob job = new TestStreamingInsertJob(); + Deencapsulation.setField(job, "lock", new ReentrantReadWriteLock(true)); + Deencapsulation.setField(job, "jobId", 9001L); + Deencapsulation.setField(job, "jobName", "test_job"); + Deencapsulation.setField(job, "jobStatus", JobStatus.RUNNING); + Deencapsulation.setField(job, "offsetProvider", provider); + Deencapsulation.setField(job, "properties", new HashMap()); + Deencapsulation.setField(job, "targetProperties", new HashMap()); + Deencapsulation.setField(job, "runningStreamTask", new NoopStreamingMultiTblTask(taskId)); + return job; + } + + private static SnapshotSplit snapshotSplit(String splitId) { + return new SnapshotSplit( + splitId, + "source_db.source_table", + Collections.singletonList("id"), + new Object[]{1L}, + new Object[]{2L}, + null); + } + + private static CommitOffsetRequest snapshotRequest(long taskId, String splitId, String tableSchemas) { + CommitOffsetRequest request = new CommitOffsetRequest(); + request.setTaskId(taskId); + request.setOffset("[{\"splitId\":\"" + splitId + "\",\"lsn\":\"100\"}]"); + request.setTableSchemas(tableSchemas); + return request; + } + + private static CommitOffsetRequest binlogRequest(long taskId, String lsn) { + CommitOffsetRequest request = new CommitOffsetRequest(); + request.setTaskId(taskId); + request.setOffset("[{\"splitId\":\"binlog-split\",\"lsn\":\"" + lsn + "\"}]"); + return request; + } + + private static class TestStreamingInsertJob extends StreamingInsertJob { + private int journalCount; + + @Override + public void logUpdateOperation() { + journalCount++; + } + } + + private static class EndJdbcSourceOffsetProvider extends JdbcSourceOffsetProvider { + @Override + public boolean hasReachedEnd() { + return true; + } + } + + private static class NoopStreamingMultiTblTask extends StreamingMultiTblTask { + NoopStreamingMultiTblTask(long taskId) { + super(9001L, taskId, null, null, null, null, null, + new StreamingJobProperties(new HashMap<>()), null, null); + Deencapsulation.setField(this, "status", TaskStatus.RUNNING); + } + + @Override + public void successCallback(CommitOffsetRequest offsetRequest) throws JobException { + } + } +} diff --git a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java index 6efb9959748bcd..f32b616eedb4c1 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/job/offset/jdbc/JdbcSourceOffsetProviderOffsetTest.java @@ -17,16 +17,47 @@ package org.apache.doris.job.offset.jdbc; +import org.apache.doris.common.Config; +import org.apache.doris.common.jmockit.Deencapsulation; import org.apache.doris.job.cdc.split.BinlogSplit; +import org.apache.doris.job.cdc.split.SnapshotSplit; +import org.apache.doris.job.extensions.insert.streaming.StreamingInsertJob; import org.junit.Assert; import org.junit.Test; import java.util.Collections; +import java.util.HashMap; import java.util.Map; public class JdbcSourceOffsetProviderOffsetTest { + @Test + public void testSnapshotOffsetUsesConfiguredPersistInterval() { + int oldInterval = Config.streaming_job_snapshot_offset_persist_interval_sec; + try { + Config.streaming_job_snapshot_offset_persist_interval_sec = 123; + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.currentOffset = new JdbcOffset( + Collections.singletonList(snapshotSplit("source_table:0"))); + + Assert.assertTrue(provider.shouldPersistOffset(0L, 1_000L)); + Assert.assertFalse(provider.shouldPersistOffset(1_000L, 123_999L)); + Assert.assertTrue(provider.shouldPersistOffset(1_000L, 124_000L)); + } finally { + Config.streaming_job_snapshot_offset_persist_interval_sec = oldInterval; + } + } + + @Test + public void testBinlogOffsetPersistsImmediately() { + JdbcSourceOffsetProvider provider = new JdbcSourceOffsetProvider(); + provider.currentOffset = new JdbcOffset(Collections.singletonList( + new BinlogSplit(Collections.singletonMap("lsn", "100")))); + + Assert.assertTrue(provider.shouldPersistOffset(1_000L, 1_001L)); + } + @Test public void testEndOffsetAdvancesWhenCurrentOffsetIsAhead() { assertEndOffsetAdvancesWhenCurrentOffsetIsAhead(new TestJdbcSourceOffsetProvider(-1)); @@ -76,6 +107,71 @@ public void testStaleCompareDoesNotOverwriteAlteredCurrentOffsetState() { ((BinlogSplit) provider.currentOffset.getSplits().get(0)).getStartingOffset()); } + @Test + public void testValidBinlogOffsetClearsSnapshotState() { + assertValidBinlogOffsetClearsSnapshotState(new TestJdbcSourceOffsetProvider(-1)); + } + + @Test + public void testTvfValidBinlogOffsetClearsSnapshotState() { + assertValidBinlogOffsetClearsSnapshotState(new TestJdbcTvfSourceOffsetProvider(-1)); + } + + @Test + public void testEmptyBinlogOffsetIsRejected() { + assertEmptyBinlogOffsetIsRejected(new TestJdbcSourceOffsetProvider(-1)); + } + + @Test + public void testTvfEmptyBinlogOffsetIsRejected() { + assertEmptyBinlogOffsetIsRejected(new TestJdbcTvfSourceOffsetProvider(-1)); + } + + @Test + public void testRepeatedValidBinlogOffsetCleanupIsIdempotent() { + assertRepeatedValidBinlogOffsetCleanupIsIdempotent(new TestJdbcSourceOffsetProvider(-1)); + } + + @Test + public void testTvfRepeatedValidBinlogOffsetCleanupIsIdempotent() { + assertRepeatedValidBinlogOffsetCleanupIsIdempotent(new TestJdbcTvfSourceOffsetProvider(-1)); + } + + @Test + public void testBinlogOffsetRestoredFromPersistInfo() throws Exception { + JdbcSourceOffsetProvider source = new TestJdbcSourceOffsetProvider(-1); + source.updateOffset(new JdbcOffset(Collections.singletonList( + new BinlogSplit(Collections.singletonMap("lsn", "200"))))); + StreamingInsertJob job = mockJobWithPersistInfo(source.getPersistInfo()); + JdbcSourceOffsetProvider restored = new JdbcSourceOffsetProvider(); + + restored.replayIfNeed(job); + + Assert.assertNotNull(restored.currentOffset); + Assert.assertFalse(restored.currentOffset.snapshotSplit()); + Assert.assertEquals("200", ((BinlogSplit) restored.currentOffset.getSplits().get(0)) + .getStartingOffset().get("lsn")); + Assert.assertTrue(restored.chunkHighWatermarkMap.isEmpty()); + } + + @Test + public void testTvfBinlogOffsetRestoredFromPersistInfo() throws Exception { + JdbcSourceOffsetProvider source = new TestJdbcTvfSourceOffsetProvider(-1); + source.updateOffset(new JdbcOffset(Collections.singletonList( + new BinlogSplit(Collections.singletonMap("lsn", "200"))))); + StreamingInsertJob job = mockJobWithPersistInfo(source.getPersistInfo()); + JdbcTvfSourceOffsetProvider restored = new JdbcTvfSourceOffsetProvider(); + + restored.restoreFromPersistInfo(source.getPersistInfo()); + restored.replayIfNeed(job); + + Assert.assertNotNull(restored.currentOffset); + Assert.assertFalse(restored.currentOffset.snapshotSplit()); + Assert.assertEquals("200", ((BinlogSplit) restored.currentOffset.getSplits().get(0)) + .getStartingOffset().get("lsn")); + Assert.assertTrue(restored.chunkHighWatermarkMap.isEmpty()); + } + private static void assertEndOffsetAdvancesWhenCurrentOffsetIsAhead(JdbcSourceOffsetProvider provider) { Map staleEndOffset = Collections.singletonMap("lsn", "100"); Map committedOffset = Collections.singletonMap("lsn", "200"); @@ -91,6 +187,114 @@ private static void assertEndOffsetAdvancesWhenCurrentOffsetIsAhead(JdbcSourceOf Assert.assertEquals("{\"lsn\":\"200\"}", provider.getShowMaxOffset()); } + private static void assertValidBinlogOffsetClearsSnapshotState(JdbcSourceOffsetProvider provider) { + seedSnapshotState(provider); + Map binlogOffset = Collections.singletonMap("lsn", "200"); + + provider.updateOffset(new JdbcOffset( + Collections.singletonList(new BinlogSplit(binlogOffset)))); + + Assert.assertTrue(provider.chunkHighWatermarkMap.isEmpty()); + Assert.assertTrue(provider.remainingSplits.isEmpty()); + Assert.assertTrue(provider.finishedSplits.isEmpty()); + assertProgressCleared(provider.committedSplitProgress); + assertProgressCleared(provider.cdcSplitProgress); + Assert.assertEquals("table-schemas", provider.tableSchemas); + Map expectedPersist = new HashMap<>(binlogOffset); + expectedPersist.put(JdbcSourceOffsetProvider.SPLIT_ID, BinlogSplit.BINLOG_SPLIT_ID); + Assert.assertEquals(expectedPersist, provider.binlogOffsetPersist); + String persistInfo = provider.getPersistInfo(); + Assert.assertFalse(persistInfo.contains("source_table:0")); + Assert.assertFalse(persistInfo.contains("source_table:1")); + Assert.assertTrue(persistInfo.contains("table-schemas")); + } + + private static void assertEmptyBinlogOffsetIsRejected(JdbcSourceOffsetProvider provider) { + seedSnapshotState(provider); + + try { + provider.updateOffset(new JdbcOffset( + Collections.singletonList(new BinlogSplit(Collections.emptyMap())))); + Assert.fail("Empty committed binlog offset should be rejected"); + } catch (IllegalArgumentException e) { + Assert.assertTrue(e.getMessage().contains("Committed binlog offset must not be empty")); + } + + Assert.assertFalse(provider.chunkHighWatermarkMap.isEmpty()); + Assert.assertFalse(provider.remainingSplits.isEmpty()); + Assert.assertFalse(provider.finishedSplits.isEmpty()); + Assert.assertEquals("source_table", provider.committedSplitProgress.getCurrentSplittingTable()); + Assert.assertEquals("source_table", provider.cdcSplitProgress.getCurrentSplittingTable()); + Assert.assertNull(provider.binlogOffsetPersist); + } + + private static void assertRepeatedValidBinlogOffsetCleanupIsIdempotent( + JdbcSourceOffsetProvider provider) { + seedSnapshotState(provider); + JdbcOffset binlogOffset = new JdbcOffset(Collections.singletonList( + new BinlogSplit(Collections.singletonMap("lsn", "200")))); + + provider.updateOffset(binlogOffset); + String firstPersistInfo = provider.getPersistInfo(); + Map>> clearedHighWatermarkMap = + provider.chunkHighWatermarkMap; + provider.updateOffset(binlogOffset); + + Assert.assertEquals(firstPersistInfo, provider.getPersistInfo()); + Assert.assertSame(clearedHighWatermarkMap, provider.chunkHighWatermarkMap); + Assert.assertTrue(provider.chunkHighWatermarkMap.isEmpty()); + Assert.assertTrue(provider.remainingSplits.isEmpty()); + Assert.assertTrue(provider.finishedSplits.isEmpty()); + assertProgressCleared(provider.committedSplitProgress); + assertProgressCleared(provider.cdcSplitProgress); + } + + private static StreamingInsertJob mockJobWithPersistInfo(String persistInfo) { + StreamingInsertJob job = new ReplayStreamingInsertJob(); + Deencapsulation.setField(job, "jobId", 9001L); + Deencapsulation.setField(job, "syncTables", Collections.emptyList()); + job.setOffsetProviderPersist(persistInfo); + return job; + } + + private static void seedSnapshotState(JdbcSourceOffsetProvider provider) { + SnapshotSplit remaining = snapshotSplit("source_table:1"); + SnapshotSplit finished = snapshotSplit("source_table:0"); + provider.remainingSplits.add(remaining); + provider.finishedSplits.add(finished); + provider.chunkHighWatermarkMap + .computeIfAbsent("source_db.source_table", key -> new HashMap<>()) + .put(finished.getSplitId(), finished.getHighWatermark()); + provider.committedSplitProgress = splitProgress(); + provider.cdcSplitProgress = splitProgress(); + provider.tableSchemas = "table-schemas"; + } + + private static SnapshotSplit snapshotSplit(String splitId) { + return new SnapshotSplit( + splitId, + "source_db.source_table", + Collections.singletonList("id"), + new Object[]{1L}, + new Object[]{2L}, + Collections.singletonMap("lsn", "100")); + } + + private static JdbcSourceOffsetProvider.SplitProgress splitProgress() { + JdbcSourceOffsetProvider.SplitProgress progress = new JdbcSourceOffsetProvider.SplitProgress(); + progress.setCurrentSplittingTable("source_table"); + progress.setNextSplitStart(new Object[]{2L}); + progress.setNextSplitId(2); + return progress; + } + + private static void assertProgressCleared(JdbcSourceOffsetProvider.SplitProgress progress) { + Assert.assertNotNull(progress); + Assert.assertNull(progress.getCurrentSplittingTable()); + Assert.assertNull(progress.getNextSplitStart()); + Assert.assertNull(progress.getNextSplitId()); + } + private static class TestJdbcSourceOffsetProvider extends JdbcSourceOffsetProvider { private final int compareResult; @@ -104,6 +308,12 @@ protected int compareOffset(Map offsetFirst, Map } } + private static class ReplayStreamingInsertJob extends StreamingInsertJob { + ReplayStreamingInsertJob() { + super(); + } + } + private static class TestJdbcTvfSourceOffsetProvider extends JdbcTvfSourceOffsetProvider { private final int compareResult; diff --git a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java index 11075ea2d81964..0ad2629ce94db6 100644 --- a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java +++ b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/mysql/MySqlSourceReader.java @@ -1039,10 +1039,11 @@ private MySqlSourceConfig generateMySqlConfig( configFactory.debeziumProperties(dbzProps); configFactory.heartbeatInterval(Duration.ofMillis(DEBEZIUM_HEARTBEAT_INTERVAL_MS)); - if (cdcConfig.containsKey(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE)) { - configFactory.splitSize( - Integer.parseInt(cdcConfig.get(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE))); - } + configFactory.splitSize( + Integer.parseInt( + cdcConfig.getOrDefault( + DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE, + DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE_DEFAULT))); // todo: Currently, only one split key is supported; future will require multiple split // keys. diff --git a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java index e9cc9e8d9859dc..330f461510bf99 100644 --- a/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java +++ b/fs_brokers/cdc_client/src/main/java/org/apache/doris/cdcclient/source/reader/postgres/PostgresSourceReader.java @@ -294,11 +294,11 @@ private PostgresSourceConfig generatePostgresConfig( throw new RuntimeException("Unknown offset " + startupMode); } - // Set split size if provided - if (cdcConfig.containsKey(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE)) { - configFactory.splitSize( - Integer.parseInt(cdcConfig.get(DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE))); - } + configFactory.splitSize( + Integer.parseInt( + cdcConfig.getOrDefault( + DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE, + DataSourceConfigKeys.SNAPSHOT_SPLIT_SIZE_DEFAULT))); if (cdcConfig.containsKey(DataSourceConfigKeys.SNAPSHOT_SPLIT_KEY)) { configFactory.chunkKeyColumn(cdcConfig.get(DataSourceConfigKeys.SNAPSHOT_SPLIT_KEY)); diff --git a/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_snapshot_finished_restart_fe.groovy b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_snapshot_finished_restart_fe.groovy new file mode 100644 index 00000000000000..6bf0a57eb8269c --- /dev/null +++ b/regression-test/suites/job_p0/streaming_job/cdc/test_streaming_mysql_job_snapshot_finished_restart_fe.groovy @@ -0,0 +1,160 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +import org.apache.doris.regression.suite.ClusterOptions +import org.awaitility.Awaitility + +import static java.util.concurrent.TimeUnit.SECONDS + +suite("test_streaming_mysql_job_snapshot_finished_restart_fe", + "docker,mysql,external_docker,external_docker_mysql,nondatalake") { + def jobName = "test_streaming_mysql_job_snapshot_finished_restart_fe" + def tableName = "snapshot_finished_restart_fe" + def mysqlDb = "test_cdc_db" + def totalRows = 5 + def options = new ClusterOptions() + options.setFeNum(1) + options.cloudMode = null + + docker(options) { + def currentDb = (sql "select database()")[0][0] + + sql """DROP JOB IF EXISTS where jobname = '${jobName}'""" + sql """DROP TABLE IF EXISTS ${currentDb}.${tableName} FORCE""" + + String enabled = context.config.otherConfigs.get("enableJdbcTest") + if (enabled != null && enabled.equalsIgnoreCase("true")) { + String mysqlPort = context.config.otherConfigs.get("mysql_57_port") + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String s3Endpoint = getS3Endpoint() + String bucket = getS3BucketName() + String driverUrl = + "https://${bucket}.${s3Endpoint}/regression/jdbc_driver/mysql-connector-j-8.4.0.jar" + + connect("root", "123456", "jdbc:mysql://${externalEnvIp}:${mysqlPort}") { + sql """CREATE DATABASE IF NOT EXISTS ${mysqlDb}""" + sql """DROP TABLE IF EXISTS ${mysqlDb}.${tableName}""" + sql """CREATE TABLE ${mysqlDb}.${tableName} ( + `id` int NOT NULL, + `name` varchar(200), + PRIMARY KEY (`id`) + ) ENGINE=InnoDB""" + sql """INSERT INTO ${mysqlDb}.${tableName} (id, name) VALUES + (1, 'name_1'), + (2, 'name_2'), + (3, 'name_3'), + (4, 'name_4'), + (5, 'name_5')""" + } + + sql """CREATE JOB ${jobName} + ON STREAMING + FROM MYSQL ( + "jdbc_url" = "jdbc:mysql://${externalEnvIp}:${mysqlPort}", + "driver_url" = "${driverUrl}", + "driver_class" = "com.mysql.cj.jdbc.Driver", + "user" = "root", + "password" = "123456", + "database" = "${mysqlDb}", + "include_tables" = "${tableName}", + "offset" = "snapshot", + "snapshot_split_size" = "1", + "snapshot_parallelism" = "1" + ) + TO DATABASE ${currentDb} ( + "table.create.properties.replication_num" = "1" + ) + """ + + try { + Awaitility.await().atMost(300, SECONDS) + .pollInterval(2, SECONDS).until( + { + def jobStatus = sql """ + SELECT Status + FROM jobs("type"="insert") + WHERE Name='${jobName}' AND ExecuteType='STREAMING' + """ + log.info("jobStatus before FE restart: " + jobStatus) + jobStatus.size() == 1 && jobStatus.get(0).get(0) == "FINISHED" + } + ) + + def jobIdRows = sql """ + SELECT Id + FROM jobs("type"="insert") + WHERE Name='${jobName}' AND ExecuteType='STREAMING' + """ + assert jobIdRows.size() == 1 + def jobId = jobIdRows.get(0).get(0).toString() + + def rowsBeforeRestart = sql """ + SELECT COUNT(*), COUNT(DISTINCT id) + FROM ${currentDb}.${tableName} + """ + assert rowsBeforeRestart.size() == 1 + assert rowsBeforeRestart.get(0).get(0) == totalRows + assert rowsBeforeRestart.get(0).get(1) == totalRows + + cluster.restartFrontends() + sleep(60000) + context.reconnectFe() + + // A terminal job must survive replay with its finish time. Otherwise + // the scheduler may treat it as an expired job and remove it. + Awaitility.await().atMost(120, SECONDS) + .pollInterval(2, SECONDS).until( + { + def jobAfterRestart = sql """ + SELECT Id, Status + FROM jobs("type"="insert") + WHERE Name='${jobName}' AND ExecuteType='STREAMING' + """ + log.info("job after FE restart: " + jobAfterRestart) + jobAfterRestart.size() == 1 + && jobAfterRestart.get(0).get(0).toString() == jobId + && jobAfterRestart.get(0).get(1) == "FINISHED" + } + ) + + def rowsAfterRestart = sql """ + SELECT COUNT(*), COUNT(DISTINCT id) + FROM ${currentDb}.${tableName} + """ + assert rowsAfterRestart.size() == 1 + assert rowsAfterRestart.get(0).get(0) == totalRows + assert rowsAfterRestart.get(0).get(1) == totalRows + } catch (Exception ex) { + def showJob = sql """ + SELECT * + FROM jobs("type"="insert") + WHERE Name='${jobName}' + """ + def showTask = sql """ + SELECT * + FROM tasks("type"="insert") + WHERE JobName='${jobName}' + """ + log.info("show job: " + showJob) + log.info("show task: " + showTask) + throw ex + } finally { + sql """DROP JOB IF EXISTS where jobname = '${jobName}'""" + } + } + } +}