Skip to content

Commit

Permalink
upgrade cdc to 3.2.0 (DataLinkDC#3878)
Browse files Browse the repository at this point in the history
Co-authored-by: gaoyan1998 <[email protected]>
  • Loading branch information
gaoyan1998 and gaoyan1998 authored Oct 18, 2024
1 parent 51ee0bb commit c081c3d
Show file tree
Hide file tree
Showing 11 changed files with 26 additions and 423 deletions.

This file was deleted.

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -23,12 +23,16 @@
import org.dinky.trans.AbstractOperation;
import org.dinky.trans.Operation;

import org.apache.flink.cdc.cli.parser.YamlPipelineDefinitionParser;
import org.apache.flink.cdc.common.configuration.Configuration;
import org.apache.flink.cdc.composer.PipelineComposer;
import org.apache.flink.cdc.composer.definition.PipelineDef;
import org.apache.flink.cdc.composer.flink.FlinkPipelineComposer;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.TableResult;
import org.apache.flink.table.api.internal.TableResultImpl;

import java.lang.reflect.Constructor;
import java.util.regex.Matcher;
import java.util.regex.Pattern;

Expand Down Expand Up @@ -86,11 +90,11 @@ public Operation create(String statement) {
public TableResult execute(Executor executor) {
String yamlText = getPipelineConfigure(statement);
Configuration globalPipelineConfig = Configuration.fromMap(executor.getSetConfig());
// Parse pipeline definition file
YamlTextPipelineDefinitionParser pipelineDefinitionParser = new YamlTextPipelineDefinitionParser();
// Create composer
PipelineComposer composer = createComposer(executor);
try {
// Parse pipeline definition file
YamlPipelineDefinitionParser pipelineDefinitionParser = new YamlPipelineDefinitionParser();
// Create composer
PipelineComposer composer = createComposer(executor);
PipelineDef pipelineDef = pipelineDefinitionParser.parse(yamlText, globalPipelineConfig);
// Compose pipeline
composer.compose(pipelineDef);
Expand All @@ -110,8 +114,16 @@ public String getPipelineConfigure(String statement) {
return "";
}

public DinkyFlinkPipelineComposer createComposer(Executor executor) {

return DinkyFlinkPipelineComposer.of(executor);
public PipelineComposer createComposer(Executor executor) {
try {
Class<FlinkPipelineComposer> clazz = (Class<FlinkPipelineComposer>)
Class.forName("org.apache.flink.cdc.composer.flink.FlinkPipelineComposer");
Constructor<FlinkPipelineComposer> constructor =
clazz.getDeclaredConstructor(StreamExecutionEnvironment.class, boolean.class);
constructor.setAccessible(true);
return constructor.newInstance(executor.getStreamExecutionEnvironment(), false);
} catch (Exception e) {
throw new RuntimeException(e);
}
}
}
Loading

0 comments on commit c081c3d

Please sign in to comment.