This is an automated email from the ASF dual-hosted git repository.
leonard pushed a change to branch master
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
from ed09e8b23 [minor][cdc][docs] Improve the indentation of the example
yaml in README file
new f0c29f64f [cdc-composer] Introduce transform definition and parser
new c92903016 [cdc-common] Introduce partitionKeys into schema and add
related util classes
new bcad5d9d1 [cdc-runtime] Introduce TransformSchemaOperator and
TransformDataOperator to support transformation
new a35b8dd44 [build] Optimize pom to solve the CI error
new 1d05abff8 [minor][cdc-common] Improve the java doc of translators
The 5 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.../cli/parser/YamlPipelineDefinitionParser.java | 71 ++-
.../parser/YamlPipelineDefinitionParserTest.java | 41 +-
.../definitions/pipeline-definition-full.yaml | 3 +
.../flink/cdc/common/event/DataChangeEvent.java | 28 +
.../org/apache/flink/cdc/common/schema/Schema.java | 60 +-
.../flink/cdc/common/utils/DateTimeUtils.java | 122 ++++
.../apache/flink/cdc/common/utils/SchemaUtils.java | 43 +-
.../apache/flink/cdc/common/utils/StringUtils.java | 26 +
.../flink/cdc/common/utils/ThreadLocalCache.java | 87 +++
.../flink/cdc/common/utils/StringUtilsTest.java | 21 +-
.../cdc/composer/definition/TransformDef.java | 135 ++++-
.../cdc/composer/flink/FlinkPipelineComposer.java | 38 +-
.../flink/translator/DataSinkTranslator.java | 3 +-
.../flink/translator/DataSourceTranslator.java | 6 +-
.../flink/translator/PartitioningTranslator.java | 5 +-
.../composer/flink/translator/RouteTranslator.java | 3 +-
.../flink/translator/SchemaOperatorTranslator.java | 2 +-
.../flink/translator/TransformTranslator.java | 82 +++
.../flink/FlinkPipelineComposerITCase.java | 125 +++++
.../values/source/ValuesDataSourceHelper.java | 102 +++-
.../flink-sql-connector-db2-cdc/pom.xml | 4 +
.../flink-cdc-source-connectors/pom.xml | 10 +
flink-cdc-dist/pom.xml | 10 +
flink-cdc-runtime/pom.xml | 38 ++
.../runtime/functions/BuiltInScalarFunction.java | 238 ++++++++
.../functions/BuiltInTimestampFunction.java | 55 ++
.../cdc/runtime/functions/SystemFunctionUtils.java | 473 ++++++++++++++++
.../operators/transform/ProjectionColumn.java | 106 ++++
.../transform/ProjectionColumnProcessor.java | 153 +++++
.../transform/SchemaMetadataTransform.java | 85 +++
.../operators/transform/TableChangeInfo.java | 150 +++++
.../cdc/runtime/operators/transform/TableInfo.java | 90 +++
.../operators/transform/TransformDataOperator.java | 406 ++++++++++++++
.../transform/TransformExpressionCompiler.java | 71 +++
.../transform/TransformExpressionKey.java | 97 ++++
.../operators/transform/TransformFilter.java | 76 +++
.../transform/TransformFilterProcessor.java | 148 +++++
.../operators/transform/TransformProjection.java | 78 +++
.../transform/TransformProjectionProcessor.java | 186 ++++++
.../transform/TransformSchemaOperator.java | 297 ++++++++++
.../flink/cdc/runtime/parser/JaninoCompiler.java | 254 +++++++++
.../flink/cdc/runtime/parser/TransformParser.java | 378 +++++++++++++
.../TransformNumericExceptFirstOperandChecker.java | 91 +++
.../runtime/parser/metadata/TransformSchema.java | 29 +-
.../parser/metadata/TransformSchemaFactory.java | 48 ++
.../parser/metadata/TransformSqlOperatorTable.java | 249 ++++++++
.../parser/metadata/TransformSqlReturnTypes.java | 191 +++++++
.../runtime/parser/metadata/TransformTable.java | 60 ++
.../serializer/schema/SchemaSerializer.java | 5 +
.../cdc/runtime/typeutils/DataTypeConverter.java | 508 +++++++++++++++++
.../transform/TransformDataOperatorTest.java | 624 +++++++++++++++++++++
.../transform/TransformSchemaOperatorTest.java | 179 ++++++
.../cdc/runtime/parser/JaninoCompilerTest.java | 142 +++++
.../cdc/runtime/parser/TransformParserTest.java | 277 +++++++++
.../serializer/schema/SchemaSerializerTest.java | 1 +
pom.xml | 124 ++++
56 files changed, 6846 insertions(+), 88 deletions(-)
create mode 100644
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/utils/DateTimeUtils.java
create mode 100644
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/utils/ThreadLocalCache.java
copy
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-debezium/src/main/java/org/apache/flink/cdc/debezium/Validator.java
=>
flink-cdc-common/src/test/java/org/apache/flink/cdc/common/utils/StringUtilsTest.java
(68%)
create mode 100644
flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/flink/translator/TransformTranslator.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/BuiltInScalarFunction.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/BuiltInTimestampFunction.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/functions/SystemFunctionUtils.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumn.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/ProjectionColumnProcessor.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/SchemaMetadataTransform.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TableChangeInfo.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TableInfo.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformDataOperator.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformExpressionCompiler.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformExpressionKey.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilter.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformFilterProcessor.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformProjection.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformProjectionProcessor.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/transform/TransformSchemaOperator.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/JaninoCompiler.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/TransformParser.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformNumericExceptFirstOperandChecker.java
copy
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-postgres-cdc/src/main/java/org/apache/flink/cdc/connectors/postgres/source/handler/PostgresSchemaChangeEventHandler.java
=>
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSchema.java
(56%)
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSchemaFactory.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSqlOperatorTable.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformSqlReturnTypes.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/parser/metadata/TransformTable.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/typeutils/DataTypeConverter.java
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/TransformDataOperatorTest.java
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/transform/TransformSchemaOperatorTest.java
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/JaninoCompilerTest.java
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/parser/TransformParserTest.java