This is an automated email from the ASF dual-hosted git repository.

guoweijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git

commit 94a6ed03bffda03cdaff99e8bdb7b45da46e2c6a
Author: Xu Huang <[email protected]>
AuthorDate: Thu Nov 28 15:00:49 2024 +0800

    [FLINK-36815][API] Introduce subtask index and attempt id in TaskInfo for 
DSv2
---
 .../flink/datastream/api/context/TaskInfo.java     | 17 ++++++++++++++++
 .../impl/context/DefaultRuntimeContext.java        |  6 +++++-
 .../datastream/impl/context/DefaultTaskInfo.java   | 23 +++++++++++++++++++++-
 .../datastream/impl/operators/ProcessOperator.java |  2 ++
 .../TwoInputBroadcastProcessOperator.java          |  2 ++
 .../TwoInputNonBroadcastProcessOperator.java       |  2 ++
 .../impl/operators/TwoOutputProcessOperator.java   |  2 ++
 .../context/DefaultNonPartitionedContextTest.java  |  4 ++++
 .../impl/context/DefaultTaskInfoTest.java          |  4 +++-
 .../DefaultTwoOutputNonPartitionedContextTest.java |  4 ++++
 10 files changed, 63 insertions(+), 3 deletions(-)

diff --git 
a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/TaskInfo.java
 
