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..3386b7c60455be 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 @@ -839,7 +839,10 @@ private void modifyPropertiesInternal(Map jobProperties, ((KafkaProgress) progress).checkPartitions(kafkaPartitionOffsets); } - if (Config.isCloudMode()) { + // Kafka client properties do not change the consumed offsets, so keep the cloud progress for replay. + if (Config.isCloudMode() + && (!Strings.isNullOrEmpty(dataSourceProperties.getTopic()) + || !kafkaPartitionOffsets.isEmpty())) { Cloud.ResetRLProgressRequest.Builder builder = Cloud.ResetRLProgressRequest.newBuilder() .setRequestIp(FrontendOptions.getLocalHostAddressCached()); builder.setCloudUniqueId(Config.cloud_unique_id); 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..9fa96233a6980b 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 @@ -25,6 +25,8 @@ import org.apache.doris.catalog.Partition; import org.apache.doris.catalog.Table; import org.apache.doris.catalog.info.PartitionNamesInfo; +import org.apache.doris.cloud.proto.Cloud; +import org.apache.doris.cloud.rpc.MetaServiceProxy; import org.apache.doris.common.Config; import org.apache.doris.common.MetaNotFoundException; import org.apache.doris.common.Pair; @@ -44,6 +46,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; @@ -248,6 +251,59 @@ public void testUpdateLagRebuildsConvertedPropertiesAfterReplay() throws UserExc } } + @Test + public void testReplayKafkaClientPropertiesKeepsCloudProgress() throws Exception { + String originalCloudUniqueId = Config.cloud_unique_id; + String originalMetaServiceEndpoint = Config.meta_service_endpoint; + try { + Config.cloud_unique_id = "test-cloud"; + Config.meta_service_endpoint = "127.0.0.1:20121"; + + KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L, + 1L, "127.0.0.1:9020", "topic1", UserIdentity.ADMIN); + Map partitionOffsets = Maps.newHashMap(); + partitionOffsets.put(1, 10L); + Deencapsulation.setField(routineLoadJob, "progress", new KafkaProgress(partitionOffsets)); + + Map alteredProperties = Maps.newHashMap(); + alteredProperties.put("property.group.id", "replayed-group"); + alteredProperties.put("property.client.id", "replayed-client"); + KafkaDataSourceProperties dataSourceProperties = new KafkaDataSourceProperties(alteredProperties); + dataSourceProperties.setAlter(true); + dataSourceProperties.setTimezone("Asia/Shanghai"); + dataSourceProperties.analyze(); + + Env env = Mockito.mock(Env.class); + MetaServiceProxy metaServiceProxy = Mockito.mock(MetaServiceProxy.class); + Cloud.ResetRLProgressResponse response = Cloud.ResetRLProgressResponse.newBuilder() + .setStatus(Cloud.MetaServiceResponseStatus.newBuilder() + .setCode(Cloud.MetaServiceCode.OK)) + .build(); + Mockito.when(metaServiceProxy.resetRLProgress(Mockito.any())).thenReturn(response); + try (MockedStatic envStatic = Mockito.mockStatic(Env.class); + MockedStatic metaServiceProxyStatic = + Mockito.mockStatic(MetaServiceProxy.class)) { + envStatic.when(Env::getCurrentEnv).thenReturn(env); + metaServiceProxyStatic.when(MetaServiceProxy::getInstance).thenReturn(metaServiceProxy); + + routineLoadJob.replayModifyProperties(new AlterRoutineLoadJobOperationLog( + routineLoadJob.getId(), Maps.newHashMap(), dataSourceProperties)); + + Mockito.verify(metaServiceProxy, Mockito.never()).resetRLProgress(Mockito.any()); + } + + Assert.assertEquals("replayed-group", + routineLoadJob.getConvertedCustomProperties().get("group.id")); + Assert.assertEquals("replayed-client", + routineLoadJob.getConvertedCustomProperties().get("client.id")); + Assert.assertEquals(partitionOffsets, + ((KafkaProgress) routineLoadJob.getProgress()).getOffsetByPartition()); + } finally { + Config.cloud_unique_id = originalCloudUniqueId; + Config.meta_service_endpoint = originalMetaServiceEndpoint; + } + } + @Test public void testUpdateProgressWarnsWhenReadCommittedTaskHasZeroRowsAndLag() throws UserException { KafkaRoutineLoadJob routineLoadJob = new KafkaRoutineLoadJob(1L, "kafka_routine_load_job", 1L,