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

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


The following commit(s) were added to refs/heads/master by this push:
     new 8c5437a1f [FLINK-34952][cdc-composer][sink] Flink CDC pipeline 
supports SinkFunction (#3204)
8c5437a1f is described below

commit 8c5437a1f110b337defa42275a16037f9f169993
Author: Hongshun Wang <[email protected]>
AuthorDate: Tue Apr 9 15:27:55 2024 +0800

    [FLINK-34952][cdc-composer][sink] Flink CDC pipeline supports SinkFunction 
(#3204)
---
 .../flink/translator/DataSinkTranslator.java       |  30 +++++
 .../flink/FlinkPipelineComposerITCase.java         |  34 ++++--
 .../values/factory/ValuesDataFactory.java          |   3 +-
 .../cdc/connectors/values/sink/ValuesDataSink.java |  22 +++-
 .../values/sink/ValuesDataSinkFunction.java        |  98 +++++++++++++++
 .../values/sink/ValuesDataSinkOptions.java         |   7 ++
 .../operators/sink/DataSinkFunctionOperator.java   | 131 +++++++++++++++++++++
 7 files changed, 311 insertions(+), 14 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 8bf3ef88d..ec165dd56 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
@@ -24,8 +24,10 @@ import org.apache.flink.cdc.common.annotation.Internal;
 import org.apache.flink.cdc.common.event.Event;
 import org.apache.flink.cdc.common.sink.DataSink;
 import org.apache.flink.cdc.common.sink.EventSinkProvider;
+import org.apache.flink.cdc.common.sink.FlinkSinkFunctionProvider;
 import org.apache.flink.cdc.common.sink.FlinkSinkProvider;
 import org.apache.flink.cdc.composer.definition.SinkDef;
+import org.apache.flink.cdc.runtime.operators.sink.DataSinkFunctionOperator;
 import 
org.apache.flink.cdc.runtime.operators.sink.DataSinkWriterOperatorFactory;
 import org.apache.flink.runtime.jobgraph.OperatorID;
 import org.apache.flink.streaming.api.connector.sink2.CommittableMessage;
@@ -34,6 +36,10 @@ import 
org.apache.flink.streaming.api.connector.sink2.WithPostCommitTopology;
 import org.apache.flink.streaming.api.connector.sink2.WithPreCommitTopology;
 import org.apache.flink.streaming.api.connector.sink2.WithPreWriteTopology;
 import org.apache.flink.streaming.api.datastream.DataStream;
+import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
+import org.apache.flink.streaming.api.functions.sink.SinkFunction;
+import org.apache.flink.streaming.api.transformations.LegacySinkTransformation;
+import org.apache.flink.streaming.api.transformations.PhysicalTransformation;
 import 
org.apache.flink.streaming.runtime.operators.sink.CommitterOperatorFactory;
 
 /** Translator used to build {@link DataSink} for given {@link DataStream}. */
@@ -56,6 +62,12 @@ public class DataSinkTranslator {
             FlinkSinkProvider sinkProvider = (FlinkSinkProvider) 
eventSinkProvider;
             Sink<Event> sink = sinkProvider.getSink();
             sinkTo(input, sink, sinkName, schemaOperatorID);
+        } else if (eventSinkProvider instanceof FlinkSinkFunctionProvider) {
+            // SinkFunction
+            FlinkSinkFunctionProvider sinkFunctionProvider =
+                    (FlinkSinkFunctionProvider) eventSinkProvider;
+            SinkFunction<Event> sinkFunction = 
sinkFunctionProvider.getSinkFunction();
+            sinkTo(input, sinkFunction, sinkName, schemaOperatorID);
         }
     }
 
@@ -80,6 +92,24 @@ public class DataSinkTranslator {
         }
     }
 
