This is an automated email from the ASF dual-hosted git repository.
kunni pushed a change to branch FLINK-38729-2
in repository https://gitbox.apache.org/repos/asf/flink-cdc.git
discard a814a86e3 Add support for Flink 2.x.
add 69dae39e7 [FLINK-37485][starrocks] Add support for TIME type (#4253)
add 31d68ac28 [chore][test] Fix flaky postgres pipeline test case (#4293)
add a1cfab9ef [FLINK-38833][paimon] Shuffle record to different subtasks
by table, partition and bucket. (#4298)
add c05af1594 [FLINK-38334][mysql] Fix MySQL CDC source stuck in
INITIAL_ASSIGNING (#4278)
add a19fdfe0f [FLINK-37586][udf] Add support for options in user-defined
functions and update related documentation (#4252)
add a3576d9e9 [FLINK-38726][fluss] Bump Fluss version to 0.9.0-incubating
add 885a0591a [FLINK-39204][pipeline-connector/fluss] Fluss yaml sink
support add column at last
add acfb3b1af Add support for Flink 2.x.
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (a814a86e3)
\
N -- N -- N refs/heads/FLINK-38729-2 (acfb3b1af)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
No new revisions were added by this update.
Summary of changes:
.../connectors/pipeline-connectors/starrocks.md | 18 +-
docs/content.zh/docs/core-concept/transform.md | 45 +++
.../connectors/pipeline-connectors/starrocks.md | 16 +-
docs/content/docs/core-concept/transform.md | 45 +++
.../cli/parser/YamlPipelineDefinitionParser.java | 13 +-
.../parser/YamlPipelineDefinitionParserTest.java | 47 +++
...l => pipeline-definition-with-udf-options.yaml} | 14 +-
.../flink/cdc/composer/definition/UdfDef.java | 30 +-
.../flink/translator/TransformTranslator.java | 3 +-
.../cdc/composer/flink/FlinkPipelineUdfITCase.java | 73 ++++
.../testsource/source/DistributedSource.java | 450 ++++++++++++++++++---
.../testsource/source/DistributedSourceReader.java | 334 ---------------
.../flink-cdc-pipeline-connector-fluss/pom.xml | 16 +-
.../fluss/factory/FlussDataSinkFactory.java | 4 +-
.../cdc/connectors/fluss/sink/FlussDataSink.java | 2 +-
.../fluss/sink/FlussEventSerializationSchema.java | 68 ++--
.../fluss/sink/FlussMetaDataApplier.java | 93 +++--
.../connectors/fluss/sink/row/CdcAsFlussArray.java | 175 ++++++++
.../connectors/fluss/sink/row/CdcAsFlussMap.java | 30 +-
.../fluss/sink/{ => row/row}/CdcAsFlussRow.java | 30 +-
.../cdc/connectors/fluss/sink/v2/FlussEvent.java | 2 +-
.../fluss/sink/v2/FlussEventSerializer.java | 2 +-
.../connectors/fluss/sink/v2/FlussRowWithOp.java | 4 +-
.../cdc/connectors/fluss/sink/v2/FlussSink.java | 2 +-
.../connectors/fluss/sink/v2/FlussSinkWriter.java | 24 +-
.../fluss/sink/v2/metrics/WrappedFlussCounter.java | 4 +-
.../fluss/sink/v2/metrics/WrappedFlussGauge.java | 4 +-
.../sink/v2/metrics/WrapperFlussHistogram.java | 8 +-
.../fluss/sink/v2/metrics/WrapperFlussMeter.java | 4 +-
.../v2/metrics/WrapperFlussMetricRegistry.java | 16 +-
.../connectors/fluss/utils/FlussConversions.java | 135 ++++---
.../cdc/connectors/fluss/FlussPipelineITCase.java | 259 ++++++++++--
.../sink/FlussEventSerializationSchemaTest.java | 85 +++-
.../fluss/sink/FlussMetadataApplierTest.java | 164 ++++++--
.../connectors/fluss/sink/v2/FlussSinkITCase.java | 55 ++-
.../fluss/utils/FlussConversionsTest.java | 94 ++---
.../connectors/paimon/sink/v2/PaimonEventSink.java | 1 +
.../sink/v2/bucket/BucketAssignOperator.java | 9 +-
.../sink/v2/bucket/BucketWrapperChangeEvent.java | 16 +-
.../v2/bucket/BucketWrapperEventSerializer.java | 5 +-
.../source/PostgresPipelineITCaseTest.java | 153 ++++---
.../connectors/starrocks/sink/StarRocksUtils.java | 57 ++-
.../sink/EventRecordSerializationSchemaTest.java | 224 ++++++++++
.../sink/StarRocksMetadataApplierITCase.java | 11 +-
.../sink/StarRocksMetadataApplierTest.java | 147 +++++++
.../assigners/MySqlSnapshotSplitAssigner.java | 12 +-
.../source/TableExclusionDuringSnapshotIT.java | 283 +++++++++++++
.../resources/ddl/table_exclusion_snapshot.sql | 16 +-
flink-cdc-dist/pom.xml | 13 +
.../flink-cdc-pipeline-e2e-tests/pom.xml | 17 +-
.../flink/cdc/pipeline/tests/FlussE2eITCase.java | 8 +-
.../src/test/resources/rules/unexpected.yaml | 2 +-
.../examples/java/ConfigurableFunctionClass.java | 49 +++
.../examples/scala/ConfigurableFunctionClass.scala | 44 ++
.../operators/transform/PreTransformOperator.java | 23 ++
55 files changed, 2613 insertions(+), 845 deletions(-)
copy
flink-cdc-cli/src/test/resources/definitions/{pipeline-definition-with-udf.yaml
=> pipeline-definition-with-udf-options.yaml} (78%)
delete mode 100644
flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/testsource/source/DistributedSourceReader.java
create mode 100644
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/sink/row/CdcAsFlussArray.java
copy
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/serializer/StringSerializerTest.java
=>
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/sink/row/CdcAsFlussMap.java
(57%)
rename
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/src/main/java/org/apache/flink/cdc/connectors/fluss/sink/{
=> row/row}/CdcAsFlussRow.java (87%)
create mode 100644
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/java/org/apache/flink/cdc/connectors/mysql/source/TableExclusionDuringSnapshotIT.java
copy
flink-cdc-e2e-tests/flink-cdc-source-e2e-tests/src/test/resources/docker/db2/startup-agent.sql
=>
flink-cdc-connect/flink-cdc-source-connectors/flink-connector-mysql-cdc/src/test/resources/ddl/table_exclusion_snapshot.sql
(82%)
create mode 100644
flink-cdc-pipeline-udf-examples/src/main/java/org/apache/flink/cdc/udf/examples/java/ConfigurableFunctionClass.java
create mode 100644
flink-cdc-pipeline-udf-examples/src/main/scala/org/apache/flink/cdc/udf/examples/scala/ConfigurableFunctionClass.scala