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;
+        }
+    }
+}

Reply via email to