diff --git a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java index 885021440351d7..6ce738a1c8d4ad 100644 --- a/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java +++ b/fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java @@ -412,6 +412,17 @@ protected void replayUpdateProgress(RLTaskTxnCommitAttachment attachment) { protected void updateCloudProgress(RLTaskTxnCommitAttachment attachment) { super.updateCloudProgress(attachment); updateProgressAndOffsetsCache(attachment); + retainCustomKafkaPartitionProgress(); + } + + private void retainCustomKafkaPartitionProgress() { + if (CollectionUtils.isEmpty(customKafkaPartitions)) { + return; + } + // Meta Service retains offsets omitted by a partial reset, but an explicit partition list is a consumption pin. + KafkaProgress kafkaProgress = (KafkaProgress) progress; + progress = new KafkaProgress(kafkaProgress.getPartitionIdToOffset(customKafkaPartitions)); + cachedPartitionWithLatestOffsets.keySet().retainAll(customKafkaPartitions); } @Override @@ -868,8 +879,7 @@ private void modifyPropertiesInternal(Map jobProperties, // modify partition offset if (!kafkaPartitionOffsets.isEmpty()) { - // we can only modify the partition that is being consumed - ((KafkaProgress) progress).modifyOffset(kafkaPartitionOffsets); + replaceKafkaPartitionsAndOffsets(kafkaPartitionOffsets); } // modify broker list @@ -897,6 +907,19 @@ private void modifyPropertiesInternal(Map jobProperties, this.id, jobProperties, dataSourceProperties); } + private void replaceKafkaPartitionsAndOffsets(List> kafkaPartitionOffsets) { + List alteredPartitions = Lists.newArrayListWithCapacity(kafkaPartitionOffsets.size()); + Map alteredProgress = Maps.newHashMapWithExpectedSize(kafkaPartitionOffsets.size()); + for (Pair partitionOffset : kafkaPartitionOffsets) { + alteredPartitions.add(partitionOffset.first); + alteredProgress.put(partitionOffset.first, partitionOffset.second); + } + customKafkaPartitions = alteredPartitions; + currentKafkaPartitions = Lists.newArrayList(alteredPartitions); + progress = new KafkaProgress(alteredProgress); + cachedPartitionWithLatestOffsets.keySet().retainAll(alteredPartitions); + } + private void resetCloudProgress(Cloud.ResetRLProgressRequest.Builder builder) throws DdlException { Cloud.ResetRLProgressResponse response; try { diff --git a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java index 7f0c8588372403..4679fcc682715e 100644 --- a/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java +++ b/fe/fe-core/src/test/java/org/apache/doris/load/routineload/KafkaRoutineLoadJobTest.java @@ -44,6 +44,7 @@ import org.apache.doris.nereids.trees.plans.commands.info.LabelNameInfo; import org.apache.doris.nereids.trees.plans.commands.load.LoadProperty; import org.apache.doris.nereids.trees.plans.commands.load.LoadSeparator; +import org.apache.doris.persist.AlterRoutineLoadJobOperationLog; import org.apache.doris.qe.ConnectContext; import org.apache.doris.thrift.TResourceInfo; import org.apache.doris.thrift.TRoutineLoadTask; @@ -61,6 +62,10 @@ import org.mockito.MockedStatic; import org.mockito.Mockito; +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.DataInputStream; +import java.io.DataOutputStream; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -248,6 +253,107 @@ public void testUpdateLagRebuildsConvertedPropertiesAfterReplay() throws UserExc } } + @Test + public void testReplayExplicitPartitionAlterPinsPartitions() throws Exception { + String originalDeployMode = Config.deploy_mode; + String originalCloudUniqueId = Config.cloud_unique_id; + try { + Config.deploy_mode = ""; + Config.cloud_unique_id = ""; + + Map alterProperties = Maps.newHashMap(); + alterProperties.put(KafkaConfiguration.KAFKA_PARTITIONS.getName(), "1"); + alterProperties.put(KafkaConfiguration.KAFKA_OFFSETS.getName(), "2"); + KafkaDataSourceProperties dataSourceProperties = new KafkaDataSourceProperties(alterProperties); + dataSourceProperties.setAlter(true); + dataSourceProperties.setTimezone("Asia/Shanghai"); + dataSourceProperties.analyze(); + + AlterRoutineLoadJobOperationLog operationLog = new AlterRoutineLoadJobOperationLog( + 1L, Maps.newHashMap(), dataSourceProperties); + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + try (DataOutputStream out = new DataOutputStream(bytes)) { + operationLog.write(out); + } + try (DataInputStream in = new DataInputStream(new ByteArrayInputStream(bytes.toByteArray()))) { + operationLog = AlterRoutineLoadJobOperationLog.read(in); + } + + KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L, + 1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN); + Deencapsulation.setField(routineLoadJob, "customKafkaPartitions", Lists.newArrayList(0, 1, 2)); + Deencapsulation.setField(routineLoadJob, "currentKafkaPartitions", Lists.newArrayList(0, 1, 2)); + Map partitionOffsets = Maps.newHashMap(); + partitionOffsets.put(0, 10L); + partitionOffsets.put(1, 11L); + partitionOffsets.put(2, 12L); + Deencapsulation.setField(routineLoadJob, "progress", new KafkaProgress(partitionOffsets)); + Deencapsulation.setField(routineLoadJob, "cachedPartitionWithLatestOffsets", + Maps.newHashMap(partitionOffsets)); + + routineLoadJob.replayModifyProperties(operationLog); + + Assert.assertEquals(Lists.newArrayList(1), + Deencapsulation.getField(routineLoadJob, "customKafkaPartitions")); + Assert.assertEquals(Lists.newArrayList(1), + Deencapsulation.getField(routineLoadJob, "currentKafkaPartitions")); + Map expectedProgress = Maps.newHashMap(); + expectedProgress.put(1, 2L); + Assert.assertEquals(expectedProgress, + ((KafkaProgress) routineLoadJob.getProgress()).getOffsetByPartition()); + Map expectedLatestOffsets = Maps.newHashMap(); + expectedLatestOffsets.put(1, 11L); + Assert.assertEquals(expectedLatestOffsets, + Deencapsulation.getField(routineLoadJob, "cachedPartitionWithLatestOffsets")); + Assert.assertTrue(routineLoadJob.dataSourcePropertiesJsonToString() + .contains("\"currentKafkaPartitions\":\"1\"")); + String showCreateInfo = routineLoadJob.getShowCreateInfo(); + Assert.assertTrue(showCreateInfo.contains("\"kafka_partitions\" = \"1\"")); + Assert.assertTrue(showCreateInfo.contains("\"kafka_offsets\" = \"2\"")); + } finally { + Config.deploy_mode = originalDeployMode; + Config.cloud_unique_id = originalCloudUniqueId; + } + } + + @Test + public void testCloudProgressKeepsExplicitPartitionPin() { + String originalCloudUniqueId = Config.cloud_unique_id; + try { + Config.cloud_unique_id = "test-cloud"; + KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L, + 1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN); + Deencapsulation.setField(routineLoadJob, "customKafkaPartitions", Lists.newArrayList(1)); + Deencapsulation.setField(routineLoadJob, "currentKafkaPartitions", Lists.newArrayList(1)); + Map pinnedProgress = Maps.newHashMap(); + pinnedProgress.put(1, 2L); + Deencapsulation.setField(routineLoadJob, "progress", new KafkaProgress(pinnedProgress)); + Map latestOffsets = Maps.newHashMap(); + latestOffsets.put(1, 11L); + Deencapsulation.setField(routineLoadJob, "cachedPartitionWithLatestOffsets", + latestOffsets); + + RLTaskTxnCommitAttachment attachment = new RLTaskTxnCommitAttachment(); + Map cloudProgress = Maps.newHashMap(); + cloudProgress.put(0, 20L); + cloudProgress.put(1, 21L); + cloudProgress.put(2, 22L); + Deencapsulation.setField(attachment, "progress", + new KafkaProgress(cloudProgress)); + + Deencapsulation.invoke(routineLoadJob, "updateCloudProgress", attachment); + + Map expectedOffsets = Maps.newHashMap(); + expectedOffsets.put(1, 22L); + Assert.assertEquals(expectedOffsets, + ((KafkaProgress) routineLoadJob.getProgress()).getOffsetByPartition()); + Assert.assertEquals(expectedOffsets, + Deencapsulation.getField(routineLoadJob, "cachedPartitionWithLatestOffsets")); + } finally { + Config.cloud_unique_id = originalCloudUniqueId; + } + } + @Test public void testUpdateProgressWarnsWhenReadCommittedTaskHasZeroRowsAndLag() throws UserException { KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L,