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(