This is an automated email from the ASF dual-hosted git repository. leonard pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
commit 9a9603413c7c6d7094a30f8c8e676b65b0d3b9ba Author: yuxiqian <[email protected]> AuthorDate: Wed Jun 5 12:04:48 2024 +0800 [FLINK-35415][base] Fix compatibility with Flink 1.19 --- .../flink/translator/DataSinkTranslator.java | 45 +++++++++++++++++++--- 1 file changed, 39 insertions(+), 6 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 a6188d2df..ec0bcb593 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 @@ -30,6 +30,7 @@ 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.core.io.SimpleVersionedSerializer; import org.apache.flink.runtime.jobgraph.OperatorID; import org.apache.flink.streaming.api.connector.sink2.CommittableMessage; import org.apache.flink.streaming.api.connector.sink2.CommittableMessageTypeInfo; @@ -39,9 +40,11 @@ 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.operators.OneInputStreamOperatorFactory; 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; + +import java.lang.reflect.InvocationTargetException; /** Translator used to build {@link DataSink} for given {@link DataStream}. */ @Internal @@ -117,10 +120,8 @@ public class DataSinkTranslator { DataStream<Event> inputStream, String sinkName, OperatorID schemaOperatorID) { - TwoPhaseCommittingSink<Event, CommT> committingSink = - (TwoPhaseCommittingSink<Event, CommT>) sink; TypeInformation<CommittableMessage<CommT>> typeInformation = - CommittableMessageTypeInfo.of(committingSink::getCommittableSerializer); + CommittableMessageTypeInfo.of(() -> getCommittableSerializer(sink)); DataStream<CommittableMessage<CommT>> written = inputStream.transform( SINK_WRITER_PREFIX + sinkName, @@ -140,8 +141,7 @@ public class DataSinkTranslator { preCommitted.transform( SINK_COMMITTER_PREFIX + sinkName, typeInformation, - new CommitterOperatorFactory<>( - committingSink, isBatchMode, isCheckpointingEnabled)); + getCommitterOperatorFactory(sink, isBatchMode, isCheckpointingEnabled)); if (sink instanceof WithPostCommitTopology) { ((WithPostCommitTopology<Event, CommT>) sink).addPostCommitTopology(committed); @@ -152,4 +152,37 @@ public class DataSinkTranslator { return sinkDef.getName() .orElse(String.format("Flink CDC Event Sink: %s", sinkDef.getType())); } + + private static <CommT> SimpleVersionedSerializer<CommT> getCommittableSerializer(Object sink) { + // FIX ME: TwoPhaseCommittingSink has been deprecated, and its signature has changed + // during Flink 1.18 to 1.19. Remove this when Flink 1.18 is no longer supported. + try { + return (SimpleVersionedSerializer<CommT>) + sink.getClass().getDeclaredMethod("getCommittableSerializer").invoke(sink); + } catch (NoSuchMethodException | IllegalAccessException | InvocationTargetException e) { + throw new RuntimeException("Failed to get CommittableSerializer", e); + } + } + + private static <CommT> + OneInputStreamOperatorFactory<CommittableMessage<CommT>, CommittableMessage<CommT>> + getCommitterOperatorFactory( + Sink<Event> sink, boolean isBatchMode, boolean isCheckpointingEnabled) { + // FIX ME: OneInputStreamOperatorFactory is an @Internal class, and its signature has + // changed during Flink 1.18 to 1.19. Remove this when Flink 1.18 is no longer supported. + try { + return (OneInputStreamOperatorFactory< + CommittableMessage<CommT>, CommittableMessage<CommT>>) + Class.forName( + "org.apache.flink.streaming.runtime.operators.sink.CommitterOperatorFactory") + .getDeclaredConstructors()[0] + .newInstance(sink, isBatchMode, isCheckpointingEnabled); + + } catch (ClassNotFoundException + | InstantiationException + | IllegalAccessException + | InvocationTargetException e) { + throw new RuntimeException("Failed to create CommitterOperatorFactory", e); + } + } }
