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