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;
+    }
+}

Reply via email to