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 f0c29f64f226265539787448352b8f281a0fbd71 Author: wenmo <[email protected]> AuthorDate: Wed Mar 20 19:45:40 2024 +0800 [cdc-composer] Introduce transform definition and parser --- .../cli/parser/YamlPipelineDefinitionParser.java | 71 ++++++++++- .../parser/YamlPipelineDefinitionParserTest.java | 41 ++++++- .../definitions/pipeline-definition-full.yaml | 3 + .../cdc/composer/definition/TransformDef.java | 135 ++++++++++++++++++++- .../cdc/composer/flink/FlinkPipelineComposer.java | 31 +++-- .../flink/translator/TransformTranslator.java | 48 ++++++++ 6 files changed, 313 insertions(+), 16 deletions(-) diff --git a/flink-cdc-cli/src/main/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParser.java b/flink-cdc-cli/src/main/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParser.java index 67db2b821..23b2c63ff 100644 --- a/flink-cdc-cli/src/main/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParser.java +++ b/flink-cdc-cli/src/main/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParser.java @@ -18,10 +18,12 @@ package org.apache.flink.cdc.cli.parser; import org.apache.flink.cdc.common.configuration.Configuration; +import org.apache.flink.cdc.common.utils.StringUtils; import org.apache.flink.cdc.composer.definition.PipelineDef; import org.apache.flink.cdc.composer.definition.RouteDef; import org.apache.flink.cdc.composer.definition.SinkDef; import org.apache.flink.cdc.composer.definition.SourceDef; +import org.apache.flink.cdc.composer.definition.TransformDef; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.type.TypeReference; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JsonNode; @@ -43,6 +45,7 @@ public class YamlPipelineDefinitionParser implements PipelineDefinitionParser { private static final String SOURCE_KEY = "source"; private static final String SINK_KEY = "sink"; private static final String ROUTE_KEY = "route"; + private static final String TRANSFORM_KEY = "transform"; private static final String PIPELINE_KEY = "pipeline"; // Source / sink keys @@ -54,6 +57,18 @@ public class YamlPipelineDefinitionParser implements PipelineDefinitionParser { private static final String ROUTE_SINK_TABLE_KEY = "sink-table"; private static final String ROUTE_DESCRIPTION_KEY = "description"; + // Transform keys + private static final String TRANSFORM_SOURCE_TABLE_KEY = "source-table"; + private static final String TRANSFORM_PROJECTION_KEY = "projection"; + private static final String TRANSFORM_FILTER_KEY = "filter"; + private static final String TRANSFORM_DESCRIPTION_KEY = "description"; + + public static final String TRANSFORM_PRIMARY_KEY_KEY = "primary-keys"; + + public static final String TRANSFORM_PARTITION_KEY_KEY = "partition-keys"; + + public static final String TRANSFORM_TABLE_OPTION_KEY = "table-options"; + private final ObjectMapper mapper = new ObjectMapper(new YAMLFactory()); /** Parse the specified pipeline definition file. */ @@ -78,6 +93,14 @@ public class YamlPipelineDefinitionParser implements PipelineDefinitionParser { "Missing required field \"%s\" in pipeline definition", SINK_KEY)); + // Transforms are optional + List<TransformDef> transformDefs = new ArrayList<>(); + Optional.ofNullable(root.get(TRANSFORM_KEY)) + .ifPresent( + node -> + node.forEach( + transform -> transformDefs.add(toTransformDef(transform)))); + // Routes are optional List<RouteDef> routeDefs = new ArrayList<>(); Optional.ofNullable(root.get(ROUTE_KEY)) @@ -91,7 +114,7 @@ public class YamlPipelineDefinitionParser implements PipelineDefinitionParser { pipelineConfig.addAll(globalPipelineConfig); pipelineConfig.addAll(userPipelineConfig); - return new PipelineDef(sourceDef, sinkDef, routeDefs, null, pipelineConfig); + return new PipelineDef(sourceDef, sinkDef, routeDefs, transformDefs, pipelineConfig); } private SourceDef toSourceDef(JsonNode sourceNode) { @@ -148,6 +171,52 @@ public class YamlPipelineDefinitionParser implements PipelineDefinitionParser { return new RouteDef(sourceTable, sinkTable, description); } + private TransformDef toTransformDef(JsonNode transformNode) { + String sourceTable = + checkNotNull( + transformNode.get(TRANSFORM_SOURCE_TABLE_KEY), + "Missing required field \"%s\" in transform configuration", + TRANSFORM_SOURCE_TABLE_KEY) + .asText(); + String projection = + Optional.ofNullable(transformNode.get(TRANSFORM_PROJECTION_KEY)) + .map(JsonNode::asText) + .orElse(null); + // When the star is in the first place, a backslash needs to be added for escape. + if (!StringUtils.isNullOrWhitespaceOnly(projection) && projection.contains("\\*")) { + projection = projection.replace("\\*", "*"); + } + String filter = + Optional.ofNullable(transformNode.get(TRANSFORM_FILTER_KEY)) + .map(JsonNode::asText) + .orElse(null); + String primaryKeys = + Optional.ofNullable(transformNode.get(TRANSFORM_PRIMARY_KEY_KEY)) + .map(JsonNode::asText) + .orElse(null); + String partitionKeys = + Optional.ofNullable(transformNode.get(TRANSFORM_PARTITION_KEY_KEY)) + .map(JsonNode::asText) + .orElse(null); + String tableOptions = + Optional.ofNullable(transformNode.get(TRANSFORM_TABLE_OPTION_KEY)) + .map(JsonNode::asText) + .orElse(null); + String description = + Optional.ofNullable(transformNode.get(TRANSFORM_DESCRIPTION_KEY)) + .map(JsonNode::asText) + .orElse(null); + + return new TransformDef( + sourceTable, + projection, + filter, + primaryKeys, + partitionKeys, + tableOptions, + description); + } + private Configuration toPipelineConfig(JsonNode pipelineConfigNode) { if (pipelineConfigNode == null || pipelineConfigNode.isNull()) { return new Configuration(); diff --git a/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java b/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java index 6a1c9c688..e29ea332a 100644 --- a/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java +++ b/flink-cdc-cli/src/test/java/org/apache/flink/cdc/cli/parser/YamlPipelineDefinitionParserTest.java @@ -22,6 +22,7 @@ import org.apache.flink.cdc.composer.definition.PipelineDef; import org.apache.flink.cdc.composer.definition.RouteDef; import org.apache.flink.cdc.composer.definition.SinkDef; import org.apache.flink.cdc.composer.definition.SourceDef; +import org.apache.flink.cdc.composer.definition.TransformDef; import org.apache.flink.shaded.guava31.com.google.common.collect.ImmutableMap; import org.apache.flink.shaded.guava31.com.google.common.io.Resources; @@ -179,7 +180,23 @@ class YamlPipelineDefinitionParserTest { "mydb.default.web_order", "odsdb.default.ods_web_order", "sync table to with given prefix ods_")), - null, + Arrays.asList( + new TransformDef( + "mydb.app_order_.*", + "id, order_id, TO_UPPER(product_name)", + "id > 10 AND order_id > 100", + "id", + "product_name", + "comment=app order", + "project fields from source table"), + new TransformDef( + "mydb.web_order_.*", + "CONCAT(id, order_id) as uniq_id, *", + "uniq_id > 10", + null, + null, + null, + "add new uniq_id for each row")), Configuration.fromMap( ImmutableMap.<String, String>builder() .put("name", "source-database-sync-pipe") @@ -223,7 +240,23 @@ class YamlPipelineDefinitionParserTest { "mydb.default.web_order", "odsdb.default.ods_web_order", "sync table to with given prefix ods_")), - null, + Arrays.asList( + new TransformDef( + "mydb.app_order_.*", + "id, order_id, TO_UPPER(product_name)", + "id > 10 AND order_id > 100", + "id", + "product_name", + "comment=app order", + "project fields from source table"), + new TransformDef( + "mydb.web_order_.*", + "CONCAT(id, order_id) as uniq_id, *", + "uniq_id > 10", + null, + null, + null, + "add new uniq_id for each row")), Configuration.fromMap( ImmutableMap.<String, String>builder() .put("name", "source-database-sync-pipe") @@ -257,7 +290,7 @@ class YamlPipelineDefinitionParserTest { Collections.singletonList( new RouteDef( "mydb.default.app_order_.*", "odsdb.default.app_order", null)), - null, + Collections.emptyList(), Configuration.fromMap( ImmutableMap.<String, String>builder() .put("parallelism", "4") @@ -268,6 +301,6 @@ class YamlPipelineDefinitionParserTest { new SourceDef("mysql", null, new Configuration()), new SinkDef("kafka", null, new Configuration()), Collections.emptyList(), - null, + Collections.emptyList(), new Configuration()); } diff --git a/flink-cdc-cli/src/test/resources/definitions/pipeline-definition-full.yaml b/flink-cdc-cli/src/test/resources/definitions/pipeline-definition-full.yaml index f6a4baa3f..5f0f6f8f9 100644 --- a/flink-cdc-cli/src/test/resources/definitions/pipeline-definition-full.yaml +++ b/flink-cdc-cli/src/test/resources/definitions/pipeline-definition-full.yaml @@ -43,6 +43,9 @@ transform: - source-table: mydb.app_order_.* projection: id, order_id, TO_UPPER(product_name) filter: id > 10 AND order_id > 100 + primary-keys: id + partition-keys: product_name + table-options: comment=app order description: project fields from source table - source-table: mydb.web_order_.* projection: CONCAT(id, order_id) as uniq_id, * diff --git a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/definition/TransformDef.java b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/definition/TransformDef.java index 6b31255d7..62491917d 100644 --- a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/definition/TransformDef.java +++ b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/definition/TransformDef.java @@ -17,9 +17,138 @@ package org.apache.flink.cdc.composer.definition; +import org.apache.flink.cdc.common.utils.StringUtils; + +import java.util.Objects; +import java.util.Optional; + /** - * Definition of transformation. + * Definition of a transformation. * - * <p>Transformation will be implemented later, therefore we left the class blank. + * <p>A transformation definition contains: + * + * <ul> + * <li>sourceTable: a regex pattern for matching input table IDs. Required for the definition. + * <li>projection: a string for projecting the row of matched table as output. Optional for the + * definition. + * <li>filter: a string for filtering the row of matched table as output. Optional for the + * definition. + * <li>primaryKeys: a string for primary key columns for matching input table IDs, seperated by + * `,`. Optional for the definition. + * <li>partitionKeys: a string for partition key columns for matching input table IDs, seperated + * by `,`. Optional for the definition. + * <li>tableOptions: a string for table options for matching input table IDs, options are + * seperated by `,`, key and value are seperated by `=`. Optional for the definition. + * <li>description: description for the transformation. Optional for the definition. + * </ul> */ -public class TransformDef {} +public class TransformDef { + private final String sourceTable; + private final String projection; + private final String filter; + private final String description; + private final String primaryKeys; + private final String partitionKeys; + private final String tableOptions; + + public TransformDef( + String sourceTable, + String projection, + String filter, + String primaryKeys, + String partitionKeys, + String tableOptions, + String description) { + this.sourceTable = sourceTable; + this.projection = projection; + this.filter = filter; + this.primaryKeys = primaryKeys; + this.partitionKeys = partitionKeys; + this.tableOptions = tableOptions; + this.description = description; + } + + public String getSourceTable() { + return sourceTable; + } + + public Optional<String> getProjection() { + return Optional.ofNullable(projection); + } + + public boolean isValidProjection() { + return !StringUtils.isNullOrWhitespaceOnly(projection); + } + + public Optional<String> getFilter() { + return Optional.ofNullable(filter); + } + + public boolean isValidFilter() { + return !StringUtils.isNullOrWhitespaceOnly(filter); + } + + public Optional<String> getDescription() { + return Optional.ofNullable(description); + } + + public String getPrimaryKeys() { + return primaryKeys; + } + + public String getPartitionKeys() { + return partitionKeys; + } + + public String getTableOptions() { + return tableOptions; + } + + @Override + public String toString() { + return "TransformDef{" + + "sourceTable='" + + sourceTable + + '\'' + + ", projection='" + + projection + + '\'' + + ", filter='" + + filter + + '\'' + + ", description='" + + description + + '\'' + + '}'; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + TransformDef that = (TransformDef) o; + return Objects.equals(sourceTable, that.sourceTable) + && Objects.equals(projection, that.projection) + && Objects.equals(filter, that.filter) + && Objects.equals(description, that.description) + && Objects.equals(primaryKeys, that.primaryKeys) + && Objects.equals(partitionKeys, that.partitionKeys) + && Objects.equals(tableOptions, that.tableOptions); + } + + @Override + public int hashCode() { + return Objects.hash( + sourceTable, + projection, + filter, + description, + primaryKeys, + partitionKeys, + tableOptions); + } +} diff --git a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java index f8a322839..fdbd37d00 100644 --- a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java +++ b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/FlinkPipelineComposer.java @@ -34,6 +34,7 @@ import org.apache.flink.cdc.composer.flink.translator.DataSourceTranslator; import org.apache.flink.cdc.composer.flink.translator.PartitioningTranslator; import org.apache.flink.cdc.composer.flink.translator.RouteTranslator; import org.apache.flink.cdc.composer.flink.translator.SchemaOperatorTranslator; +import org.apache.flink.cdc.composer.flink.translator.TransformTranslator; import org.apache.flink.cdc.composer.utils.FactoryDiscoveryUtils; import org.apache.flink.cdc.runtime.serializer.event.EventSerializer; import org.apache.flink.configuration.DeploymentOptions; @@ -100,12 +101,9 @@ public class FlinkPipelineComposer implements PipelineComposer { DataStream<Event> stream = sourceTranslator.translate(pipelineDef.getSource(), env, pipelineDef.getConfig()); - // Route - RouteTranslator routeTranslator = new RouteTranslator(); - stream = routeTranslator.translate(stream, pipelineDef.getRoute()); - - // Create sink in advance as schema operator requires MetadataApplier - DataSink dataSink = createDataSink(pipelineDef.getSink(), pipelineDef.getConfig()); + // Transform Schema + TransformTranslator transformTranslator = new TransformTranslator(); + stream = transformTranslator.translateSchema(stream, pipelineDef.getTransforms()); // Schema operator SchemaOperatorTranslator schemaOperatorTranslator = @@ -114,11 +112,28 @@ public class FlinkPipelineComposer implements PipelineComposer { .getConfig() .get(PipelineOptions.PIPELINE_SCHEMA_CHANGE_BEHAVIOR), pipelineDef.getConfig().get(PipelineOptions.PIPELINE_SCHEMA_OPERATOR_UID)); + + OperatorIDGenerator schemaOperatorIDGenerator = + new OperatorIDGenerator(schemaOperatorTranslator.getSchemaOperatorUid()); + + // Transform Data + stream = + transformTranslator.translateData( + stream, + pipelineDef.getTransforms(), + schemaOperatorIDGenerator.generate(), + pipelineDef.getConfig().get(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE)); + + // Route + RouteTranslator routeTranslator = new RouteTranslator(); + stream = routeTranslator.translate(stream, pipelineDef.getRoute()); + + // Create sink in advance as schema operator requires MetadataApplier + DataSink dataSink = createDataSink(pipelineDef.getSink(), pipelineDef.getConfig()); + stream = schemaOperatorTranslator.translate( stream, parallelism, dataSink.getMetadataApplier()); - OperatorIDGenerator schemaOperatorIDGenerator = - new OperatorIDGenerator(schemaOperatorTranslator.getSchemaOperatorUid()); // Add partitioner PartitioningTranslator partitioningTranslator = new PartitioningTranslator(); diff --git a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java new file mode 100644 index 000000000..866154ae5 --- /dev/null +++ b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java @@ -0,0 +1,48 @@ +/* + * 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.composer.flink.translator; + +import org.apache.flink.cdc.common.event.Event; +import org.apache.flink.cdc.composer.definition.TransformDef; +import org.apache.flink.runtime.jobgraph.OperatorID; +import org.apache.flink.streaming.api.datastream.DataStream; + +import java.util.List; + +/** Translator for transform schema. */ +public class TransformTranslator { + + public DataStream<Event> translateSchema( + DataStream<Event> input, List<TransformDef> transforms) { + if (transforms.isEmpty()) { + return input; + } + return input; + } + + public DataStream<Event> translateData( + DataStream<Event> input, + List<TransformDef> transforms, + OperatorID schemaOperatorID, + String timezone) { + if (transforms.isEmpty()) { + return input; + } + return input; + } +}
