This is an automated email from the ASF dual-hosted git repository. arvid pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 5d336edcf2ae397a34cde2b55784695415020dda Author: Arvid Heise <[email protected]> AuthorDate: Mon Nov 25 14:35:47 2024 +0100 [FLINK-36788] Move GlobalCommitter to flink-runtime Small refactor for next commit --- .../sink2 => runtime/operators/sink}/GlobalCommittableWrapper.java | 2 +- .../sink2 => runtime/operators/sink}/GlobalCommitterOperator.java | 3 ++- .../sink2 => runtime/operators/sink}/GlobalCommitterSerializer.java | 2 +- .../runtime/translators/GlobalCommitterTransformationTranslator.java | 2 +- .../api/connector/sink2/CommittableMessageSerializerTest.java | 2 ++ .../api/connector/sink2/CommittableMessageTypeInfoTest.java | 1 + .../operators/sink}/GlobalCommitterOperatorTest.java | 5 ++++- .../operators/sink}/GlobalCommitterSerializerTest.java | 4 +++- .../sink2 => runtime/operators/sink}/IntegerSerializer.java | 2 +- .../sink/committables/CommittableCollectorSerializerTest.java | 2 +- 10 files changed, 17 insertions(+), 8 deletions(-) diff --git a/flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommittableWrapper.java b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommittableWrapper.java similarity index 96% rename from flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommittableWrapper.java rename to flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommittableWrapper.java index c9f6d081510..a74f7c379cc 100644 --- a/flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommittableWrapper.java +++ b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommittableWrapper.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.api.connector.sink2; +package org.apache.flink.streaming.runtime.operators.sink; import org.apache.flink.annotation.Internal; import org.apache.flink.streaming.runtime.operators.sink.committables.CommittableCollector; diff --git a/flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterOperator.java b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterOperator.java similarity index 98% rename from flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterOperator.java rename to flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterOperator.java index 7dc4492e40a..3b80a0870ab 100644 --- a/flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterOperator.java +++ b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterOperator.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.api.connector.sink2; +package org.apache.flink.streaming.runtime.operators.sink; import org.apache.flink.annotation.Internal; import org.apache.flink.api.common.state.ListState; @@ -29,6 +29,7 @@ import org.apache.flink.metrics.groups.SinkCommitterMetricGroup; import org.apache.flink.runtime.metrics.groups.InternalSinkCommitterMetricGroup; import org.apache.flink.runtime.state.StateInitializationContext; import org.apache.flink.runtime.state.StateSnapshotContext; +import org.apache.flink.streaming.api.connector.sink2.CommittableMessage; import org.apache.flink.streaming.api.graph.StreamConfig; import org.apache.flink.streaming.api.operators.AbstractStreamOperator; import org.apache.flink.streaming.api.operators.OneInputStreamOperator; diff --git a/flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterSerializer.java b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterSerializer.java similarity index 99% rename from flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterSerializer.java rename to flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterSerializer.java index a8c5ae54978..4fe1cd35733 100644 --- a/flink-runtime/src/main/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterSerializer.java +++ b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterSerializer.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.api.connector.sink2; +package org.apache.flink.streaming.runtime.operators.sink; import org.apache.flink.annotation.Internal; import org.apache.flink.core.io.SimpleVersionedSerialization; diff --git a/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/translators/GlobalCommitterTransformationTranslator.java b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/translators/GlobalCommitterTransformationTranslator.java index d46e6a29510..e4b3449337d 100644 --- a/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/translators/GlobalCommitterTransformationTranslator.java +++ b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/translators/GlobalCommitterTransformationTranslator.java @@ -22,7 +22,6 @@ import org.apache.flink.annotation.Internal; import org.apache.flink.api.common.typeinfo.Types; import org.apache.flink.api.dag.Transformation; import org.apache.flink.streaming.api.connector.sink2.CommittableMessage; -import org.apache.flink.streaming.api.connector.sink2.GlobalCommitterOperator; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.graph.TransformationTranslator; import org.apache.flink.streaming.api.operators.ChainingStrategy; @@ -31,6 +30,7 @@ import org.apache.flink.streaming.api.transformations.GlobalCommitterTransform; import org.apache.flink.streaming.api.transformations.OneInputTransformation; import org.apache.flink.streaming.api.transformations.PhysicalTransformation; import org.apache.flink.streaming.runtime.operators.sink.CommitterOperatorFactory; +import org.apache.flink.streaming.runtime.operators.sink.GlobalCommitterOperator; import org.apache.flink.streaming.runtime.operators.sink.SinkWriterOperatorFactory; import java.util.ArrayDeque; diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageSerializerTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageSerializerTest.java index 0f84c1c1200..379200339c5 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageSerializerTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageSerializerTest.java @@ -18,6 +18,8 @@ package org.apache.flink.streaming.api.connector.sink2; +import org.apache.flink.streaming.runtime.operators.sink.IntegerSerializer; + import org.junit.jupiter.api.Test; import java.io.IOException; diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageTypeInfoTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageTypeInfoTest.java index 43d64ad0918..c318ab5a4bd 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageTypeInfoTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/CommittableMessageTypeInfoTest.java @@ -19,6 +19,7 @@ package org.apache.flink.streaming.api.connector.sink2; import org.apache.flink.api.common.typeutils.TypeInformationTestBase; +import org.apache.flink.streaming.runtime.operators.sink.IntegerSerializer; /** Test for {@link CommittableMessageTypeInfo}. */ class CommittableMessageTypeInfoTest diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterOperatorTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterOperatorTest.java similarity index 96% rename from flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterOperatorTest.java rename to flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterOperatorTest.java index 5dc646240fc..b2bd505fcff 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterOperatorTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterOperatorTest.java @@ -16,10 +16,13 @@ * limitations under the License. */ -package org.apache.flink.streaming.api.connector.sink2; +package org.apache.flink.streaming.runtime.operators.sink; import org.apache.flink.api.connector.sink2.Committer; import org.apache.flink.runtime.checkpoint.OperatorSubtaskState; +import org.apache.flink.streaming.api.connector.sink2.CommittableMessage; +import org.apache.flink.streaming.api.connector.sink2.CommittableSummary; +import org.apache.flink.streaming.api.connector.sink2.CommittableWithLineage; import org.apache.flink.streaming.runtime.streamrecord.StreamRecord; import org.apache.flink.streaming.util.OneInputStreamOperatorTestHarness; diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterSerializerTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterSerializerTest.java similarity index 96% rename from flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterSerializerTest.java rename to flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterSerializerTest.java index 8c5cc11b905..eb2416386a7 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/GlobalCommitterSerializerTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/GlobalCommitterSerializerTest.java @@ -16,13 +16,15 @@ * limitations under the License. */ -package org.apache.flink.streaming.api.connector.sink2; +package org.apache.flink.streaming.runtime.operators.sink; import org.apache.flink.core.io.SimpleVersionedSerializer; import org.apache.flink.core.memory.DataInputDeserializer; import org.apache.flink.core.memory.DataOutputSerializer; import org.apache.flink.metrics.groups.SinkCommitterMetricGroup; import org.apache.flink.runtime.metrics.groups.MetricsGroupTestUtils; +import org.apache.flink.streaming.api.connector.sink2.CommittableSummary; +import org.apache.flink.streaming.api.connector.sink2.CommittableWithLineage; import org.apache.flink.streaming.runtime.operators.sink.committables.CheckpointCommittableManager; import org.apache.flink.streaming.runtime.operators.sink.committables.CommittableCollector; import org.apache.flink.streaming.runtime.operators.sink.committables.CommittableCollectorSerializer; diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/IntegerSerializer.java b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/IntegerSerializer.java similarity index 96% rename from flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/IntegerSerializer.java rename to flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/IntegerSerializer.java index 4fa4e3a650b..dc59acd349b 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/api/connector/sink2/IntegerSerializer.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/IntegerSerializer.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.api.connector.sink2; +package org.apache.flink.streaming.runtime.operators.sink; import org.apache.flink.core.io.SimpleVersionedSerializer; import org.apache.flink.core.memory.DataInputDeserializer; diff --git a/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorSerializerTest.java b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorSerializerTest.java index 938427ba2d5..a69c7c64355 100644 --- a/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorSerializerTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorSerializerTest.java @@ -26,7 +26,7 @@ import org.apache.flink.metrics.groups.SinkCommitterMetricGroup; import org.apache.flink.runtime.metrics.groups.MetricsGroupTestUtils; import org.apache.flink.streaming.api.connector.sink2.CommittableSummary; import org.apache.flink.streaming.api.connector.sink2.CommittableWithLineage; -import org.apache.flink.streaming.api.connector.sink2.IntegerSerializer; +import org.apache.flink.streaming.runtime.operators.sink.IntegerSerializer; import org.assertj.core.api.InstanceOfAssertFactories; import org.junit.jupiter.api.Test;
