This is an automated email from the ASF dual-hosted git repository.
dwysakowicz pushed a change to branch benchmark-request
in repository https://gitbox.apache.org/repos/asf/flink.git.
discard 37dcf10 [FLINK-18934][runtime] Drop StreamStatusMaintainer &
StreamStatusProvider
discard 4358472 [squash/discuss] Remove test for not forwarding watermarks
discard e944549 [FLINK-18934][runtime] Idle stream does not advance watermark
in connected stream
omit f7530b8 [FLINK-18934][core] Extract reusable CombinedWatermark class
add d09745a [FLINK-22650][python][table-planner-blink] Support
StreamExecPythonCorrelate json serialization/deserialization
add 5eebab4 [FLINK-22652][python][table-planner-blink] Support
StreamExecPythonGroupWindowAggregate json serialization/deserialization
add 099a18e [FLINK-22666][table] Make structured type's fields more
lenient during casting
add 5746385 [hotfix][table-planner-blink] Give more helpful exception for
codegen structured types
add 80ef986 [FLINK-22620][orc] Drop BatchTableSource OrcTableSource and
related classes
add 06dec01 [FLINK-22651][python][table-planner-blink] Support
StreamExecPythonGroupAggregate json serialization/deserialization
add 4d33e85 [FLINK-22653][python][table-planner-blink] Support
StreamExecPythonOverAggregate json serialization/deserialization
add 057f86b [FLINK-12295][table] Fix comments in MinAggFunction and
MaxAggFunction
add 8cf6f84 [FLINK-20695][ha] Clean ha data for job if globally terminated
add 0528f84 [FLINK-22067][tests] Wait for vertices to using API
add 7399bb4 [hotfix][tests] Remove try/catch from
SavepointWindowReaderITCase.testApplyEvictorWindowStateReader
add ce3631a [FLINK-22622][parquet] Drop BatchTableSource
ParquetTableSource and related classes
add aa9f474 [hotfix] Ignore failing KinesisITCase traacked in FLINK-22613
add f15d277 [FLINK-19545][e2e] Add e2e test for native Kubernetes HA
add 1846bed [FLINK-22656] Fix typos
add 73fcaa1 [hotfix][runtime] Fixes JavaDoc for
RetrievableStateStorageHelper
add 15eafd9 [hotfix][runtime] Cleans up unnecessary annotations
add 9d2e2d9 [FLINK-22494][kubernetes] Introduces
PossibleInconsistentStateException
add cc59ad5 [FLINK-22494][ha] Refactors TestingLongStateHandleHelper to
operate on references
add 417cf78 [FLINK-22494][ha] Introduces PossibleInconsistentState to
StateHandleStore
add e632591 [FLINK-22494][runtime] Refactors CheckpointsCleaner to handle
also discardOnFailedStoring
add fa0d1dc [FLINK-22494][runtime] Adds
PossibleInconsistentStateException handling to CheckpointCoordinator
add 7d2a325 [FLINK-22515][docs] Add documentation for GSR-Flink
Integration
add 1a1a729 [FLINK-20487][table-planner-blink] Remove restriction on
StreamPhysicalGroupWindowAggregate which only supports insert-only input node
add 9761019 [FLINK-22623][hbase] Drop BatchTableSource/Sink
HBaseTableSource/Sink and related classes
add 927a21c [hotfix][hbase] Fix warnings around decimals in HBaseTestBase
add 79936be [FLINK-22636][zk] Group job specific zNodes under /jobs zNode
add cf3c8b1 [FLINK-22696][tests] Enable Confluent Schema Registry e2e
test on jdk 11
add a3b4a08 [FLINK-22487][python] Support `print` to print logs in PyFlink
add bcf8bd7 [FLINK-22697][examples-table] Update StreamSQLExample to
current API
add 54ba156 [FLINK-22697][examples-table] Remove TPCHQuery3Table example
add 217d703 [FLINK-22697][examples-table] Update WordCountSQLExample to
current API
add 25d1aeb [FLINK-22697][examples-table] Remove WordCountTable example
add 1310e78 [FLINK-22697][examples-table] Update flink-examples-table
pom.xml
add 938680c [hotfix][docs] Update SQL grammar for WITH and sub queries
add 3c732c7 [FLINK-22699][table-common] Declare ConstantArgumentCount as
@PublicEvolving
add 74f4b7f [FLINK-22661][hive] HiveInputFormatPartitionReader can return
invalid data
add 3558ff1 [FLINK-22704][tests] Harden
ZooKeeperHaServicesTest.testCleanupJobData
add b43ffda [FLINK-22692][tests] Disable CheckpointStoreITCase with
adaptive scheduler
add f1132f3 [FLINK-22706][legal] Update NOTICE regarding docs/ contents
add a85d7f4 [FLINK-22451][table] Support (*) as argument of UDF in Table
API (#15768)
add 4cf91cd [FLINK-22708][config] Propagate savepoint settings from
StreamExecutionEnvironment to StreamGraph
add e4b7228 [FLINK-22694][e2e] Use SQL file in TPCH end to end tests
(#15944)
add 1106542 [FLINK-22266] Fix stop-with-savepoint operation in
AdaptiveScheduler
add afe0fdd [FLINK-22155][table] EXPLAIN statement should validate insert
and query separately
add 7753542 [FLINK-22721][ha] Add default implementation for
HighAvailabilityServices.cleanupJobData
add 17d5643 [FLINK-22203][table-runtime-blink] Fix
ConcurrentModificationException for testing values sink functions (#15978)
add 1938f639 [FLINK-22713][docs] Correct document syntax errors in Kafka
page (#15957)
add f041516 [FLINK-22638] Keep channels blocked on alignment timeout
add 09168af [FLINK-22733][python] DataStream.union should handle properly
for KeyedStream in Python DataStream API
add 0a7ec4f [FLINK-22688][coordination] Eases assertion on
ExceptionHistoryEntry
add 5fd02d2 [FLINK-22434] Store suspended execution graphs on termination
to keep them accessible
add 36bd609 [FLINK-22434] Retry in clusterclient if job state is unknown
add 5a64f54 [FLINK-22376][runtime] Replaced MemorySegment inside of
BufferBuilder by Buffer
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 (37dcf10)
\
N -- N -- N refs/heads/benchmark-request (5a64f54)
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:
NOTICE | 19 +-
.../content.zh/docs/connectors/datastream/kafka.md | 19 +-
.../docs/connectors/datastream/kinesis.md | 34 +
docs/content.zh/docs/dev/python/debugging.md | 7 +-
docs/content.zh/docs/dev/table/functions/udfs.md | 56 +
.../docs/dev/table/sql/queries/overview.md | 116 +-
docs/content/docs/connectors/datastream/kafka.md | 13 +-
docs/content/docs/connectors/datastream/kinesis.md | 34 +
docs/content/docs/dev/python/debugging.md | 7 +-
docs/content/docs/dev/table/functions/udfs.md | 56 +
.../content/docs/dev/table/sql/queries/overview.md | 116 +-
.../expert_high_availability_zk_section.html | 24 -
.../generated/high_availability_configuration.html | 24 -
.../client/program/rest/RestClusterClient.java | 57 +-
.../client/program/rest/RestClusterClientTest.java | 80 +
.../base/source/reader/SourceReaderBase.java | 4 +-
flink-connectors/flink-connector-hbase-1.4/pom.xml | 53 +-
.../flink/addons/hbase1/TableInputFormat.java | 37 -
.../flink/connector/hbase1/HBase1TableFactory.java | 228 ---
.../hbase1/sink/HBaseUpsertTableSink.java | 140 --
.../connector/hbase1/source/HBaseInputFormat.java | 99 --
.../hbase1/source/HBaseRowInputFormat.java | 114 --
.../connector/hbase1/source/HBaseTableSource.java | 83 --
.../org.apache.flink.table.factories.TableFactory | 16 -
.../connector/hbase1/HBaseConnectorITCase.java | 401 ++---
.../connector/hbase1/HBaseDescriptorTest.java | 182 ---
.../connector/hbase1/HBaseTableFactoryTest.java | 214 ---
.../hbase1/example/HBaseWriteExample.java | 215 ---
.../hbase1/example/HBaseWriteStreamExample.java | 105 --
.../flink/connector/hbase1/util/HBaseTestBase.java | 31 +-
.../flink-connector-hbase-2.2/README.md | 89 --
flink-connectors/flink-connector-hbase-2.2/pom.xml | 63 +-
.../flink/connector/hbase2/HBase2TableFactory.java | 226 ---
.../hbase2/sink/HBaseUpsertTableSink.java | 140 --
.../connector/hbase2/source/HBaseInputFormat.java | 107 --
.../hbase2/source/HBaseRowInputFormat.java | 117 --
.../connector/hbase2/source/HBaseTableSource.java | 83 --
.../org.apache.flink.table.factories.TableFactory | 16 -
.../connector/hbase2/HBaseConnectorITCase.java | 402 ++---
.../connector/hbase2/HBaseDescriptorTest.java | 152 --
.../connector/hbase2/HBaseTableFactoryTest.java | 214 ---
.../hbase2/example/HBaseWriteExample.java | 215 ---
.../HBaseRowDataAsyncLookupFunctionTest.java | 7 -
.../flink/connector/hbase2/util/HBaseTestBase.java | 31 +-
.../flink-connector-hbase-base/pom.xml | 32 +-
.../hbase/sink/LegacyMutationConverter.java | 56 -
.../hbase/source/AbstractHBaseTableSource.java | 194 ---
.../hbase/source/HBaseLookupFunction.java | 158 --
.../connector/hbase/util/HBaseReadWriteHelper.java | 225 ---
.../table/descriptors/AbstractHBaseValidator.java | 75 -
.../org/apache/flink/table/descriptors/HBase.java | 173 ---
.../flink/connector/hbase/util/PlannerType.java | 25 -
.../hive/read/HiveInputFormatPartitionReader.java | 22 +-
.../read/HiveInputFormatPartitionReaderITCase.java | 101 ++
.../connectors/kinesis/FlinkKinesisITCase.java | 2 +
.../common/eventtime/CombinedWatermarkStatus.java | 130 --
.../eventtime/IndexedCombinedWatermarkStatus.java | 87 --
.../eventtime/WatermarkOutputMultiplexer.java | 100 +-
.../flink/configuration/ConfigConstants.java | 18 +-
.../configuration/HighAvailabilityOptions.java | 33 -
.../flink/table/tpch/TpchResultComparator.java | 19 +-
flink-end-to-end-tests/run-nightly-tests.sh | 5 +-
flink-end-to-end-tests/test-scripts/common.sh | 2 +-
.../test-scripts/common_kubernetes.sh | 41 +-
.../container-scripts/kubernetes-pod-template.yaml | 40 +-
.../test-scripts/test-data/tpch/sink/q1.sql | 18 +
.../test-scripts/test-data/tpch/sink/q1.yaml | 51 -
.../test-scripts/test-data/tpch/sink/q10.sql | 16 +
.../test-scripts/test-data/tpch/sink/q10.yaml | 43 -
.../test-scripts/test-data/tpch/sink/q11.sql | 10 +
.../test-scripts/test-data/tpch/sink/q11.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q12.sql | 11 +
.../test-scripts/test-data/tpch/sink/q12.yaml | 23 -
.../test-scripts/test-data/tpch/sink/q13.sql | 10 +
.../test-scripts/test-data/tpch/sink/q13.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q14.sql | 9 +
.../test-scripts/test-data/tpch/sink/q14.yaml | 15 -
.../test-scripts/test-data/tpch/sink/q15.sql | 13 +
.../test-scripts/test-data/tpch/sink/q15.yaml | 31 -
.../test-scripts/test-data/tpch/sink/q16.sql | 12 +
.../test-scripts/test-data/tpch/sink/q16.yaml | 27 -
.../test-scripts/test-data/tpch/sink/q17.sql | 9 +
.../test-scripts/test-data/tpch/sink/q17.yaml | 15 -
.../test-scripts/test-data/tpch/sink/q18.sql | 14 +
.../test-scripts/test-data/tpch/sink/q18.yaml | 35 -
.../test-scripts/test-data/tpch/sink/q19.sql | 9 +
.../test-scripts/test-data/tpch/sink/q19.yaml | 15 -
.../test-scripts/test-data/tpch/sink/q2.sql | 16 +
.../test-scripts/test-data/tpch/sink/q2.yaml | 43 -
.../test-scripts/test-data/tpch/sink/q20.sql | 10 +
.../test-scripts/test-data/tpch/sink/q20.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q21.sql | 10 +
.../test-scripts/test-data/tpch/sink/q21.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q22.sql | 11 +
.../test-scripts/test-data/tpch/sink/q22.yaml | 23 -
.../test-scripts/test-data/tpch/sink/q3.sql | 12 +
.../test-scripts/test-data/tpch/sink/q3.yaml | 27 -
.../test-scripts/test-data/tpch/sink/q4.sql | 10 +
.../test-scripts/test-data/tpch/sink/q4.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q5.sql | 10 +
.../test-scripts/test-data/tpch/sink/q5.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q6.sql | 9 +
.../test-scripts/test-data/tpch/sink/q6.yaml | 15 -
.../test-scripts/test-data/tpch/sink/q7.sql | 12 +
.../test-scripts/test-data/tpch/sink/q7.yaml | 27 -
.../test-scripts/test-data/tpch/sink/q8.sql | 10 +
.../test-scripts/test-data/tpch/sink/q8.yaml | 19 -
.../test-scripts/test-data/tpch/sink/q9.sql | 11 +
.../test-scripts/test-data/tpch/sink/q9.yaml | 23 -
.../test-scripts/test-data/tpch/source.sql | 133 ++
.../test-scripts/test_kubernetes_application.sh | 2 +-
...cation.sh => test_kubernetes_application_ha.sh} | 37 +-
.../test_kubernetes_pyflink_application.sh | 2 +-
flink-end-to-end-tests/test-scripts/test_tpch.sh | 31 +-
flink-examples/flink-examples-table/pom.xml | 92 +-
.../examples/java/basics/StreamSQLExample.java | 78 +-
.../table/examples/java/basics/WordCountSQL.java | 84 --
.../examples/java/basics/WordCountSQLExample.java | 48 +
.../table/examples/java/basics/WordCountTable.java | 81 -
.../examples/scala/basics/StreamSQLExample.scala | 75 +-
.../examples/scala/basics/TPCHQuery3Table.scala | 180 ---
.../table/examples/scala/basics/WordCountSQL.scala | 62 -
.../scala/basics/WordCountSQLExample.scala | 48 +
.../examples/scala/basics/WordCountTable.scala | 62 -
.../java/basics/StreamSQLExampleITCase.java | 36 +-
.../java/basics/WordCountSQLExampleITCase.java | 37 +-
.../scala/basics/StreamSQLExampleITCase.java | 39 +-
.../scala/basics/WordCountSQLExampleITCase.java | 36 +-
flink-formats/flink-orc-nohive/pom.xml | 11 +-
flink-formats/flink-orc/pom.xml | 29 +-
.../java/org/apache/flink/orc/OrcBatchReader.java | 1559 --------------------
.../java/org/apache/flink/orc/OrcInputFormat.java | 220 ---
.../org/apache/flink/orc/OrcRowInputFormat.java | 110 --
.../org/apache/flink/orc/OrcRowSplitReader.java | 79 -
.../java/org/apache/flink/orc/OrcTableSource.java | 576 --------
.../org/apache/flink/orc/OrcBatchReaderTest.java | 156 --
.../apache/flink/orc/OrcRowInputFormatTest.java | 1103 --------------
.../org/apache/flink/orc/OrcTableSourceITCase.java | 118 --
.../org/apache/flink/orc/OrcTableSourceTest.java | 384 -----
.../flink/orc/util/OrcTestFileGenerator.java | 374 -----
flink-formats/flink-parquet/pom.xml | 30 +-
.../flink/formats/parquet/ParquetInputFormat.java | 323 ----
.../formats/parquet/ParquetMapInputFormat.java | 143 --
.../formats/parquet/ParquetPojoInputFormat.java | 138 --
.../formats/parquet/ParquetRowInputFormat.java | 49 -
.../flink/formats/parquet/ParquetTableSource.java | 616 --------
.../formats/parquet/utils/ParquetRecordReader.java | 311 ----
.../parquet/utils/ParquetSchemaConverter.java | 510 -------
.../parquet/utils/ParquetTimestampUtils.java | 69 -
.../flink/formats/parquet/utils/RowConverter.java | 541 -------
.../formats/parquet/utils/RowMaterializer.java | 49 -
.../formats/parquet/utils/RowReadSupport.java | 54 -
.../formats/parquet/ParquetMapInputFormatTest.java | 125 --
.../parquet/ParquetPojoInputFormatTest.java | 115 --
.../formats/parquet/ParquetRowInputFormatTest.java | 460 ------
.../formats/parquet/ParquetTableSourceITCase.java | 117 --
.../formats/parquet/ParquetTableSourceTest.java | 252 ----
.../parquet/utils/ParquetRecordReaderTest.java | 380 -----
.../parquet/utils/ParquetSchemaConverterTest.java | 189 ---
.../flink/formats/parquet/utils/TestUtil.java | 302 ----
.../src/test/resources/avro/nested.avsc | 35 -
.../src/test/resources/avro/simple.avsc | 12 -
.../highavailability/KubernetesHaServices.java | 19 +-
.../KubernetesStateHandleStore.java | 74 +-
.../kubeclient/Fabric8FlinkKubeClient.java | 11 +-
.../kubernetes/kubeclient/FlinkKubeClient.java | 8 +-
.../flink/kubernetes/KubernetesClientTestBase.java | 35 +-
.../flink/kubernetes/MixedKubernetesServer.java | 6 +-
.../highavailability/KubernetesHaServicesTest.java | 34 +
.../KubernetesStateHandleStoreITCase.java | 10 +-
.../KubernetesStateHandleStoreTest.java | 482 +++---
.../kubeclient/Fabric8FlinkKubeClientTest.java | 49 +-
.../flink/state/api/output/BoundedStreamTask.java | 4 -
.../operators/StateBootstrapWrapperOperator.java | 9 -
.../flink/state/api/SavepointReaderITTestBase.java | 21 +-
.../state/api/SavepointReaderKeyedStateITCase.java | 44 +-
.../state/api/SavepointWindowReaderITCase.java | 368 ++---
.../flink/state/api/utils/SavepointTestBase.java | 45 +-
.../flink/state/api/utils/WaitingFunction.java | 25 -
.../flink/state/api/utils/WaitingSource.java | 37 +-
.../state/api/utils/WaitingWindowAssigner.java | 72 -
flink-python/pyflink/datastream/data_stream.py | 5 +-
.../pyflink/datastream/tests/test_data_stream.py | 23 +-
.../fn_execution/beam/beam_sdk_worker_main.py | 38 +
flink-python/pyflink/table/descriptors.py | 151 --
.../pyflink/table/tests/test_descriptor.py | 110 +-
.../pyflink/table/tests/test_pandas_udaf.py | 70 +
flink-python/pyflink/table/tests/test_udaf.py | 105 ++
flink-python/pyflink/table/tests/test_udtf.py | 50 +-
.../PythonTimestampsAndWatermarksOperator.java | 4 +-
.../runtime/checkpoint/CheckpointCoordinator.java | 26 +-
.../runtime/checkpoint/CheckpointsCleaner.java | 46 +-
.../DefaultCompletedCheckpointStore.java | 4 +
.../checkpoint/ZooKeeperCheckpointIDCounter.java | 14 +-
.../ZooKeeperCheckpointRecoveryFactory.java | 10 +-
.../channel/RecoveredChannelStateHandler.java | 47 +-
.../channel/SequentialChannelStateReaderImpl.java | 3 +-
.../flink/runtime/dispatcher/Dispatcher.java | 52 +-
.../flink/runtime/dispatcher/MiniDispatcher.java | 14 +-
.../highavailability/AbstractHaServices.java | 58 +-
.../highavailability/HighAvailabilityServices.java | 10 +-
.../HighAvailabilityServicesUtils.java | 2 +-
.../zookeeper/ZooKeeperClientHAServices.java | 10 +-
.../zookeeper/ZooKeeperHaServices.java | 99 +-
.../network/api/serialization/EventSerializer.java | 4 +-
.../runtime/io/network/buffer/BufferBuilder.java | 35 +-
.../runtime/io/network/buffer/BufferConsumer.java | 32 +-
.../runtime/io/network/buffer/BufferProvider.java | 5 +
.../runtime/io/network/buffer/LocalBufferPool.java | 16 +-
.../partition/BufferWritingResultPartition.java | 2 +
.../network/partition/PartitionSortedBuffer.java | 6 +-
.../partition/SortMergeResultPartition.java | 3 +-
.../runtime/jobgraph/SavepointRestoreSettings.java | 7 +-
.../JobMasterServiceLeadershipRunner.java | 21 +-
.../apache/flink/runtime/jobmaster/JobResult.java | 4 +-
.../ZooKeeperLeaderElectionDriver.java | 37 +-
.../ZooKeeperLeaderElectionDriverFactory.java | 17 +-
.../ZooKeeperLeaderRetrievalDriver.java | 23 +-
.../PossibleInconsistentStateException.java | 31 +-
.../persistence/RetrievableStateStorageHelper.java | 7 +-
.../runtime/persistence/StateHandleStore.java | 13 +-
.../rest/handler/job/JobExceptionsHandler.java | 3 -
.../runtime/scheduler/adaptive/Executing.java | 7 +-
.../exceptionhistory/ExceptionHistoryEntry.java | 3 +
.../flink/runtime/util/StateHandleStoreUtils.java | 76 +
.../apache/flink/runtime/util/ZooKeeperUtils.java | 215 +--
.../zookeeper/ZooKeeperStateHandleStore.java | 152 +-
.../CheckpointCoordinatorFailureTest.java | 96 +-
.../CheckpointCoordinatorTestingUtils.java | 6 +
.../ZKCheckpointIDCounterMultiServersTest.java | 2 +-
.../ZooKeeperCheckpointIDCounterITCase.java | 18 +-
.../channel/ChannelStateChunkReaderTest.java | 16 +
.../channel/ChannelStateSerializerImplTest.java | 5 +-
.../channel/RecoveredChannelStateHandlerTest.java | 234 +++
.../dispatcher/DispatcherResourceCleanupTest.java | 46 +-
.../flink/runtime/dispatcher/DispatcherTest.java | 134 +-
.../runtime/dispatcher/TestingDispatcher.java | 2 +-
.../ZooKeeperDefaultDispatcherRunnerTest.java | 4 +-
.../highavailability/AbstractHaServicesTest.java | 50 +-
.../TestingHighAvailabilityServices.java | 12 +
.../zookeeper/ZooKeeperHaServicesTest.java | 68 +-
.../SpanningRecordSerializationTest.java | 8 +-
.../buffer/BufferBuilderAndConsumerTest.java | 58 +-
.../io/network/buffer/BufferBuilderTestUtils.java | 29 +-
.../BufferConsumerWithPartialRecordLengthTest.java | 5 +-
.../io/network/buffer/LocalBufferPoolTest.java | 13 +-
.../runtime/io/network/buffer/NoOpBufferPool.java | 10 +
.../io/network/buffer/UnpooledBufferPool.java | 10 +-
.../BoundedBlockingSubpartitionWriteReadTest.java | 5 +-
.../partition/consumer/LocalInputChannelTest.java | 11 +-
.../runtime/io/network/util/TestBufferFactory.java | 11 +-
.../io/network/util/TestPooledBufferProvider.java | 62 +-
.../io/network/util/TestSubpartitionProducer.java | 5 +-
.../JobMasterServiceLeadershipRunnerTest.java | 17 +-
.../runtime/jobmaster/TestingJobManagerRunner.java | 9 +-
.../LeaderChangeClusterComponentsTest.java | 18 +-
.../runtime/leaderelection/LeaderElectionTest.java | 2 +-
...KeeperLeaderElectionConnectionHandlingTest.java | 32 +-
.../ZooKeeperLeaderElectionTest.java | 42 +-
.../persistence/TestingLongStateHandleHelper.java | 104 +-
.../rest/handler/job/JobExceptionsHandlerTest.java | 27 +
.../runtime/scheduler/DefaultSchedulerTest.java | 42 +
.../runtime/scheduler/adaptive/CancelingTest.java | 4 +-
.../runtime/scheduler/adaptive/ExecutingTest.java | 153 +-
.../runtime/scheduler/adaptive/FailingTest.java | 13 +-
.../runtime/scheduler/adaptive/RestartingTest.java | 9 +-
.../scheduler/adaptive/StopWithSavepointTest.java | 38 +-
.../ExceptionHistoryEntryTest.java | 37 +
.../runtime/util/StateHandleStoreUtilsTest.java | 114 ++
.../flink/runtime/util/ZooKeeperUtilsTest.java | 47 +
.../zookeeper/ZooKeeperStateHandleStoreTest.java | 416 ++++--
.../source/ContinuousFileReaderOperator.java | 1 +
.../streaming/api/graph/StreamGraphGenerator.java | 3 +-
.../streaming/api/operators/AbstractInput.java | 6 -
.../api/operators/AbstractStreamOperator.java | 54 +-
.../api/operators/AbstractStreamOperatorV2.java | 30 +-
.../streaming/api/operators/CountingOutput.java | 6 -
.../flink/streaming/api/operators/Input.java | 3 -
.../flink/streaming/api/operators/Output.java | 3 -
.../streaming/api/operators/StreamSource.java | 10 +-
.../api/operators/StreamSourceContexts.java | 53 +-
.../api/operators/TimestampedCollector.java | 6 -
.../api/operators/TwoInputStreamOperator.java | 5 -
.../streaming/runtime/io/AbstractDataOutput.java | 46 +
.../streaming/runtime/io/RecordWriterOutput.java | 49 +-
.../io/StreamMultipleInputProcessorFactory.java | 108 +-
.../runtime/io/StreamTwoInputProcessorFactory.java | 62 +-
...tractAlternatingAlignedBarrierHandlerState.java | 9 +-
.../AlternatingCollectingBarriers.java | 7 +-
...=> AlternatingCollectingBarriersUnaligned.java} | 33 +-
.../AlternatingWaitingForFirstBarrier.java | 10 +-
...lternatingWaitingForFirstBarrierUnaligned.java} | 44 +-
.../runtime/io/checkpointing/ChannelState.java | 11 +
.../io/checkpointing/CollectingBarriers.java | 1 +
.../SingleCheckpointBarrierHandler.java | 4 +-
.../io/checkpointing/WaitingForFirstBarrier.java | 1 +
.../operators/TimestampsAndWatermarksOperator.java | 12 +-
.../streamstatus/StreamStatusMaintainer.java | 20 +-
.../runtime/streamstatus/StreamStatusProvider.java | 16 +-
.../runtime/tasks/BroadcastingOutputCollector.java | 20 +-
.../streaming/runtime/tasks/ChainingOutput.java | 30 +-
.../tasks/CopyingBroadcastingOutputCollector.java | 6 +-
.../runtime/tasks/CopyingChainingOutput.java | 6 +-
.../runtime/tasks/MultipleInputStreamTask.java | 1 +
.../runtime/tasks/OneInputStreamTask.java | 15 +-
.../streaming/runtime/tasks/OperatorChain.java | 43 +-
.../runtime/tasks/SourceOperatorStreamTask.java | 17 +-
.../streaming/runtime/tasks/SourceStreamTask.java | 2 +-
.../runtime/tasks/StreamIterationTail.java | 4 -
.../flink/streaming/runtime/tasks/StreamTask.java | 5 +
.../runtime/tasks/TwoInputStreamTask.java | 1 +
.../api/graph/StreamGraphGeneratorTest.java | 43 +
.../AbstractUdfStreamOperatorLifecycleTest.java | 4 +-
.../api/operators/MockStreamStatusMaintainer.java} | 22 +-
.../StreamSourceContextIdleDetectionTests.java | 116 +-
.../checkpointing/AlternatingCheckpointsTest.java | 39 +-
.../DemultiplexingRecordDeserializerTest.java | 20 +-
.../StreamSourceOperatorLatencyMetricsTest.java | 7 +-
.../StreamSourceOperatorWatermarksTest.java | 8 +
.../runtime/tasks/MultipleInputStreamTaskTest.java | 45 +-
.../runtime/tasks/OneInputStreamTaskTest.java | 209 +++
.../streaming/runtime/tasks/OperatorChainTest.java | 7 +-
.../streaming/runtime/tasks/StreamTaskTest.java | 2 +
.../tasks/SubtaskCheckpointCoordinatorTest.java | 4 -
.../runtime/tasks/TwoInputStreamTaskTest.java | 4 +-
.../util/AbstractStreamOperatorTestHarness.java | 10 +-
.../flink/streaming/util/CollectorOutput.java | 6 -
.../apache/flink/streaming/util/MockOutput.java | 6 -
.../flink/streaming/util/MockStreamTask.java | 9 +
.../streaming/util/MockStreamTaskBuilder.java | 10 +
.../src/main/codegen/includes/parserImpls.ftl | 9 +-
.../src/main/codegen/includes/parserImpls.ftl | 9 +-
.../flink/sql/parser/dql/SqlRichExplain.java | 2 +-
.../resolver/rules/ExpandColumnFunctionsRule.java | 6 +
.../rules/StarReferenceFlatteningRule.java | 10 +
.../resolver/ExpressionResolverTest.java | 18 +
.../table/types/extraction/ExtractionUtils.java | 2 +-
.../flink/table/types/inference/ArgumentCount.java | 2 +
.../types/inference/ConstantArgumentCount.java | 4 +-
.../types/logical/utils/LogicalTypeCasts.java | 58 +-
.../table/types/LogicalTypeCastAvoidanceTest.java | 42 +-
.../flink/table/types/LogicalTypeCastsTest.java | 36 +
.../functions/aggfunctions/MaxAggFunction.java | 9 +-
.../functions/aggfunctions/MinAggFunction.java | 4 +-
.../operations/SqlToOperationConverter.java | 13 +-
.../nodes/exec/batch/BatchExecPythonCorrelate.java | 12 +-
.../exec/common/CommonExecPythonCorrelate.java | 26 +-
.../exec/stream/StreamExecPythonCorrelate.java | 28 +-
.../stream/StreamExecPythonGroupAggregate.java | 53 +-
.../StreamExecPythonGroupWindowAggregate.java | 74 +-
.../exec/stream/StreamExecPythonOverAggregate.java | 30 +-
.../rules/logical/PythonCorrelateSplitRule.java | 66 +-
.../table/planner/calcite/FlinkPlannerImpl.scala | 13 +-
.../batch/BatchPhysicalPythonCorrelate.scala | 10 +-
.../stream/StreamPhysicalPythonCorrelate.scala | 13 +-
.../StreamPhysicalPythonGroupWindowAggregate.scala | 1 -
.../FlinkChangelogModeInferenceProgram.scala | 18 +-
.../plan/rules/logical/PythonCalcSplitRule.scala | 10 +-
.../factories/TestValuesRuntimeFunctions.java | 116 +-
.../operations/SqlToOperationConverterTest.java | 19 +
.../nodes/exec/stream/JsonSerdeCoverageTest.java | 4 -
.../exec/stream/PythonCorrelateJsonPlanTest.java | 92 ++
.../stream/PythonGroupAggregateJsonPlanTest.java | 70 +
.../PythonGroupWindowAggregateJsonPlanTest.java | 161 ++
.../stream/PythonOverAggregateJsonPlanTest.java | 147 ++
.../explain/testExecuteSqlWithExplainInsert.out | 2 +-
.../testJoinWithFilter.out | 431 ++++++
.../testPythonTableFunction.out | 329 +++++
.../tesPythonAggCallsWithGroupBy.out | 317 ++++
.../testEventTimeHopWindow.out | 425 ++++++
.../testEventTimeSessionWindow.out | 423 ++++++
.../testEventTimeTumbleWindow.out | 599 ++++++++
.../testProcTimeHopWindow.out | 464 ++++++
.../testProcTimeSessionWindow.out | 462 ++++++
.../testProcTimeTumbleWindow.out | 593 ++++++++
.../testProcTimeBoundedNonPartitionedRangeOver.out | 490 ++++++
.../testProcTimeBoundedPartitionedRangeOver.out | 504 +++++++
...undedPartitionedRowsOverWithBuiltinProctime.out | 420 ++++++
.../testProcTimeUnboundedPartitionedRangeOver.out | 484 ++++++
.../testRowTimeBoundedPartitionedRowsOver.out | 415 ++++++
.../CalcPythonCorrelateTransposeRuleTest.xml | 8 +-
.../rules/logical/PythonCorrelateSplitRuleTest.xml | 17 +-
.../planner/plan/stream/sql/DeduplicateTest.xml | 33 +
.../planner/plan/stream/sql/TableScanTest.xml | 51 +-
.../plan/stream/sql/agg/GroupWindowTest.xml | 256 ++--
.../planner/plan/stream/sql/DeduplicateTest.scala | 3 -
.../planner/plan/stream/sql/TableScanTest.scala | 6 +-
.../plan/stream/sql/agg/GroupWindowTest.scala | 42 +
.../runtime/stream/sql/DataStreamScalaITCase.scala | 119 ++
.../runtime/stream/sql/GroupWindowITCase.scala | 144 +-
.../planner/runtime/stream/table/CalcITCase.scala | 55 +
.../table/sqlexec/SqlToOperationConverter.java | 13 +-
.../flink/table/calcite/FlinkPlannerImpl.scala | 12 +-
.../table/sqlexec/SqlToOperationConverterTest.java | 23 +
.../resources/testExecuteSqlWithExplainInsert0.out | 2 +-
.../resources/testExecuteSqlWithExplainInsert1.out | 13 +-
.../data/conversion/StructuredObjectConverter.java | 29 +-
.../multipleinput/input/FirstInputOfTwoInput.java | 6 -
.../operators/multipleinput/input/OneInput.java | 6 -
.../multipleinput/input/SecondInputOfTwoInput.java | 6 -
.../multipleinput/output/BroadcastingOutput.java | 8 -
...gSecondInputOfTwoInputStreamOperatorOutput.java | 10 -
.../FirstInputOfTwoInputStreamOperatorOutput.java | 10 -
.../output/OneInputStreamOperatorOutput.java | 10 -
.../SecondInputOfTwoInputStreamOperatorOutput.java | 10 -
.../wmassigners/WatermarkAssignerOperator.java | 10 +-
.../multipleinput/output/BlackHoleOutput.java | 6 -
.../over/NonBufferOverWindowOperatorTest.java | 6 -
.../wmassigners/WatermarkAssignerOperatorTest.java | 32 +-
.../WatermarkAssignerOperatorTestBase.java | 8 -
.../test/checkpointing/CheckpointStoreITCase.java | 17 +-
.../runtime/SortingBoundedInputITCase.java | 4 -
412 files changed, 14090 insertions(+), 18315 deletions(-)
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/addons/hbase1/TableInputFormat.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/HBase1TableFactory.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/sink/HBaseUpsertTableSink.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/source/HBaseInputFormat.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/source/HBaseRowInputFormat.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/source/HBaseTableSource.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/test/java/org/apache/flink/connector/hbase1/HBaseDescriptorTest.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/test/java/org/apache/flink/connector/hbase1/HBaseTableFactoryTest.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/test/java/org/apache/flink/connector/hbase1/example/HBaseWriteExample.java
delete mode 100644
flink-connectors/flink-connector-hbase-1.4/src/test/java/org/apache/flink/connector/hbase1/example/HBaseWriteStreamExample.java
delete mode 100644 flink-connectors/flink-connector-hbase-2.2/README.md
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/main/java/org/apache/flink/connector/hbase2/HBase2TableFactory.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/main/java/org/apache/flink/connector/hbase2/sink/HBaseUpsertTableSink.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/main/java/org/apache/flink/connector/hbase2/source/HBaseInputFormat.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/main/java/org/apache/flink/connector/hbase2/source/HBaseRowInputFormat.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/main/java/org/apache/flink/connector/hbase2/source/HBaseTableSource.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/main/resources/META-INF/services/org.apache.flink.table.factories.TableFactory
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/test/java/org/apache/flink/connector/hbase2/HBaseDescriptorTest.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/test/java/org/apache/flink/connector/hbase2/HBaseTableFactoryTest.java
delete mode 100644
flink-connectors/flink-connector-hbase-2.2/src/test/java/org/apache/flink/connector/hbase2/example/HBaseWriteExample.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/main/java/org/apache/flink/connector/hbase/sink/LegacyMutationConverter.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/main/java/org/apache/flink/connector/hbase/source/AbstractHBaseTableSource.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/main/java/org/apache/flink/connector/hbase/source/HBaseLookupFunction.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/main/java/org/apache/flink/connector/hbase/util/HBaseReadWriteHelper.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/main/java/org/apache/flink/table/descriptors/AbstractHBaseValidator.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/main/java/org/apache/flink/table/descriptors/HBase.java
delete mode 100644
flink-connectors/flink-connector-hbase-base/src/test/java/org/apache/flink/connector/hbase/util/PlannerType.java
create mode 100644
flink-connectors/flink-connector-hive/src/test/java/org/apache/flink/connectors/hive/read/HiveInputFormatPartitionReaderITCase.java
delete mode 100644
flink-core/src/main/java/org/apache/flink/api/common/eventtime/CombinedWatermarkStatus.java
delete mode 100644
flink-core/src/main/java/org/apache/flink/api/common/eventtime/IndexedCombinedWatermarkStatus.java
copy flink-python/pyflink/fn_execution/beam/beam_sdk_worker_main.py =>
flink-end-to-end-tests/test-scripts/container-scripts/kubernetes-pod-template.yaml
(56%)
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q1.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q1.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q10.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q10.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q11.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q11.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q12.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q12.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q13.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q13.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q14.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q14.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q15.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q15.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q16.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q16.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q17.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q17.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q18.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q18.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q19.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q19.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q2.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q2.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q20.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q20.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q21.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q21.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q22.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q22.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q3.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q3.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q4.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q4.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q5.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q5.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q6.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q6.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q7.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q7.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q8.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q8.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q9.sql
delete mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/sink/q9.yaml
create mode 100644
flink-end-to-end-tests/test-scripts/test-data/tpch/source.sql
copy flink-end-to-end-tests/test-scripts/{test_kubernetes_application.sh =>
test_kubernetes_application_ha.sh} (63%)
delete mode 100644
flink-examples/flink-examples-table/src/main/java/org/apache/flink/table/examples/java/basics/WordCountSQL.java
create mode 100644
flink-examples/flink-examples-table/src/main/java/org/apache/flink/table/examples/java/basics/WordCountSQLExample.java
delete mode 100644
flink-examples/flink-examples-table/src/main/java/org/apache/flink/table/examples/java/basics/WordCountTable.java
delete mode 100644
flink-examples/flink-examples-table/src/main/scala/org/apache/flink/table/examples/scala/basics/TPCHQuery3Table.scala
delete mode 100644
flink-examples/flink-examples-table/src/main/scala/org/apache/flink/table/examples/scala/basics/WordCountSQL.scala
create mode 100644
flink-examples/flink-examples-table/src/main/scala/org/apache/flink/table/examples/scala/basics/WordCountSQLExample.scala
delete mode 100644
flink-examples/flink-examples-table/src/main/scala/org/apache/flink/table/examples/scala/basics/WordCountTable.scala
copy
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/HBaseValidator.java
=>
flink-examples/flink-examples-table/src/test/java/org/apache/flink/table/examples/java/basics/StreamSQLExampleITCase.java
(50%)
copy
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/HBaseValidator.java
=>
flink-examples/flink-examples-table/src/test/java/org/apache/flink/table/examples/java/basics/WordCountSQLExampleITCase.java
(51%)
rename
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/pojo/PojoSimpleRecord.java
=>
flink-examples/flink-examples-table/src/test/java/org/apache/flink/table/examples/scala/basics/StreamSQLExampleITCase.java
(53%)
rename
flink-connectors/flink-connector-hbase-2.2/src/main/java/org/apache/flink/connector/hbase2/HBaseValidator.java
=>
flink-examples/flink-examples-table/src/test/java/org/apache/flink/table/examples/scala/basics/WordCountSQLExampleITCase.java
(52%)
delete mode 100644
flink-formats/flink-orc/src/main/java/org/apache/flink/orc/OrcBatchReader.java
delete mode 100644
flink-formats/flink-orc/src/main/java/org/apache/flink/orc/OrcInputFormat.java
delete mode 100644
flink-formats/flink-orc/src/main/java/org/apache/flink/orc/OrcRowInputFormat.java
delete mode 100644
flink-formats/flink-orc/src/main/java/org/apache/flink/orc/OrcRowSplitReader.java
delete mode 100644
flink-formats/flink-orc/src/main/java/org/apache/flink/orc/OrcTableSource.java
delete mode 100644
flink-formats/flink-orc/src/test/java/org/apache/flink/orc/OrcBatchReaderTest.java
delete mode 100644
flink-formats/flink-orc/src/test/java/org/apache/flink/orc/OrcRowInputFormatTest.java
delete mode 100644
flink-formats/flink-orc/src/test/java/org/apache/flink/orc/OrcTableSourceITCase.java
delete mode 100644
flink-formats/flink-orc/src/test/java/org/apache/flink/orc/OrcTableSourceTest.java
delete mode 100644
flink-formats/flink-orc/src/test/java/org/apache/flink/orc/util/OrcTestFileGenerator.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/ParquetInputFormat.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/ParquetMapInputFormat.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/ParquetPojoInputFormat.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/ParquetRowInputFormat.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/ParquetTableSource.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/utils/ParquetRecordReader.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/utils/ParquetTimestampUtils.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/utils/RowConverter.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/utils/RowMaterializer.java
delete mode 100644
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/utils/RowReadSupport.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/ParquetMapInputFormatTest.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/ParquetPojoInputFormatTest.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/ParquetRowInputFormatTest.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/ParquetTableSourceITCase.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/ParquetTableSourceTest.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/utils/ParquetRecordReaderTest.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/utils/ParquetSchemaConverterTest.java
delete mode 100644
flink-formats/flink-parquet/src/test/java/org/apache/flink/formats/parquet/utils/TestUtil.java
delete mode 100644
flink-formats/flink-parquet/src/test/resources/avro/nested.avsc
delete mode 100644
flink-formats/flink-parquet/src/test/resources/avro/simple.avsc
delete mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/utils/WaitingFunction.java
delete mode 100644
flink-libraries/flink-state-processing-api/src/test/java/org/apache/flink/state/api/utils/WaitingWindowAssigner.java
rename
flink-connectors/flink-connector-hbase-1.4/src/main/java/org/apache/flink/connector/hbase1/HBaseValidator.java
=>
flink-runtime/src/main/java/org/apache/flink/runtime/persistence/PossibleInconsistentStateException.java
(52%)
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/util/StateHandleStoreUtils.java
create mode 100644
flink-runtime/src/test/java/org/apache/flink/runtime/checkpoint/channel/RecoveredChannelStateHandlerTest.java
create mode 100644
flink-runtime/src/test/java/org/apache/flink/runtime/util/StateHandleStoreUtilsTest.java
create mode 100644
flink-runtime/src/test/java/org/apache/flink/runtime/util/ZooKeeperUtilsTest.java
create mode 100644
flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/AbstractDataOutput.java
rename
flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/checkpointing/{CollectingBarriersUnaligned.java
=> AlternatingCollectingBarriersUnaligned.java} (68%)
rename
flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/io/checkpointing/{WaitingForFirstBarrierUnaligned.java
=> AlternatingWaitingForFirstBarrierUnaligned.java} (66%)
rename
flink-connectors/flink-connector-hbase-base/src/test/java/org/apache/flink/connector/hbase/example/HBaseFlinkTestConstants.java
=>
flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/streamstatus/StreamStatusMaintainer.java
(58%)
rename
flink-formats/flink-parquet/src/main/java/org/apache/flink/formats/parquet/utils/ParentDataHolder.java
=>
flink-streaming-java/src/main/java/org/apache/flink/streaming/runtime/streamstatus/StreamStatusProvider.java
(69%)
copy
flink-streaming-java/src/{main/java/org/apache/flink/streaming/runtime/io/checkpointing/WaitingForFirstBarrier.java
=>
test/java/org/apache/flink/streaming/api/operators/MockStreamStatusMaintainer.java}
(51%)
create mode 100644
flink-table/flink-table-planner-blink/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonCorrelateJsonPlanTest.java
create mode 100644
flink-table/flink-table-planner-blink/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupAggregateJsonPlanTest.java
create mode 100644
flink-table/flink-table-planner-blink/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest.java
create mode 100644
flink-table/flink-table-planner-blink/src/test/java/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonOverAggregateJsonPlanTest.java
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonCorrelateJsonPlanTest_jsonplan/testJoinWithFilter.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonCorrelateJsonPlanTest_jsonplan/testPythonTableFunction.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupAggregateJsonPlanTest_jsonplan/tesPythonAggCallsWithGroupBy.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest_jsonplan/testEventTimeHopWindow.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest_jsonplan/testEventTimeSessionWindow.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest_jsonplan/testEventTimeTumbleWindow.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest_jsonplan/testProcTimeHopWindow.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest_jsonplan/testProcTimeSessionWindow.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonGroupWindowAggregateJsonPlanTest_jsonplan/testProcTimeTumbleWindow.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonOverAggregateJsonPlanTest_jsonplan/testProcTimeBoundedNonPartitionedRangeOver.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonOverAggregateJsonPlanTest_jsonplan/testProcTimeBoundedPartitionedRangeOver.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonOverAggregateJsonPlanTest_jsonplan/testProcTimeBoundedPartitionedRowsOverWithBuiltinProctime.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonOverAggregateJsonPlanTest_jsonplan/testProcTimeUnboundedPartitionedRangeOver.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/resources/org/apache/flink/table/planner/plan/nodes/exec/stream/PythonOverAggregateJsonPlanTest_jsonplan/testRowTimeBoundedPartitionedRowsOver.out
create mode 100644
flink-table/flink-table-planner-blink/src/test/scala/org/apache/flink/table/planner/runtime/stream/sql/DataStreamScalaITCase.scala