Skip to content

Commit

Permalink
upgrade cdc to 3.2.0
Browse files Browse the repository at this point in the history
  • Loading branch information
gaoyan1998 committed Oct 18, 2024
1 parent 202a13e commit 7dc4504
Show file tree
Hide file tree
Showing 11 changed files with 50 additions and 451 deletions.

This file was deleted.

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -19,54 +19,56 @@

package org.dinky.trans.pipeline;

import org.dinky.executor.Executor;
import org.dinky.trans.AbstractOperation;
import org.dinky.trans.Operation;

import java.lang.reflect.Constructor;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
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.util.regex.Matcher;
import java.util.regex.Pattern;

import org.dinky.executor.Executor;
import org.dinky.trans.AbstractOperation;
import org.dinky.trans.Operation;
import org.jetbrains.annotations.Nullable;

/**
* FlinkCDCPipelineOperation
*
* <p>
* ################################################################################
* # Description: Sync MySQL all tables to Doris
* ################################################################################
* source:
* type: mysql
* hostname: localhost
* port: 3306
* username: root
* password: 123456
* tables: app_db.\.*
* server-id: 5400-5404
* server-time-zone: UTC
*
* type: mysql
* hostname: localhost
* port: 3306
* username: root
* password: 123456
* tables: app_db.\.*
* server-id: 5400-5404
* server-time-zone: UTC
* <p>
* sink:
* type: doris
* fenodes: 127.0.0.1:8030
* username: root
* password: ""
* table.create.properties.light_schema_change: true
* table.create.properties.replication_num: 1
*
* type: doris
* fenodes: 127.0.0.1:8030
* username: root
* password: ""
* table.create.properties.light_schema_change: true
* table.create.properties.replication_num: 1
* <p>
* pipeline:
* name: Sync MySQL Database to Doris
* parallelism: 2
* name: Sync MySQL Database to Doris
* parallelism: 2
*/
public class FlinkCDCPipelineOperation extends AbstractOperation implements Operation {

private static final String KEY_WORD = "EXECUTE PIPELINE";

public FlinkCDCPipelineOperation() {}
public FlinkCDCPipelineOperation() {
}

public FlinkCDCPipelineOperation(String statement) {
super(statement);
Expand All @@ -86,11 +88,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 +112,14 @@ 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 7dc4504

Please sign in to comment.