github-actions[bot] commented on code in PR #67517:
URL: https://github.com/apache/doris/pull/67517#discussion_r4227663750


##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java:
##########
@@ -868,8 +879,7 @@ private void modifyPropertiesInternal(Map<String, String> 
jobProperties,
 
             // modify partition offset
             if (!kafkaPartitionOffsets.isEmpty()) {
-                // we can only modify the partition that is being consumed
-                ((KafkaProgress) progress).modifyOffset(kafkaPartitionOffsets);
+                replaceKafkaPartitionsAndOffsets(kafkaPartitionOffsets);

Review Comment:
   [P1] Install the pin when replaying an ALTER on a follower with empty 
progress. For a cloud job created without explicit partitions, the create 
journal contains an empty progress map; only the leader scheduler later 
populates it. If the leader pauses and alters to partition 1 before a 
checkpoint, a follower replaying the ALTER fails the earlier checkPartitions 
call against its empty map, and replayModifyProperties only logs the exception, 
so this replacement is never reached. After promotion and resume it 
auto-discovers all Kafka partitions and can consume omitted partitions again. 
Replay must install the explicit definition independently of follower-local 
progress; test create-log plus ALTER-log replay with no preseeded progress.



##########
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));

Review Comment:
   [P2] Filter cloud progress with set membership. For a job pinned to 
thousands of partitions, getPartitionIdToOffset(customKafkaPartitions) scans 
every progress entry against the whole ArrayList, and the following cache 
retainAll also checks membership in that list. Recovery or rescheduling runs 
this under the job write lock in the single RoutineLoadScheduler loop, so one 
large job can delay unrelated jobs quadratically. Build one set of pinned IDs 
and filter both maps in linear time.



##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java:
##########
@@ -897,6 +907,19 @@ private void modifyPropertiesInternal(Map<String, String> 
jobProperties,
                 this.id, jobProperties, dataSourceProperties);
     }
 
+    private void replaceKafkaPartitionsAndOffsets(List<Pair<Integer, Long>> 
kafkaPartitionOffsets) {
+        List<Integer> alteredPartitions = 
Lists.newArrayListWithCapacity(kafkaPartitionOffsets.size());
+        Map<Integer, Long> alteredProgress = 
Maps.newHashMapWithExpectedSize(kafkaPartitionOffsets.size());
+        for (Pair<Integer, Long> partitionOffset : kafkaPartitionOffsets) {
+            alteredPartitions.add(partitionOffset.first);
+            alteredProgress.put(partitionOffset.first, partitionOffset.second);
+        }
+        customKafkaPartitions = alteredPartitions;

Review Comment:
   [P2] Validate the pin against the effective Kafka source. If an 
auto-discovered job on topic A has partitions {0,1}, ALTER can set topic B 
(only partition 0) with kafka_partitions="1" and numeric kafka_offsets="2": 
numeric conversion makes no metadata request, while checkPartitions sees 1 in 
A's progress. This assignment pins [1], so resume skips discovery and schedules 
nonexistent partition 1 in B. A later topic or broker change also keeps the 
pin. Validate pinned IDs against the post-ALTER topic and broker whenever the 
source changes, and test a changed-topic ALTER.



##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java:
##########
@@ -897,6 +907,19 @@ private void modifyPropertiesInternal(Map<String, String> 
jobProperties,
                 this.id, jobProperties, dataSourceProperties);
     }
 
+    private void replaceKafkaPartitionsAndOffsets(List<Pair<Integer, Long>> 
kafkaPartitionOffsets) {
+        List<Integer> alteredPartitions = 
Lists.newArrayListWithCapacity(kafkaPartitionOffsets.size());
+        Map<Integer, Long> alteredProgress = 
Maps.newHashMapWithExpectedSize(kafkaPartitionOffsets.size());
+        for (Pair<Integer, Long> partitionOffset : kafkaPartitionOffsets) {
+            alteredPartitions.add(partitionOffset.first);

Review Comment:
   [P1] Reject repeated partition IDs before scheduling them. An ALTER with 
kafka_partitions="1,1" and kafka_offsets="2,2" passes parsing and 
checkPartitions, then this list stores both copies while alteredProgress has 
one key. calculateCurrentConcurrentTaskNum counts two entries, so 
divideRoutineLoadJob creates two tasks reading partition 1 from the same 
offset; both can load the same records. Reject duplicates or normalize the list 
with explicit handling for conflicting offsets, and cover this ALTER case.



##########
fe/fe-core/src/main/java/org/apache/doris/load/routineload/kafka/KafkaRoutineLoadJob.java:
##########
@@ -897,6 +907,19 @@ private void modifyPropertiesInternal(Map<String, String> 
jobProperties,
                 this.id, jobProperties, dataSourceProperties);
     }
 
+    private void replaceKafkaPartitionsAndOffsets(List<Pair<Integer, Long>> 
kafkaPartitionOffsets) {
+        List<Integer> alteredPartitions = 
Lists.newArrayListWithCapacity(kafkaPartitionOffsets.size());
+        Map<Integer, Long> alteredProgress = 
Maps.newHashMapWithExpectedSize(kafkaPartitionOffsets.size());
+        for (Pair<Integer, Long> partitionOffset : kafkaPartitionOffsets) {
+            alteredPartitions.add(partitionOffset.first);
+            alteredProgress.put(partitionOffset.first, partitionOffset.second);
+        }
+        customKafkaPartitions = alteredPartitions;
+        currentKafkaPartitions = Lists.newArrayList(alteredPartitions);
+        progress = new KafkaProgress(alteredProgress);

Review Comment:
   [P2] Keep a way to reselect partitions in a later ALTER. After ALTER narrows 
{0,1,2} to {1}, replacing progress here removes 0 and 2. A subsequent valid 
ALTER requesting {1,2} fails earlier in modifyPropertiesInternal because 
KafkaProgress.checkPartitions requires 2 to exist in the current progress map. 
This makes the new partition pin irreversible through ALTER. Validate requested 
IDs against the Kafka topic, or retain membership separately from active 
progress, and test two successive ALTERs.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to