+    private void sinkTo(
+            DataStream<Event> input,
+            SinkFunction<Event> sinkFunction,
+            String sinkName,
+            OperatorID schemaOperatorID) {
+        DataSinkFunctionOperator sinkOperator =
+                new DataSinkFunctionOperator(sinkFunction, schemaOperatorID);
+        final StreamExecutionEnvironment executionEnvironment = 
input.getExecutionEnvironment();
+        PhysicalTransformation<Event> transformation =
+                new LegacySinkTransformation<>(
+                        input.getTransformation(),
+                        SINK_WRITER_PREFIX + sinkName,
+                        sinkOperator,
+                        executionEnvironment.getParallelism(),
+                        false);
+        executionEnvironment.addOperator(transformation);
+    }
+
     private <CommT> void addCommittingTopology(
             Sink<Event> sink,
             DataStream<Event> inputStream,
diff --git 
a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposerITCase.java
 
b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposerITCase.java
index 2b3d18ca5..ed7eaee3c 100644
--- 
a/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposerITCase.java
+++ 
b/flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposerITCase.java
@@ -26,6 +26,7 @@ import org.apache.flink.cdc.composer.definition.SourceDef;
 import org.apache.flink.cdc.composer.definition.TransformDef;
 import org.apache.flink.cdc.connectors.values.ValuesDatabase;
 import org.apache.flink.cdc.connectors.values.factory.ValuesDataFactory;
+import org.apache.flink.cdc.connectors.values.sink.ValuesDataSink;
 import org.apache.flink.cdc.connectors.values.sink.ValuesDataSinkOptions;
 import org.apache.flink.cdc.connectors.values.source.ValuesDataSourceHelper;
 import org.apache.flink.cdc.connectors.values.source.ValuesDataSourceOptions;
@@ -34,8 +35,9 @@ import org.apache.flink.test.junit5.MiniClusterExtension;
 
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
-import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.RegisterExtension;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.EnumSource;
 
 import java.io.ByteArrayOutputStream;
 import java.io.PrintStream;
@@ -94,8 +96,9 @@ class FlinkPipelineComposerITCase {
         System.setOut(standardOut);
     }
 
-    @Test
-    void testSingleSplitSingleTable() throws Exception {
+    @ParameterizedTest
+    @EnumSource
+    void testSingleSplitSingleTable(ValuesDataSink.SinkApi sinkApi) throws 
Exception {
         FlinkPipelineComposer composer = FlinkPipelineComposer.ofMiniCluster();
 
         // Setup value source
@@ -109,6 +112,7 @@ class FlinkPipelineComposerITCase {
         // Setup value sink
         Configuration sinkConfig = new Configuration();
         sinkConfig.set(ValuesDataSinkOptions.MATERIALIZED_IN_MEMORY, true);
+        sinkConfig.set(ValuesDataSinkOptions.SINK_API, sinkApi);
         SinkDef sinkDef = new SinkDef(ValuesDataFactory.IDENTIFIER, "Value 
Sink", sinkConfig);
 
         // Setup pipeline
@@ -148,8 +152,9 @@ class FlinkPipelineComposerITCase {
                         
"DataChangeEvent{tableId=default_namespace.default_schema.table1, before=[2, ], 
after=[2, x], op=UPDATE, meta=()}");
     }
 
-    @Test
-    void testSingleSplitMultipleTables() throws Exception {
+    @ParameterizedTest
+    @EnumSource
+    void testSingleSplitMultipleTables(ValuesDataSink.SinkApi sinkApi) throws 
Exception {
         FlinkPipelineComposer composer = FlinkPipelineComposer.ofMiniCluster();
 
         // Setup value source
@@ -163,6 +168,7 @@ class FlinkPipelineComposerITCase {
         // Setup value sink
         Configuration sinkConfig = new Configuration();
         sinkConfig.set(ValuesDataSinkOptions.MATERIALIZED_IN_MEMORY, true);
+        sinkConfig.set(ValuesDataSinkOptions.SINK_API, sinkApi);
         SinkDef sinkDef = new SinkDef(ValuesDataFactory.IDENTIFIER, "Value 
Sink", sinkConfig);
 
         // Setup pipeline
@@ -212,8 +218,9 @@ class FlinkPipelineComposerITCase {
                         
"DataChangeEvent{tableId=default_namespace.default_schema.table1, before=[2, 
2], after=[2, x], op=UPDATE, meta=()}");
     }
 
-    @Test
-    void testMultiSplitsSingleTable() throws Exception {
+    @ParameterizedTest
+    @EnumSource
+    void testMultiSplitsSingleTable(ValuesDataSink.SinkApi sinkApi) throws 
Exception {
         FlinkPipelineComposer composer = FlinkPipelineComposer.ofMiniCluster();
 
         // Setup value source
@@ -227,6 +234,7 @@ class FlinkPipelineComposerITCase {
         // Setup value sink
         Configuration sinkConfig = new Configuration();
         sinkConfig.set(ValuesDataSinkOptions.MATERIALIZED_IN_MEMORY, true);
+        sinkConfig.set(ValuesDataSinkOptions.SINK_API, sinkApi);
         SinkDef sinkDef = new SinkDef(ValuesDataFactory.IDENTIFIER, "Value 
Sink", sinkConfig);
 
         // Setup pipeline
@@ -253,8 +261,9 @@ class FlinkPipelineComposerITCase {
                         
"default_namespace.default_schema.table1:col1=5;col2=5;col3=");
     }
 
-    @Test
-    void testTransform() throws Exception {
+    @ParameterizedTest
+    @EnumSource
+    void testTransform(ValuesDataSink.SinkApi sinkApi) throws Exception {
         FlinkPipelineComposer composer = FlinkPipelineComposer.ofMiniCluster();
 
         // Setup value source
@@ -268,6 +277,7 @@ class FlinkPipelineComposerITCase {
         // Setup value sink
         Configuration sinkConfig = new Configuration();
         sinkConfig.set(ValuesDataSinkOptions.MATERIALIZED_IN_MEMORY, true);
+        sinkConfig.set(ValuesDataSinkOptions.SINK_API, sinkApi);
         SinkDef sinkDef = new SinkDef(ValuesDataFactory.IDENTIFIER, "Value 
Sink", sinkConfig);
 
         // Setup transform
@@ -310,8 +320,9 @@ class FlinkPipelineComposerITCase {
                         
"DataChangeEvent{tableId=default_namespace.default_schema.table1, before=[2, , 
20], after=[2, x, 20], op=UPDATE, meta=()}");
     }
 
-    @Test
-    void testTransformTwice() throws Exception {
+    @ParameterizedTest
+    @EnumSource
+    void testTransformTwice(ValuesDataSink.SinkApi sinkApi) throws Exception {
         FlinkPipelineComposer composer = FlinkPipelineComposer.ofMiniCluster();
 
         // Setup value source
@@ -325,6 +336,7 @@ class FlinkPipelineComposerITCase {
         // Setup value sink
         Configuration sinkConfig = new Configuration();
         sinkConfig.set(ValuesDataSinkOptions.MATERIALIZED_IN_MEMORY, true);
+        sinkConfig.set(ValuesDataSinkOptions.SINK_API, sinkApi);
         SinkDef sinkDef = new SinkDef(ValuesDataFactory.IDENTIFIER, "Value 
Sink", sinkConfig);
 
         // Setup transform
diff --git 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/factory/ValuesDataFactory.java
 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/factory/ValuesDataFactory.java
index 31a947061..ee8411d2b 100644
--- 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/factory/ValuesDataFactory.java
+++ 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/factory/ValuesDataFactory.java
@@ -53,7 +53,8 @@ public class ValuesDataFactory implements DataSourceFactory, 
DataSinkFactory {
     public DataSink createDataSink(Context context) {
         return new ValuesDataSink(
                 
context.getFactoryConfiguration().get(ValuesDataSinkOptions.MATERIALIZED_IN_MEMORY),
-                
context.getFactoryConfiguration().get(ValuesDataSinkOptions.PRINT_ENABLED));
+                
context.getFactoryConfiguration().get(ValuesDataSinkOptions.PRINT_ENABLED),
+                
context.getFactoryConfiguration().get(ValuesDataSinkOptions.SINK_API));
     }
 
     @Override
diff --git 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSink.java
 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSink.java
index 2f0b3627e..d6789452d 100644
--- 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSink.java
+++ 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSink.java
@@ -30,6 +30,7 @@ import org.apache.flink.cdc.common.event.TableId;
 import org.apache.flink.cdc.common.schema.Schema;
 import org.apache.flink.cdc.common.sink.DataSink;
 import org.apache.flink.cdc.common.sink.EventSinkProvider;
+import org.apache.flink.cdc.common.sink.FlinkSinkFunctionProvider;
 import org.apache.flink.cdc.common.sink.FlinkSinkProvider;
 import org.apache.flink.cdc.common.sink.MetadataApplier;
 import org.apache.flink.cdc.common.utils.SchemaUtils;
@@ -49,14 +50,22 @@ public class ValuesDataSink implements DataSink, 
Serializable {
 
     private final boolean print;
 
-    public ValuesDataSink(boolean materializedInMemory, boolean print) {
+    private final SinkApi sinkApi;
+
+    public ValuesDataSink(boolean materializedInMemory, boolean print, SinkApi 
sinkApi) {
         this.materializedInMemory = materializedInMemory;
         this.print = print;
+        this.sinkApi = sinkApi;
     }
 
     @Override
     public EventSinkProvider getEventSinkProvider() {
-        return FlinkSinkProvider.of(new ValuesSink(materializedInMemory, 
print));
+        if (SinkApi.SINK_V2.equals(sinkApi)) {
+            return FlinkSinkProvider.of(new ValuesSink(materializedInMemory, 
print));
+        } else {
+            return FlinkSinkFunctionProvider.of(
+                    new ValuesDataSinkFunction(materializedInMemory, print));
+        }
     }
 
     @Override
@@ -157,4 +166,13 @@ public class ValuesDataSink implements DataSink, 
Serializable {
         @Override
         public void close() {}
     }
+
+    /** SinkApi which sink based on. */
+    public enum SinkApi {
+        /** Sink based on SinkFunction. */
+        SINK_FUNCTION,
+
+        /** Sink based on SinkV2. */
+        SINK_V2;
+    }
 }
diff --git 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkFunction.java
 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkFunction.java
new file mode 100644
index 000000000..f4876d5e0
--- /dev/null
+++ 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkFunction.java
@@ -0,0 +1,98 @@
+/*
+ * 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.connectors.values.sink;
+
+import org.apache.flink.api.common.eventtime.Watermark;
+import org.apache.flink.cdc.common.data.RecordData;
+import org.apache.flink.cdc.common.event.ChangeEvent;
+import org.apache.flink.cdc.common.event.CreateTableEvent;
+import org.apache.flink.cdc.common.event.DataChangeEvent;
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.cdc.common.event.SchemaChangeEvent;
+import org.apache.flink.cdc.common.event.TableId;
+import org.apache.flink.cdc.common.schema.Schema;
+import org.apache.flink.cdc.common.utils.SchemaUtils;
+import org.apache.flink.cdc.connectors.values.ValuesDatabase;
+import org.apache.flink.streaming.api.functions.sink.SinkFunction;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/** An e2e {@link SinkFunction} implementation that print all {@link 
DataChangeEvent} out. */
+public class ValuesDataSinkFunction implements SinkFunction<Event> {
+    private final boolean materializedInMemory;
+
+    private final boolean print;
+
+    /**
+     * keep the relationship of TableId and Schema as write method may rely on 
the schema
+     * information of DataChangeEvent.
+     */
+    private final Map<TableId, Schema> schemaMaps;
+
+    private final Map<TableId, List<RecordData.FieldGetter>> fieldGetterMaps;
+
+    public ValuesDataSinkFunction(boolean materializedInMemory, boolean print) 
{
+        this.materializedInMemory = materializedInMemory;
+        this.print = print;
+        schemaMaps = new HashMap<>();
+        fieldGetterMaps = new HashMap<>();
+    }
+
+    @Override
+    public void invoke(Event event, Context context) throws Exception {
+        if (event instanceof SchemaChangeEvent) {
+            SchemaChangeEvent schemaChangeEvent = (SchemaChangeEvent) event;
+            TableId tableId = schemaChangeEvent.tableId();
+            if (event instanceof CreateTableEvent) {
+                Schema schema = ((CreateTableEvent) event).getSchema();
+                schemaMaps.put(tableId, schema);
+                fieldGetterMaps.put(tableId, 
SchemaUtils.createFieldGetters(schema));
+            } else {
+                if (!schemaMaps.containsKey(tableId)) {
+                    throw new RuntimeException("schema of " + tableId + " is 
not existed.");
+                }
+                Schema schema =
+                        SchemaUtils.applySchemaChangeEvent(
+                                schemaMaps.get(tableId), schemaChangeEvent);
+                schemaMaps.put(tableId, schema);
+                fieldGetterMaps.put(tableId, 
SchemaUtils.createFieldGetters(schema));
+            }
+        } else if (materializedInMemory && event instanceof DataChangeEvent) {
+            ValuesDatabase.applyDataChangeEvent((DataChangeEvent) event);
+        }
+
+        if (print) {
+            // print the detail message to console for verification.
+            System.out.println(
+                    ValuesDataSinkHelper.convertEventToStr(
+                            event, fieldGetterMaps.get(((ChangeEvent) 
event).tableId())));
+        }
+    }
+
+    @Override
+    public void writeWatermark(Watermark watermark) throws Exception {
+        SinkFunction.super.writeWatermark(watermark);
+    }
+
+    @Override
+    public void finish() throws Exception {
+        SinkFunction.super.finish();
+    }
+}
diff --git 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkOptions.java
 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkOptions.java
index 822271db1..11b132a2d 100644
--- 
a/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkOptions.java
+++ 
b/flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-values/src/main/java/org/apache/flink/cdc/connectors/values/sink/ValuesDataSinkOptions.java
@@ -35,4 +35,11 @@ public class ValuesDataSinkOptions {
                     .booleanType()
                     .defaultValue(true)
                     .withDescription("True if the Event should be print to 
console.");
+
+    public static final ConfigOption<ValuesDataSink.SinkApi> SINK_API =
+            ConfigOptions.key("sink.api")
+                    .enumType(ValuesDataSink.SinkApi.class)
+                    .defaultValue(ValuesDataSink.SinkApi.SINK_V2)
+                    .withDescription(
+                            "The sink api on which the sink is based: 
SinkFunction or SinkV2.");
 }
diff --git 
a/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/sink/DataSinkFunctionOperator.java
 
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/sink/DataSinkFunctionOperator.java
new file mode 100644
index 000000000..438c3f302
--- /dev/null
+++ 
b/flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/sink/DataSinkFunctionOperator.java
@@ -0,0 +1,131 @@
+/*
+ * 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.runtime.operators.sink;
+
+import org.apache.flink.cdc.common.annotation.Internal;
+import org.apache.flink.cdc.common.event.ChangeEvent;
+import org.apache.flink.cdc.common.event.CreateTableEvent;
+import org.apache.flink.cdc.common.event.Event;
+import org.apache.flink.cdc.common.event.FlushEvent;
+import org.apache.flink.cdc.common.event.TableId;
+import org.apache.flink.cdc.common.schema.Schema;
+import org.apache.flink.runtime.jobgraph.OperatorID;
+import org.apache.flink.runtime.state.StateInitializationContext;
+import org.apache.flink.streaming.api.functions.sink.SinkFunction;
+import org.apache.flink.streaming.api.graph.StreamConfig;
+import org.apache.flink.streaming.api.operators.ChainingStrategy;
+import org.apache.flink.streaming.api.operators.Output;
+import org.apache.flink.streaming.api.operators.StreamSink;
+import org.apache.flink.streaming.runtime.streamrecord.StreamRecord;
+import org.apache.flink.streaming.runtime.tasks.StreamTask;
+
+import java.util.HashSet;
+import java.util.Optional;
+import java.util.Set;
+
+/**
+ * An operator that processes records to be written into a {@link
+ * org.apache.flink.streaming.api.functions.sink.SinkFunction}.
+ *
+ * <p>The operator is a proxy of {@link 
org.apache.flink.streaming.api.operators.StreamSink} in
+ * Flink.
+ *
+ * <p>The operator is always part of a sink pipeline and is the first operator.
+ */
+@Internal
+public class DataSinkFunctionOperator extends StreamSink<Event> {
+
+    private SchemaEvolutionClient schemaEvolutionClient;
+    private final OperatorID schemaOperatorID;
+    /** A set of {@link TableId} that already processed {@link 
CreateTableEvent}. */
+    private final Set<TableId> processedTableIds;
+
+    public DataSinkFunctionOperator(SinkFunction<Event> userFunction, 
OperatorID schemaOperatorID) {
+        super(userFunction);
+        this.schemaOperatorID = schemaOperatorID;
+        processedTableIds = new HashSet<>();
+        this.chainingStrategy = ChainingStrategy.ALWAYS;
+    }
+
+    @Override
+    public void setup(
+            StreamTask<?, ?> containingTask,
+            StreamConfig config,
+            Output<StreamRecord<Object>> output) {
+        super.setup(containingTask, config, output);
+        schemaEvolutionClient =
+                new SchemaEvolutionClient(
+                        
containingTask.getEnvironment().getOperatorCoordinatorEventGateway(),
+                        schemaOperatorID);
+    }
+
+    @Override
+    public void initializeState(StateInitializationContext context) throws 
Exception {
+        
schemaEvolutionClient.registerSubtask(getRuntimeContext().getIndexOfThisSubtask());
+        super.initializeState(context);
+    }
+
+    @Override
+    public void processElement(StreamRecord<Event> element) throws Exception {
+        Event event = element.getValue();
+
+        // FlushEvent triggers flush
+        if (event instanceof FlushEvent) {
+            handleFlushEvent(((FlushEvent) event));
+            return;
+        }
+
+        // CreateTableEvent marks the table as processed directly
+        if (event instanceof CreateTableEvent) {
+            processedTableIds.add(((CreateTableEvent) event).tableId());
+            super.processElement(element);
+            return;
+        }
+
+        // Check if the table is processed before emitting all other events, 
because we have to make
+        // sure that sink have a view of the full schema before processing any 
change events,
+        // including schema changes.
+        ChangeEvent changeEvent = (ChangeEvent) event;
+        if (!processedTableIds.contains(changeEvent.tableId())) {
+            emitLatestSchema(changeEvent.tableId());
+            processedTableIds.add(changeEvent.tableId());
+        }
+        processedTableIds.add(changeEvent.tableId());
+        super.processElement(element);
+    }
+
+    // ----------------------------- Helper functions 
-------------------------------
+    private void handleFlushEvent(FlushEvent event) throws Exception {
+        userFunction.finish();
+        schemaEvolutionClient.notifyFlushSuccess(
+                getRuntimeContext().getIndexOfThisSubtask(), 
event.getTableId());
+    }
+
+    private void emitLatestSchema(TableId tableId) throws Exception {
+        Optional<Schema> schema = 
schemaEvolutionClient.getLatestSchema(tableId);
+        if (schema.isPresent()) {
+            // request and process CreateTableEvent because SinkFunction need 
to retrieve
+            // Schema to deserialize RecordData after resuming job.
+            super.processElement(new StreamRecord<>(new 
CreateTableEvent(tableId, schema.get())));
+            processedTableIds.add(tableId);
+        } else {
+            throw new RuntimeException(
+                    "Could not find schema message from SchemaRegistry for " + 
tableId);
+        }
+    }
+}

Reply via email to