This is an automated email from the ASF dual-hosted git repository.
jiabaosun pushed a commit to branch release-3.1
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
The following commit(s) were added to refs/heads/release-3.1 by this push:
new fa9fb0b1c [FLINK-35149][cdc-composer] Fix DataSinkTranslator#sinkTo
ignoring pre-write topology if not TwoPhaseCommittingSink
fa9fb0b1c is described below
commit fa9fb0b1c49848e77c211a5913d7f28c33e04ff0
Author: Hongshun Wang <[email protected]>
AuthorDate: Thu Jun 6 17:05:49 2024 +0800
[FLINK-35149][cdc-composer] Fix DataSinkTranslator#sinkTo ignoring
pre-write topology if not TwoPhaseCommittingSink
---
.../flink/translator/DataSinkTranslator.java | 6 +-
.../flink/translator/DataSinkTranslatorTest.java | 89 ++++++++++++++++++++++
2 files changed, 93 insertions(+), 2 deletions(-)
diff --git
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
index ec165dd56..a6188d2df 100644
---
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
+++
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
@@ -21,6 +21,7 @@ import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.connector.sink2.Sink;
import org.apache.flink.api.connector.sink2.TwoPhaseCommittingSink;
import org.apache.flink.cdc.common.annotation.Internal;
+import org.apache.flink.cdc.common.annotation.VisibleForTesting;
import org.apache.flink.cdc.common.event.Event;
import org.apache.flink.cdc.common.sink.DataSink;
import org.apache.flink.cdc.common.sink.EventSinkProvider;
@@ -71,7 +72,8 @@ public class DataSinkTranslator {
}
}
- private void sinkTo(
+ @VisibleForTesting
+ void sinkTo(
DataStream<Event> input,
Sink<Event> sink,
String sinkName,
@@ -85,7 +87,7 @@ public class DataSinkTranslator {
if (sink instanceof TwoPhaseCommittingSink) {
addCommittingTopology(sink, stream, sinkName, schemaOperatorID);
} else {
- input.transform(
+ stream.transform(
SINK_WRITER_PREFIX + sinkName,
CommittableMessageTypeInfo.noOutput(),
new DataSinkWriterOperatorFactory<>(sink,
schemaOperatorID));
diff --git
a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslatorTest.java
b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslatorTest.java
new file mode 100644
index 000000000..d2b8206d5
--- /dev/null
+++
b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslatorTest.java
@@ -0,0 +1,89 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.cdc.composer.flink.translator;
+
+import org.apache.flink.api.connector.sink2.SinkWriter;
+import org.apache.flink.api.dag.Transformation;
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.runtime.jobgraph.OperatorID;
+import org.apache.flink.streaming.api.connector.sink2.WithPreWriteTopology;
+import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.datastream.DataStreamSource;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.transformations.OneInputTransformation;
+
+import org.apache.flink.shaded.guava31.com.google.common.collect.Lists;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.io.IOException;
+import java.util.ArrayList;
+
+/** A test for {@link DataSinkTranslator}. */
+class DataSinkTranslatorTest {
+
+ @Test
+ void testPreWriteWithoutCommitSink() {
+ StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
+ ArrayList<Event> mockEvents = Lists.newArrayList(new EmptyEvent(), new
EmptyEvent());
+ DataStreamSource<Event> inputStream = env.fromCollection(mockEvents);
+ DataSinkTranslator translator = new DataSinkTranslator();
+
+ // Node hash must be a 32 character String that describes a hex code
+ String uid = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
+ MockPreWriteWithoutCommitSink mockPreWriteWithoutCommitSink =
+ new MockPreWriteWithoutCommitSink(uid);
+ translator.sinkTo(
+ inputStream,
+ mockPreWriteWithoutCommitSink,
+ "testPreWriteWithoutCommitSink",
+ new OperatorID());
+
+ // Check if the `addPreWriteTopology` is called, and the uid is set
when the transformation
+ // added
+ OneInputTransformation<Event, Event> oneInputTransformation =
+ (OneInputTransformation) env.getTransformations().get(0);
+ Transformation<?> reblanceTransformation =
oneInputTransformation.getInputs().get(0);
+ Assertions.assertEquals(uid,
reblanceTransformation.getUserProvidedNodeHash());
+ }
+
+ private static class EmptyEvent implements Event {}
+
+ private static class MockPreWriteWithoutCommitSink implements
WithPreWriteTopology<Event> {
+
+ private final String uid;
+
+ public MockPreWriteWithoutCommitSink(String uid) {
+ this.uid = uid;
+ }
+
+ @Override
+ public DataStream<Event> addPreWriteTopology(DataStream<Event>
inputDataStream) {
+ // return a new DataSteam with specified uid
+ DataStream<Event> rebalance = inputDataStream.rebalance();
+ rebalance.getTransformation().setUidHash(uid);
+ return rebalance;
+ }
+
+ @Override
+ public SinkWriter<Event> createWriter(InitContext context) throws
IOException {
+ return null;
+ }
+ }
+}