b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/TaskInfo.java
index 04b8d67bcbd..ad0f6890d85 100644
--- 
a/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/TaskInfo.java
+++ 
b/flink-datastream-api/src/main/java/org/apache/flink/datastream/api/context/TaskInfo.java
@@ -43,4 +43,21 @@ public interface TaskInfo {
      * @return The name of current task.
      */
     String getTaskName();
+
+    /**
+     * Gets the number of this parallel subtask. The numbering starts from 0 
and goes up to
+     * parallelism-1 (parallelism as returned by {@link #getParallelism()}).
+     *
+     * @return The index of the parallel subtask.
+     */
+    int getIndexOfThisSubtask();
+
+    /**
+     * Gets the attempt number of this parallel subtask. First attempt is 
numbered 0. The attempt
+     * number corresponds to the number of times this task has been 
restarted(after
+     * failure/cancellation) since the job was initially started.
+     *
+     * @return The attempt number of the subtask.
+     */
+    int getAttemptNumber();
 }
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultRuntimeContext.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultRuntimeContext.java
index 54d76d7a8d3..53ba8f5605c 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultRuntimeContext.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultRuntimeContext.java
@@ -38,9 +38,13 @@ public class DefaultRuntimeContext implements RuntimeContext 
{
             int parallelism,
             int maxParallelism,
             String taskName,
+            int indexOfSubtask,
+            int attemptNumber,
             MetricGroup metricGroup) {
         this.jobInfo = new DefaultJobInfo(jobName, jobType);
-        this.taskInfo = new DefaultTaskInfo(parallelism, maxParallelism, 
taskName);
+        this.taskInfo =
+                new DefaultTaskInfo(
+                        parallelism, maxParallelism, taskName, indexOfSubtask, 
attemptNumber);
         this.metricGroup = metricGroup;
     }
 
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultTaskInfo.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultTaskInfo.java
index c2589e93148..fe621d113fb 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultTaskInfo.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/context/DefaultTaskInfo.java
@@ -28,10 +28,21 @@ public class DefaultTaskInfo implements TaskInfo {
 
     private final String taskName;
 
-    public DefaultTaskInfo(int parallelism, int maxParallelism, String 
taskName) {
+    private final int indexOfSubtask;
+
+    private final int attemptNumber;
+
+    public DefaultTaskInfo(
+            int parallelism,
+            int maxParallelism,
+            String taskName,
+            int indexOfSubtask,
+            int attemptNumber) {
         this.parallelism = parallelism;
         this.maxParallelism = maxParallelism;
         this.taskName = taskName;
+        this.indexOfSubtask = indexOfSubtask;
+        this.attemptNumber = attemptNumber;
     }
 
     @Override
@@ -48,4 +59,14 @@ public class DefaultTaskInfo implements TaskInfo {
     public String getTaskName() {
         return taskName;
     }
+
+    @Override
+    public int getIndexOfThisSubtask() {
+        return indexOfSubtask;
+    }
+
+    @Override
+    public int getAttemptNumber() {
+        return attemptNumber;
+    }
 }
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java
index 5f35a570214..10876fbeaf2 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/ProcessOperator.java
@@ -63,6 +63,8 @@ public class ProcessOperator<IN, OUT>
                         taskInfo.getNumberOfParallelSubtasks(),
                         taskInfo.getMaxNumberOfParallelSubtasks(),
                         taskInfo.getTaskName(),
+                        taskInfo.getIndexOfThisSubtask(),
+                        taskInfo.getAttemptNumber(),
                         operatorContext.getMetricGroup());
         partitionedContext =
                 new DefaultPartitionedContext(
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java
index df3d7465ee4..dea6d2cd2d0 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputBroadcastProcessOperator.java
@@ -68,6 +68,8 @@ public class TwoInputBroadcastProcessOperator<IN1, IN2, OUT>
                         taskInfo.getNumberOfParallelSubtasks(),
                         taskInfo.getMaxNumberOfParallelSubtasks(),
                         taskInfo.getTaskName(),
+                        taskInfo.getIndexOfThisSubtask(),
+                        taskInfo.getAttemptNumber(),
                         operatorContext.getMetricGroup());
         this.partitionedContext =
                 new DefaultPartitionedContext(
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java
index 9501aaacf08..c3cfb438792 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoInputNonBroadcastProcessOperator.java
@@ -71,6 +71,8 @@ public class TwoInputNonBroadcastProcessOperator<IN1, IN2, 
OUT>
                         taskInfo.getNumberOfParallelSubtasks(),
                         taskInfo.getMaxNumberOfParallelSubtasks(),
                         taskInfo.getTaskName(),
+                        taskInfo.getIndexOfThisSubtask(),
+                        taskInfo.getAttemptNumber(),
                         operatorContext.getMetricGroup());
         this.partitionedContext =
                 new DefaultPartitionedContext(
diff --git 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java
 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java
index bae23887c24..3f4f8e05bec 100644
--- 
a/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java
+++ 
b/flink-datastream/src/main/java/org/apache/flink/datastream/impl/operators/TwoOutputProcessOperator.java
@@ -80,6 +80,8 @@ public class TwoOutputProcessOperator<IN, OUT_MAIN, OUT_SIDE>
                         taskInfo.getNumberOfParallelSubtasks(),
                         taskInfo.getMaxNumberOfParallelSubtasks(),
                         taskInfo.getTaskName(),
+                        taskInfo.getIndexOfThisSubtask(),
+                        taskInfo.getAttemptNumber(),
                         operatorContext.getMetricGroup());
         this.partitionedContext =
                 new DefaultPartitionedContext(
diff --git 
a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultNonPartitionedContextTest.java
 
b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultNonPartitionedContextTest.java
index b4fa8c62531..e5e36937a4d 100644
--- 
a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultNonPartitionedContextTest.java
+++ 
b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultNonPartitionedContextTest.java
@@ -55,6 +55,8 @@ class DefaultNonPartitionedContextTest {
                         1,
                         2,
                         "mock-task",
+                        0,
+                        0,
                         operatorRuntimeContext.getMetricGroup());
 
         DefaultNonPartitionedContext<Integer> nonPartitionedContext =
@@ -105,6 +107,8 @@ class DefaultNonPartitionedContextTest {
                         1,
                         2,
                         "mock-task",
+                        0,
+                        0,
                         operatorRuntimeContext.getMetricGroup());
         DefaultNonPartitionedContext<Integer> nonPartitionedContext =
                 new DefaultNonPartitionedContext<>(
diff --git 
a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTaskInfoTest.java
 
b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTaskInfoTest.java
index 2f90511675d..1a1264507d0 100644
--- 
a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTaskInfoTest.java
+++ 
b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTaskInfoTest.java
@@ -26,9 +26,11 @@ import static org.assertj.core.api.Assertions.assertThat;
 class DefaultTaskInfoTest {
     @Test
     void testTaskInfo() {
-        DefaultTaskInfo taskInfo = new DefaultTaskInfo(1, 2, "taskName");
+        DefaultTaskInfo taskInfo = new DefaultTaskInfo(1, 2, "taskName", 0, 0);
         assertThat(taskInfo.getParallelism()).isEqualTo(1);
         assertThat(taskInfo.getMaxParallelism()).isEqualTo(2);
         assertThat(taskInfo.getTaskName()).isEqualTo("taskName");
+        assertThat(taskInfo.getIndexOfThisSubtask()).isEqualTo(0);
+        assertThat(taskInfo.getAttemptNumber()).isEqualTo(0);
     }
 }
diff --git 
a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTwoOutputNonPartitionedContextTest.java
 
b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTwoOutputNonPartitionedContextTest.java
index 137527b6542..cd101bfbb10 100644
--- 
a/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTwoOutputNonPartitionedContextTest.java
+++ 
b/flink-datastream/src/test/java/org/apache/flink/datastream/impl/context/DefaultTwoOutputNonPartitionedContextTest.java
@@ -60,6 +60,8 @@ class DefaultTwoOutputNonPartitionedContextTest {
                         1,
                         2,
                         "mock-task",
+                        0,
+                        0,
                         operatorRuntimeContext.getMetricGroup());
         DefaultTwoOutputNonPartitionedContext<Integer, Long> 
nonPartitionedContext =
                 new DefaultTwoOutputNonPartitionedContext<>(
@@ -118,6 +120,8 @@ class DefaultTwoOutputNonPartitionedContextTest {
                         1,
                         2,
                         "mock-task",
+                        0,
+                        0,
                         operatorRuntimeContext.getMetricGroup());
         DefaultTwoOutputNonPartitionedContext<Integer, Long> 
nonPartitionedContext =
                 new DefaultTwoOutputNonPartitionedContext<>(

Reply via email to