Skip to content

Commit 1055e42

Browse files
authored
[flink] Add RuntimeContext adapter for Flink 2.x compatibility (apache#2241)
1 parent f7062ea commit 1055e42

10 files changed

Lines changed: 304 additions & 2 deletions

File tree

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.adapter;
19+
20+
import org.apache.flink.api.common.functions.RuntimeContext;
21+
22+
/**
23+
* An adapter for Flink {@link RuntimeContext} class. The {@link RuntimeContext} class added the
24+
* `getJobInfo` and `getTaskInfo` methods in version 1.19 and deprecated many methods, such as
25+
* `getAttemptNumber`.
26+
*
27+
* <p>TODO: remove this class when no longer support flink 1.18.
28+
*/
29+
public class RuntimeContextAdapter {
30+
31+
public static int getAttemptNumber(RuntimeContext runtimeContext) {
32+
return runtimeContext.getAttemptNumber();
33+
}
34+
}
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.adapter;
19+
20+
import org.apache.flink.api.common.operators.MailboxExecutor;
21+
import org.apache.flink.runtime.operators.coordination.OperatorEventDispatcher;
22+
import org.apache.flink.streaming.api.graph.StreamConfig;
23+
import org.apache.flink.streaming.api.operators.Output;
24+
import org.apache.flink.streaming.api.operators.StreamOperatorParameters;
25+
import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
26+
import org.apache.flink.streaming.runtime.tasks.ProcessingTimeService;
27+
import org.apache.flink.streaming.runtime.tasks.StreamTask;
28+
29+
import java.util.function.Supplier;
30+
31+
/**
32+
* Adapter for {@link StreamOperatorParameters} because the constructor is compatibility in flink
33+
* 1.18 and 1.19. However, this constructor only used in test.
34+
*
35+
* <p>TODO: remove this class when no longer support flink 1.18 and 1.19.
36+
*/
37+
public class StreamOperatorParametersAdapter {
38+
39+
public static <OUT> StreamOperatorParameters<OUT> create(
40+
StreamTask<?, ?> containingTask,
41+
StreamConfig config,
42+
Output<StreamRecord<OUT>> output,
43+
Supplier<ProcessingTimeService> processingTimeServiceFactory,
44+
OperatorEventDispatcher operatorEventDispatcher,
45+
MailboxExecutor mailboxExecutor) {
46+
return new StreamOperatorParameters<>(
47+
containingTask,
48+
config,
49+
output,
50+
processingTimeServiceFactory,
51+
operatorEventDispatcher);
52+
}
53+
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.tiering.committer;
19+
20+
/**
21+
* UT for {@link TieringCommitOperator}. Test the compatibility of the `getAttemptNumber` method in
22+
* flink 1.18.
23+
*/
24+
public class Flink118TieringCommitOperatorTest extends TieringCommitOperatorTest {}
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.adapter;
19+
20+
import org.apache.flink.api.common.operators.MailboxExecutor;
21+
import org.apache.flink.runtime.operators.coordination.OperatorEventDispatcher;
22+
import org.apache.flink.streaming.api.graph.StreamConfig;
23+
import org.apache.flink.streaming.api.operators.Output;
24+
import org.apache.flink.streaming.api.operators.StreamOperatorParameters;
25+
import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
26+
import org.apache.flink.streaming.runtime.tasks.ProcessingTimeService;
27+
import org.apache.flink.streaming.runtime.tasks.StreamTask;
28+
29+
import java.util.function.Supplier;
30+
31+
/**
32+
* Adapter for {@link StreamOperatorParameters} because the constructor is compatibility in flink
33+
* 1.18 and 1.19. However, this constructor only used in test.
34+
*
35+
* <p>TODO: remove this class when no longer support flink 1.18 and 1.19.
36+
*/
37+
public class StreamOperatorParametersAdapter {
38+
39+
public static <OUT> StreamOperatorParameters<OUT> create(
40+
StreamTask<?, ?> containingTask,
41+
StreamConfig config,
42+
Output<StreamRecord<OUT>> output,
43+
Supplier<ProcessingTimeService> processingTimeServiceFactory,
44+
OperatorEventDispatcher operatorEventDispatcher,
45+
MailboxExecutor mailboxExecutor) {
46+
return new StreamOperatorParameters<>(
47+
containingTask,
48+
config,
49+
output,
50+
processingTimeServiceFactory,
51+
operatorEventDispatcher);
52+
}
53+
}
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.tiering.committer;
19+
20+
/**
21+
* UT for {@link TieringCommitOperator}. Test the compatibility of the `getAttemptNumber` method in
22+
* flink 1.19.
23+
*/
24+
public class Flink119TieringCommitOperatorTest extends TieringCommitOperatorTest {}
Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,24 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.tiering.committer;
19+
20+
/**
21+
* UT for {@link TieringCommitOperator}. Test the compatibility of the `getAttemptNumber` method in
22+
* flink 2.2.
23+
*/
24+
public class Flink22TieringCommitOperatorTest extends TieringCommitOperatorTest {}
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.adapter;
19+
20+
import org.apache.flink.api.common.functions.RuntimeContext;
21+
22+
/**
23+
* An adapter for Flink {@link RuntimeContext} class. The {@link RuntimeContext} class added the
24+
* `getJobInfo` and `getTaskInfo` methods in version 1.19 and deprecated many methods, such as
25+
* `getAttemptNumber`.
26+
*
27+
* <p>TODO: remove this class when no longer support flink 1.18.
28+
*/
29+
public class RuntimeContextAdapter {
30+
31+
public static int getAttemptNumber(RuntimeContext runtimeContext) {
32+
return runtimeContext.getTaskInfo().getAttemptNumber();
33+
}
34+
}

fluss-flink/fluss-flink-common/src/main/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperator.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import org.apache.fluss.client.metadata.LakeSnapshot;
2424
import org.apache.fluss.config.Configuration;
2525
import org.apache.fluss.exception.LakeTableSnapshotNotExistException;
26+
import org.apache.fluss.flink.adapter.RuntimeContextAdapter;
2627
import org.apache.fluss.flink.tiering.event.FailedTieringEvent;
2728
import org.apache.fluss.flink.tiering.event.FinishedTieringEvent;
2829
import org.apache.fluss.flink.tiering.event.TieringFailOverEvent;
@@ -132,7 +133,7 @@ public void setup(
132133
StreamConfig config,
133134
Output<StreamRecord<CommittableMessage<Committable>>> output) {
134135
super.setup(containingTask, config, output);
135-
int attemptNumber = getRuntimeContext().getAttemptNumber();
136+
int attemptNumber = RuntimeContextAdapter.getAttemptNumber(getRuntimeContext());
136137
if (attemptNumber > 0) {
137138
LOG.info("Send TieringFailoverEvent, current attempt number: {}", attemptNumber);
138139
// attempt number is greater than zero, the job must failover
Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with
4+
* this work for additional information regarding copyright ownership.
5+
* The ASF licenses this file to You under the Apache License, Version 2.0
6+
* (the "License"); you may not use this file except in compliance with
7+
* the License. You may obtain a copy of the License at
8+
*
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
*
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS,
13+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14+
* See the License for the specific language governing permissions and
15+
* limitations under the License.
16+
*/
17+
18+
package org.apache.fluss.flink.adapter;
19+
20+
import org.apache.flink.api.common.operators.MailboxExecutor;
21+
import org.apache.flink.runtime.operators.coordination.OperatorEventDispatcher;
22+
import org.apache.flink.streaming.api.graph.StreamConfig;
23+
import org.apache.flink.streaming.api.operators.Output;
24+
import org.apache.flink.streaming.api.operators.StreamOperatorParameters;
25+
import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
26+
import org.apache.flink.streaming.runtime.tasks.ProcessingTimeService;
27+
import org.apache.flink.streaming.runtime.tasks.StreamTask;
28+
29+
import java.util.function.Supplier;
30+
31+
/**
32+
* Adapter for {@link StreamOperatorParameters} because the constructor is compatibility in flink
33+
* 1.18 and 1.19. However, this constructor only used in test.
34+
*
35+
* <p>TODO: remove this class when no longer support flink 1.18 and 1.19.
36+
*/
37+
public class StreamOperatorParametersAdapter {
38+
39+
public static <OUT> StreamOperatorParameters<OUT> create(
40+
StreamTask<?, ?> containingTask,
41+
StreamConfig config,
42+
Output<StreamRecord<OUT>> output,
43+
Supplier<ProcessingTimeService> processingTimeServiceFactory,
44+
OperatorEventDispatcher operatorEventDispatcher,
45+
MailboxExecutor mailboxExecutor) {
46+
return new StreamOperatorParameters<>(
47+
containingTask,
48+
config,
49+
output,
50+
processingTimeServiceFactory,
51+
operatorEventDispatcher,
52+
mailboxExecutor);
53+
}
54+
}

fluss-flink/fluss-flink-common/src/test/java/org/apache/fluss/flink/tiering/committer/TieringCommitOperatorTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import org.apache.fluss.client.metadata.LakeSnapshot;
2121
import org.apache.fluss.exception.LakeTableSnapshotNotExistException;
22+
import org.apache.fluss.flink.adapter.StreamOperatorParametersAdapter;
2223
import org.apache.fluss.flink.tiering.TestingLakeTieringFactory;
2324
import org.apache.fluss.flink.tiering.TestingWriteResult;
2425
import org.apache.fluss.flink.tiering.event.FailedTieringEvent;
@@ -76,7 +77,7 @@ void beforeEach() throws Exception {
7677
MockOperatorEventDispatcher mockOperatorEventDispatcher =
7778
new MockOperatorEventDispatcher(mockOperatorEventGateway);
7879
parameters =
79-
new StreamOperatorParameters<>(
80+
StreamOperatorParametersAdapter.create(
8081
new SourceOperatorStreamTask<String>(new DummyEnvironment()),
8182
new MockStreamConfig(new Configuration(), 1),
8283
new MockOutput<>(new ArrayList<>()),

0 commit comments

Comments
 (0)