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 dfa6de919 [hotfix][docs] Fix ToC to include H1
new 88d1d6fbd [tests][build] Update migration test matrix to 3.2.0 and
later
new a2e433ebd [cdc-common] Extract column / schema type merging utility
methods to `SchemaMergingUtils`
new 2e03f9ad4 [FLINK-36763][cdc-runtime] Introduce distributed schema
evolution topology for sources with parallelized metadata
new 842446723 [FLINK-36690][cdc-runtime] Fix schema operator hanging under
extreme parallelized pressure This closes #3680
The 4 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:
.github/workflows/flink_cdc_base.yml | 7 +-
.github/workflows/flink_cdc_ci.yml | 6 +-
.github/workflows/flink_cdc_ci_nightly.yml | 6 +-
.../flink/cdc/common/event/AddColumnEvent.java | 5 +
.../cdc/common/event/AlterColumnTypeEvent.java | 5 +
.../flink/cdc/common/event/CreateTableEvent.java | 5 +
.../flink/cdc/common/event/DataChangeEvent.java | 10 +
.../flink/cdc/common/event/DropColumnEvent.java | 5 +
.../flink/cdc/common/event/DropTableEvent.java | 5 +
.../apache/flink/cdc/common/event/FlushEvent.java | 21 +-
.../flink/cdc/common/event/RenameColumnEvent.java | 5 +
.../flink/cdc/common/event/SchemaChangeEvent.java | 3 +
.../org/apache/flink/cdc/common/event/TableId.java | 6 +-
.../flink/cdc/common/event/TruncateTableEvent.java | 5 +
.../apache/flink/cdc/common/route/RouteRule.java | 4 +
.../org/apache/flink/cdc/common/schema/Schema.java | 11 +-
.../apache/flink/cdc/common/source/DataSource.java | 15 +
.../flink/cdc/common/utils/SchemaMergingUtils.java | 859 ++++++++++++++
.../apache/flink/cdc/common/utils/SchemaUtils.java | 498 ++++----
.../cdc/common/utils/SchemaMergingUtilsTest.java | 1184 ++++++++++++++++++++
.../cdc/composer/flink/FlinkPipelineComposer.java | 132 ++-
.../flink/translator/DataSourceTranslator.java | 38 +-
.../flink/translator/PartitioningTranslator.java | 27 +-
.../flink/translator/SchemaOperatorTranslator.java | 61 +-
.../flink/translator/TransformTranslator.java | 16 +-
.../flink/FlinkParallelizedPipelineITCase.java | 1011 +++++++++++++++++
.../flink/FlinkPipelineComposerITCase.java | 43 +-
.../flink/FlinkPipelineComposerLenientITCase.java | 8 +-
.../factory/DistributedDataSourceFactory.java} | 54 +-
.../testsource/source/DistributedDataSource.java | 56 +
.../source/DistributedSourceFunction.java | 275 +++++
.../source/DistributedSourceOptions.java} | 27 +-
.../org.apache.flink.cdc.common.factories.Factory | 1 +
.../resources/ref-output/distributed-ignore.txt | 172 +++
.../src/test/resources/ref-output/distributed.txt | 300 +++++
.../src/test/resources/ref-output/regular.txt | 300 +++++
.../doris/sink/DorisMetadataApplierITCase.java | 2 +-
.../pom.xml | 74 +-
.../flink-cdc-pipeline-connector-mysql/pom.xml | 14 +
.../connectors/mysql/source/MySqlDataSource.java | 7 +
.../source/MySqlParallelizedPipelineITCase.java | 207 ++++
.../mysql/testutils/MySqSourceTestUtils.java | 17 +
.../sink/v2/bucket/BucketAssignOperator.java | 2 +-
.../v2/bucket/BucketWrapperEventSerializer.java | 5 +-
.../sink/v2/bucket/BucketWrapperFlushEvent.java | 11 +-
.../sink/StarRocksMetadataApplierITCase.java | 2 +-
.../values/sink/ValuesDataSinkHelper.java | 14 +-
.../flink/cdc/pipeline/tests/MysqlE2eITCase.java | 14 +-
.../flink/cdc/pipeline/tests/RouteE2eITCase.java | 125 ++-
.../cdc/pipeline/tests/SchemaEvolveE2eITCase.java | 35 +-
.../tests/SchemaEvolvingTransformE2eITCase.java | 31 +-
.../cdc/pipeline/tests/TransformE2eITCase.java | 10 +-
.../tests/utils/PipelineTestEnvironment.java | 13 +-
.../flink-cdc-migration-testcases/pom.xml | 16 +-
.../cdc/migration/tests/MigrationTestBase.java | 45 +-
.../tests/SchemaManagerMigrationTest.java | 12 +-
.../tests/SchemaRegistryMigrationTest.java | 14 +-
.../tests/TableChangeInfoMigrationTest.java | 20 +-
.../flink-cdc-release-3.0.0/pom.xml | 86 --
.../cdc/migration/tests/MigrationMockBase.java | 27 -
.../tests/SchemaManagerMigrationMock.java | 65 --
.../tests/SchemaRegistryMigrationMock.java | 95 --
.../flink-cdc-release-3.0.1/pom.xml | 86 --
.../cdc/migration/tests/MigrationMockBase.java | 27 -
.../tests/SchemaManagerMigrationMock.java | 65 --
.../tests/SchemaRegistryMigrationMock.java | 95 --
.../tests/SchemaRegistryMigrationMock.java | 116 --
.../tests/TableChangeInfoMigrationMock.java | 60 -
.../tests/SchemaRegistryMigrationMock.java | 116 --
.../tests/TableChangeInfoMigrationMock.java | 60 -
.../pom.xml | 13 +-
.../cdc/migration/tests/MigrationMockBase.java | 0
.../tests/SchemaManagerMigrationMock.java | 53 +-
.../tests/SchemaRegistryMigrationMock.java | 181 +++
.../tests/TableChangeInfoMigrationMock.java | 0
.../pom.xml | 13 +-
.../cdc/migration/tests/MigrationMockBase.java | 0
.../tests/SchemaManagerMigrationMock.java | 53 +-
.../tests/SchemaRegistryMigrationMock.java | 181 +++
.../tests/TableChangeInfoMigrationMock.java | 0
.../tests/SchemaManagerMigrationMock.java | 58 +-
.../tests/SchemaRegistryMigrationMock.java | 227 ++--
flink-cdc-migration-tests/pom.xml | 6 +-
.../runtime/operators/schema/SchemaOperator.java | 725 ------------
.../CoordinationResponseUtils.java | 2 +-
.../common/CoordinatorExecutorThreadFactory.java | 60 +
.../operators/schema/common/SchemaDerivator.java | 345 ++++++
.../{coordinator => common}/SchemaManager.java | 164 +--
.../operators/schema/common/SchemaRegistry.java | 397 +++++++
.../operators/schema/common/TableIdRouter.java | 98 ++
.../{ => common}/event/FlushSuccessEvent.java | 39 +-
.../event/GetEvolvedSchemaRequest.java | 4 +-
.../event/GetEvolvedSchemaResponse.java | 4 +-
.../event/GetOriginalSchemaRequest.java | 4 +-
.../event/GetOriginalSchemaResponse.java | 4 +-
.../event/SinkWriterRegisterEvent.java | 4 +-
.../metrics/SchemaOperatorMetrics.java | 4 +-
.../schema/coordinator/SchemaDerivation.java | 361 ------
.../schema/coordinator/SchemaRegistry.java | 444 --------
.../coordinator/SchemaRegistryRequestHandler.java | 517 ---------
.../schema/distributed/SchemaCoordinator.java | 452 ++++++++
.../SchemaCoordinatorProvider.java} | 60 +-
.../schema/distributed/SchemaOperator.java | 231 ++++
.../{ => distributed}/SchemaOperatorFactory.java | 23 +-
.../distributed/event/SchemaChangeRequest.java | 68 ++
.../distributed/event/SchemaChangeResponse.java | 49 +-
.../event/SchemaChangeProcessingResponse.java | 32 -
.../schema/event/SchemaChangeResultRequest.java | 31 -
.../schema/event/SchemaChangeResultResponse.java | 74 --
.../schema/regular/SchemaCoordinator.java | 499 +++++++++
.../SchemaCoordinatorProvider.java} | 60 +-
.../operators/schema/regular/SchemaOperator.java | 296 +++++
.../{ => regular}/SchemaOperatorFactory.java | 18 +-
.../{ => regular}/event/SchemaChangeRequest.java | 23 +-
.../{ => regular}/event/SchemaChangeResponse.java | 72 +-
.../operators/sink/DataSinkFunctionOperator.java | 2 +-
.../operators/sink/DataSinkWriterOperator.java | 2 +-
.../operators/sink/SchemaEvolutionClient.java | 30 +-
.../operators/transform/PostTransformOperator.java | 3 +-
.../operators/transform/PreTransformOperator.java | 46 +-
...r.java => DistributedPrePartitionOperator.java} | 104 +-
.../runtime/partitioning/PartitioningEvent.java | 32 +-
...rator.java => RegularPrePartitionOperator.java} | 15 +-
.../runtime/serializer/event/EventSerializer.java | 7 +-
.../event/PartitioningEventSerializer.java | 12 +-
.../event/SchemaChangeEventSerializer.java | 5 +
.../schema/common/SchemaDerivatorTest.java | 602 ++++++++++
.../{coordinator => common}/SchemaManagerTest.java | 11 +-
.../operators/schema/common/SchemaTestBase.java | 121 ++
.../operators/schema/common/TableIdRouterTest.java | 84 ++
.../schema/coordinator/SchemaDerivationTest.java | 412 -------
.../schema/distributed/SchemaEvolveTest.java | 381 +++++++
.../schema/{ => regular}/SchemaEvolveTest.java | 214 ++--
.../schema/{ => regular}/SchemaOperatorTest.java | 60 +-
.../transform/PostTransformOperatorTest.java | 82 +-
.../transform/PreTransformOperatorTest.java | 42 +-
.../TransformOperatorWithSchemaEvolveTest.java | 14 +-
.../transform/UnifiedTransformOperatorTest.java | 14 +-
.../partitioning/PrePartitionOperatorTest.java | 43 +-
.../serializer/event/EventSerializerTest.java | 8 +-
.../event/PartitioningEventSerializerTest.java | 26 +-
...va => DistributedEventOperatorTestHarness.java} | 122 +-
...s.java => RegularEventOperatorTestHarness.java} | 175 +--
.../schema/TestingSchemaRegistryGateway.java | 2 +-
tools/mig-test/datastream/compile_jobs.rb | 2 +-
tools/mig-test/datastream/datastream-2.4.2/pom.xml | 180 ---
.../src/main/java/DataStreamJob.java | 54 -
tools/mig-test/datastream/datastream-3.0.0/pom.xml | 180 ---
.../src/main/java/DataStreamJob.java | 54 -
.../datastream/datastream-3.0.1/.gitignore | 38 -
tools/mig-test/datastream/datastream-3.0.1/pom.xml | 180 ---
.../src/main/java/DataStreamJob.java | 54 -
.../datastream/datastream-3.1.0/.gitignore | 38 -
.../datastream/datastream-3.1.1/.gitignore | 38 -
.../.gitignore | 0
.../{datastream-3.1.1 => datastream-3.2.0}/pom.xml | 4 +-
.../src/main/java/DataStreamJob.java | 0
.../.gitignore | 0
.../{datastream-3.1.0 => datastream-3.2.1}/pom.xml | 8 +-
.../src/main/java/DataStreamJob.java | 0
tools/mig-test/datastream/run_migration_test.rb | 2 +-
tools/mig-test/prepare_libs.rb | 54 +-
tools/mig-test/run_migration_test.rb | 43 +-
163 files changed, 10411 insertions(+), 5912 deletions(-)
create mode 100644
flink-cdc-common/src/main/java/org/apache/flink/cdc/common/utils/SchemaMergingUtils.java
create mode 100644
flink-cdc-common/src/test/java/org/apache/flink/cdc/common/utils/SchemaMergingUtilsTest.java
create mode 100644
flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/flink/FlinkParallelizedPipelineITCase.java
copy
flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/{utils/factory/DataSourceFactory1.java
=> testsource/factory/DistributedDataSourceFactory.java} (55%)
create mode 100644
flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/testsource/source/DistributedDataSource.java
create mode 100644
flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/testsource/source/DistributedSourceFunction.java
copy
flink-cdc-composer/src/test/java/org/apache/flink/cdc/composer/{utils/factory/TestOptions.java
=> testsource/source/DistributedSourceOptions.java} (53%)
create mode 100644
flink-cdc-composer/src/test/resources/ref-output/distributed-ignore.txt
create mode 100644
flink-cdc-composer/src/test/resources/ref-output/distributed.txt
create mode 100644 flink-cdc-composer/src/test/resources/ref-output/regular.txt
create mode 100644
flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-mysql/src/test/java/org/apache/flink/cdc/connectors/mysql/source/MySqlParallelizedPipelineITCase.java
delete mode 100644 flink-cdc-migration-tests/flink-cdc-release-3.0.0/pom.xml
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.0.0/src/main/java/com/ververica/cdc/migration/tests/MigrationMockBase.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.0.0/src/main/java/com/ververica/cdc/migration/tests/SchemaManagerMigrationMock.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.0.0/src/main/java/com/ververica/cdc/migration/tests/SchemaRegistryMigrationMock.java
delete mode 100644 flink-cdc-migration-tests/flink-cdc-release-3.0.1/pom.xml
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.0.1/src/main/java/com/ververica/cdc/migration/tests/MigrationMockBase.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.0.1/src/main/java/com/ververica/cdc/migration/tests/SchemaManagerMigrationMock.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.0.1/src/main/java/com/ververica/cdc/migration/tests/SchemaRegistryMigrationMock.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.1.0/src/main/java/org/apache/flink/cdc/migration/tests/SchemaRegistryMigrationMock.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.1.0/src/main/java/org/apache/flink/cdc/migration/tests/TableChangeInfoMigrationMock.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.1.1/src/main/java/org/apache/flink/cdc/migration/tests/SchemaRegistryMigrationMock.java
delete mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.1.1/src/main/java/org/apache/flink/cdc/migration/tests/TableChangeInfoMigrationMock.java
rename flink-cdc-migration-tests/{flink-cdc-release-3.1.1 =>
flink-cdc-release-3.2.0}/pom.xml (93%)
rename flink-cdc-migration-tests/{flink-cdc-release-3.1.0 =>
flink-cdc-release-3.2.0}/src/main/java/org/apache/flink/cdc/migration/tests/MigrationMockBase.java
(100%)
rename flink-cdc-migration-tests/{flink-cdc-release-3.1.0 =>
flink-cdc-release-3.2.0}/src/main/java/org/apache/flink/cdc/migration/tests/SchemaManagerMigrationMock.java
(54%)
create mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.2.0/src/main/java/org/apache/flink/cdc/migration/tests/SchemaRegistryMigrationMock.java
copy flink-cdc-migration-tests/{flink-cdc-release-snapshot =>
flink-cdc-release-3.2.0}/src/main/java/org/apache/flink/cdc/migration/tests/TableChangeInfoMigrationMock.java
(100%)
rename flink-cdc-migration-tests/{flink-cdc-release-3.1.0 =>
flink-cdc-release-3.2.1}/pom.xml (93%)
rename flink-cdc-migration-tests/{flink-cdc-release-3.1.1 =>
flink-cdc-release-3.2.1}/src/main/java/org/apache/flink/cdc/migration/tests/MigrationMockBase.java
(100%)
rename flink-cdc-migration-tests/{flink-cdc-release-3.1.1 =>
flink-cdc-release-3.2.1}/src/main/java/org/apache/flink/cdc/migration/tests/SchemaManagerMigrationMock.java
(54%)
create mode 100644
flink-cdc-migration-tests/flink-cdc-release-3.2.1/src/main/java/org/apache/flink/cdc/migration/tests/SchemaRegistryMigrationMock.java
copy flink-cdc-migration-tests/{flink-cdc-release-snapshot =>
flink-cdc-release-3.2.1}/src/main/java/org/apache/flink/cdc/migration/tests/TableChangeInfoMigrationMock.java
(100%)
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/SchemaOperator.java
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{event
=> common}/CoordinationResponseUtils.java (99%)
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/common/CoordinatorExecutorThreadFactory.java
create mode 100755
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/common/SchemaDerivator.java
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{coordinator
=> common}/SchemaManager.java (69%)
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/common/SchemaRegistry.java
create mode 100755
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/common/TableIdRouter.java
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/event/FlushSuccessEvent.java (60%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/event/GetEvolvedSchemaRequest.java (93%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/event/GetEvolvedSchemaResponse.java (91%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/event/GetOriginalSchemaRequest.java (93%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/event/GetOriginalSchemaResponse.java (91%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/event/SinkWriterRegisterEvent.java (92%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> common}/metrics/SchemaOperatorMetrics.java (96%)
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/coordinator/SchemaDerivation.java
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/coordinator/SchemaRegistry.java
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/coordinator/SchemaRegistryRequestHandler.java
create mode 100755
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/distributed/SchemaCoordinator.java
copy
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{coordinator/SchemaRegistryProvider.java
=> distributed/SchemaCoordinatorProvider.java} (56%)
mode change 100644 => 100755
create mode 100755
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/distributed/SchemaOperator.java
copy
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> distributed}/SchemaOperatorFactory.java (77%)
mode change 100644 => 100755
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/distributed/event/SchemaChangeRequest.java
copy
flink-cdc-composer/src/main/java/org/apache/flink/cdc/composer/definition/UdfDef.java
=>
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/distributed/event/SchemaChangeResponse.java
(50%)
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/event/SchemaChangeProcessingResponse.java
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/event/SchemaChangeResultRequest.java
delete mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/event/SchemaChangeResultResponse.java
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/regular/SchemaCoordinator.java
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{coordinator/SchemaRegistryProvider.java
=> regular/SchemaCoordinatorProvider.java} (56%)
create mode 100644
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/regular/SchemaOperator.java
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> regular}/SchemaOperatorFactory.java (84%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> regular}/event/SchemaChangeRequest.java (79%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/operators/schema/{
=> regular}/event/SchemaChangeResponse.java (58%)
copy
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/partitioning/{PrePartitionOperator.java
=> DistributedPrePartitionOperator.java} (52%)
rename
flink-cdc-runtime/src/main/java/org/apache/flink/cdc/runtime/partitioning/{PrePartitionOperator.java
=> RegularPrePartitionOperator.java} (93%)
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/common/SchemaDerivatorTest.java
rename
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/{coordinator
=> common}/SchemaManagerTest.java (97%)
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/common/SchemaTestBase.java
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/common/TableIdRouterTest.java
delete mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/coordinator/SchemaDerivationTest.java
create mode 100644
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/distributed/SchemaEvolveTest.java
rename
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/{
=> regular}/SchemaEvolveTest.java (95%)
rename
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/operators/schema/{
=> regular}/SchemaOperatorTest.java (73%)
copy
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/testutils/operators/{EventOperatorTestHarness.java
=> DistributedEventOperatorTestHarness.java} (62%)
rename
flink-cdc-runtime/src/test/java/org/apache/flink/cdc/runtime/testutils/operators/{EventOperatorTestHarness.java
=> RegularEventOperatorTestHarness.java} (66%)
delete mode 100644 tools/mig-test/datastream/datastream-2.4.2/pom.xml
delete mode 100644
tools/mig-test/datastream/datastream-2.4.2/src/main/java/DataStreamJob.java
delete mode 100644 tools/mig-test/datastream/datastream-3.0.0/pom.xml
delete mode 100644
tools/mig-test/datastream/datastream-3.0.0/src/main/java/DataStreamJob.java
delete mode 100644 tools/mig-test/datastream/datastream-3.0.1/.gitignore
delete mode 100644 tools/mig-test/datastream/datastream-3.0.1/pom.xml
delete mode 100644
tools/mig-test/datastream/datastream-3.0.1/src/main/java/DataStreamJob.java
delete mode 100644 tools/mig-test/datastream/datastream-3.1.0/.gitignore
delete mode 100644 tools/mig-test/datastream/datastream-3.1.1/.gitignore
rename tools/mig-test/datastream/{datastream-2.4.2 =>
datastream-3.2.0}/.gitignore (100%)
rename tools/mig-test/datastream/{datastream-3.1.1 =>
datastream-3.2.0}/pom.xml (98%)
rename tools/mig-test/datastream/{datastream-3.1.0 =>
datastream-3.2.0}/src/main/java/DataStreamJob.java (100%)
rename tools/mig-test/datastream/{datastream-3.0.0 =>
datastream-3.2.1}/.gitignore (100%)
rename tools/mig-test/datastream/{datastream-3.1.0 =>
datastream-3.2.1}/pom.xml (97%)
rename tools/mig-test/datastream/{datastream-3.1.1 =>
datastream-3.2.1}/src/main/java/DataStreamJob.java (100%)