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 1d05abff873a8a807aa67465332ff2100f0805cd
Author: Leonard Xu <@users.noreply.github.com>
AuthorDate: Tue Apr 2 19:34:01 2024 +0800

    [minor][cdc-common] Improve the java doc of translators
    
    This closes #2937.
---
 .../flink/cdc/composer/flink/FlinkPipelineComposer.java | 17 +++++++----------
 .../composer/flink/translator/DataSinkTranslator.java   |  3 ++-
 .../composer/flink/translator/DataSourceTranslator.java |  6 ++----
 .../flink/translator/PartitioningTranslator.java        |  5 ++++-
 .../cdc/composer/flink/translator/RouteTranslator.java  |  3 ++-
 .../flink/translator/SchemaOperatorTranslator.java      |  2 +-
 .../composer/flink/translator/TransformTranslator.java  |  5 ++++-
 7 files changed, 22 insertions(+), 19 deletions(-)

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 fdbd37d00..ceb285f85 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
@@ -96,27 +96,24 @@ public class FlinkPipelineComposer implements 
PipelineComposer {
         int parallelism = 
pipelineDef.getConfig().get(PipelineOptions.PIPELINE_PARALLELISM);
         env.getConfig().setParallelism(parallelism);
 
-        // Source
+        // Build Source Operator
         DataSourceTranslator sourceTranslator = new DataSourceTranslator();
         DataStream<Event> stream =
                 sourceTranslator.translate(pipelineDef.getSource(), env, 
pipelineDef.getConfig());
 
-        // Transform Schema
+        // Build TransformSchemaOperator for processing Schema Event
         TransformTranslator transformTranslator = new TransformTranslator();
         stream = transformTranslator.translateSchema(stream, 
pipelineDef.getTransforms());
-
-        // Schema operator
         SchemaOperatorTranslator schemaOperatorTranslator =
                 new SchemaOperatorTranslator(
                         pipelineDef
                                 .getConfig()
                                 
.get(PipelineOptions.PIPELINE_SCHEMA_CHANGE_BEHAVIOR),
                         
pipelineDef.getConfig().get(PipelineOptions.PIPELINE_SCHEMA_OPERATOR_UID));
-
         OperatorIDGenerator schemaOperatorIDGenerator =
                 new 
OperatorIDGenerator(schemaOperatorTranslator.getSchemaOperatorUid());
 
-        // Transform Data
+        // Build TransformDataOperator for processing Data Event
         stream =
                 transformTranslator.translateData(
                         stream,
@@ -124,24 +121,24 @@ public class FlinkPipelineComposer implements 
PipelineComposer {
                         schemaOperatorIDGenerator.generate(),
                         
pipelineDef.getConfig().get(PipelineOptions.PIPELINE_LOCAL_TIME_ZONE));
 
-        // Route
+        // Build Router used to route Event
         RouteTranslator routeTranslator = new RouteTranslator();
         stream = routeTranslator.translate(stream, pipelineDef.getRoute());
 
-        // Create sink in advance as schema operator requires MetadataApplier
+        // Build DataSink in advance as schema operator requires 
MetadataApplier
         DataSink dataSink = createDataSink(pipelineDef.getSink(), 
pipelineDef.getConfig());
 
         stream =
                 schemaOperatorTranslator.translate(
                         stream, parallelism, dataSink.getMetadataApplier());
 
-        // Add partitioner
+        // Build Partitioner used to shuffle Event
         PartitioningTranslator partitioningTranslator = new 
PartitioningTranslator();
         stream =
                 partitioningTranslator.translate(
                         stream, parallelism, parallelism, 
schemaOperatorIDGenerator.generate());
 
-        // Sink
+        // Build Sink Operator
         DataSinkTranslator sinkTranslator = new DataSinkTranslator();
         sinkTranslator.translate(
                 pipelineDef.getSink(), stream, dataSink, 
schemaOperatorIDGenerator.generate());
diff --git 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
index 8be9a1c5d..8bf3ef88d 100644
--- 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
+++ 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSinkTranslator.java
@@ -36,9 +36,10 @@ import 
org.apache.flink.streaming.api.connector.sink2.WithPreWriteTopology;
 import org.apache.flink.streaming.api.datastream.DataStream;
 import 
org.apache.flink.streaming.runtime.operators.sink.CommitterOperatorFactory;
 
-/** Translator for building sink into the DataStream. */
+/** Translator used to build {@link DataSink} for given {@link DataStream}. */
 @Internal
 public class DataSinkTranslator {
+
     private static final String SINK_WRITER_PREFIX = "Sink Writer: ";
     private static final String SINK_COMMITTER_PREFIX = "Sink Committer: ";
 
diff --git 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSourceTranslator.java
 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSourceTranslator.java
index f3c7b8372..ee9b17d7b 100644
--- 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSourceTranslator.java
+++ 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/DataSourceTranslator.java
@@ -32,13 +32,11 @@ import org.apache.flink.cdc.composer.definition.SourceDef;
 import org.apache.flink.cdc.composer.flink.FlinkEnvironmentUtils;
 import org.apache.flink.cdc.composer.utils.FactoryDiscoveryUtils;
 import org.apache.flink.cdc.runtime.typeutils.EventTypeInfo;
+import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.datastream.DataStreamSource;
 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
 
-/**
- * Translator for building source and generate a {@link
- * org.apache.flink.streaming.api.datastream.DataStream}.
- */
+/** Translator used to build {@link DataSource} which will generate a {@link 
DataStream}. */
 @Internal
 public class DataSourceTranslator {
 
diff --git 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/PartitioningTranslator.java
 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/PartitioningTranslator.java
index 3723c8e71..4f076685d 100644
--- 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/PartitioningTranslator.java
+++ 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/PartitioningTranslator.java
@@ -28,7 +28,10 @@ import 
org.apache.flink.cdc.runtime.typeutils.PartitioningEventTypeInfo;
 import org.apache.flink.runtime.jobgraph.OperatorID;
 import org.apache.flink.streaming.api.datastream.DataStream;
 
-/** Translator for building partitioning related transformations. */
+/**
+ * Translator used to build {@link PrePartitionOperator}, {@link 
EventPartitioner} and {@link
+ * PostPartitionProcessor} which are responsible for events partition.
+ */
 @Internal
 public class PartitioningTranslator {
 
diff --git 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/RouteTranslator.java
 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/RouteTranslator.java
index 28dde2d5c..0ad0c0dd3 100644
--- 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/RouteTranslator.java
+++ 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/RouteTranslator.java
@@ -26,8 +26,9 @@ import org.apache.flink.streaming.api.datastream.DataStream;
 
 import java.util.List;
 
-/** Translator for router. */
+/** Translator used to build {@link RouteFunction}. */
 public class RouteTranslator {
+
     public DataStream<Event> translate(DataStream<Event> input, List<RouteDef> 
routes) {
         if (routes.isEmpty()) {
             return input;
diff --git 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/SchemaOperatorTranslator.java
 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/SchemaOperatorTranslator.java
index 405f7fc10..bb434581b 100644
--- 
a/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/SchemaOperatorTranslator.java
+++ 
b/flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/SchemaOperatorTranslator.java
@@ -28,7 +28,7 @@ import org.apache.flink.cdc.runtime.typeutils.EventTypeInfo;
 import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
 
-/** Translator for building {@link SchemaOperator} into DataStream. */
+/** Translator used to build {@link SchemaOperator} for schema event process. 
*/
 @Internal
 public class SchemaOperatorTranslator {
     private final SchemaChangeBehavior schemaChangeBehavior;
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
index 74b3d33ef..53400f628 100644
--- 
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
@@ -27,7 +27,10 @@ import org.apache.flink.streaming.api.datastream.DataStream;
 
 import java.util.List;
 
-/** Translator for transform schema. */
+/**
+ * Translator used to build {@link TransformSchemaOperator} and {@link 
TransformDataOperator} for
+ * event transform.
+ */
 public class TransformTranslator {
 
     public DataStream<Event> translateSchema(

Reply via email to