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
2 people authored and donotcoffee committed Nov 12, 2024
1 parent 99ddc5b commit 1b6d647
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 1b6d647

Please sign in to comment.