This is an automated email from the ASF dual-hosted git repository.
pnowojski pushed a change to branch release-1.11
in repository https://gitbox.apache.org/repos/asf/flink.git.
from 4fba374 [FLINK-17763][dist] Properly handle log properties and spaces
in scala-shell.sh
new 9fedade [FLINK-17547][task][hotfix] Improve error handling 1 catch
one more invalid input in DataOutputSerializer.write 2 more informative error
messages
new 39f5f1b [FLINK-17547][task][hotfix] Extract NonSpanningWrapper from
SpillingAdaptiveSpanningRecordDeserializer (static inner class) As it is, no
logical changes.
new 5fd01ea [FLINK-17547][task][hotfix] Extract SpanningWrapper from
SpillingAdaptiveSpanningRecordDeserializer (static inner class). As it is, no
logical changes.
new c6bdeb4 [FLINK-17547][task][hotfix] Fix compiler warnings in
NonSpanningWrapper
new 9ebdaf4 [FLINK-17547][task][hotfix] Extract methods from
RecordsDeserializer
new 1c9bf03 [FLINK-17547][task] Use iterator for unconsumed buffers.
Motivation: support spilled records Changes: 1. change
SpillingAdaptiveSpanningRecordDeserializer.getUnconsumedBuffer signature 2.
adapt channel state persistence to new types
new 4e323c3 [FLINK-17547][task][hotfix] Extract RefCountedFileWithStream
from RefCountedFile Motivation: use RefCountedFile for reading as well.
new 1379548 [FLINK-17547][task][hotfix] Move RefCountedFile to flink-core
to use it in SpanningWrapper
new 3dacffe [FLINK-17547][task] Use RefCountedFile in SpanningWrapper
(todo: merge with next?)
new 2ed38c1 [FLINK-17547][task] Implement getUnconsumedSegment for
spilled buffers
The 10 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:
.../org/apache/flink/core/fs}/RefCountedFile.java | 61 +-
.../flink/core/memory/DataOutputSerializer.java | 4 +-
.../flink/core/memory/HybridMemorySegment.java | 3 +-
.../flink/core/memory/MemorySegmentFactory.java | 28 +-
.../org/apache/flink/util/CloseableIterator.java | 109 +++-
.../main/java/org/apache/flink/util/IOUtils.java | 16 +
.../java/org/apache/flink/util}/RefCounted.java | 2 +-
.../apache/flink/core/fs}/RefCountedFileTest.java | 61 +-
.../core/memory/MemorySegmentFactoryTest.java | 64 ++
.../apache/flink/util/CloseableIteratorTest.java | 82 +++
.../flink/fs/s3/common/FlinkS3FileSystem.java | 4 +-
.../utils/RefCountedBufferingFileStream.java | 10 +-
.../s3/common/utils/RefCountedFSOutputStream.java | 1 +
...ntedFile.java => RefCountedFileWithStream.java} | 69 +--
.../s3/common/utils/RefCountedTmpFileCreator.java | 10 +-
.../writer/S3RecoverableFsDataOutputStream.java | 12 +-
.../S3RecoverableMultipartUploadFactory.java | 8 +-
.../fs/s3/common/writer/S3RecoverableWriter.java | 8 +-
.../utils/RefCountedBufferingFileStreamTest.java | 4 +-
...Test.java => RefCountedFileWithStreamTest.java} | 60 +-
.../writer/RecoverableMultiPartUploadImplTest.java | 4 +-
.../S3RecoverableFsDataOutputStreamTest.java | 10 +-
.../channel/ChannelStateWriteRequest.java | 33 +-
.../ChannelStateWriteRequestDispatcherImpl.java | 6 +-
.../ChannelStateWriteRequestExecutorImpl.java | 20 +-
.../checkpoint/channel/ChannelStateWriter.java | 6 +-
.../checkpoint/channel/ChannelStateWriterImpl.java | 17 +-
.../runtime/io/disk/FileBasedBufferIterator.java | 90 +++
.../api/serialization/NonSpanningWrapper.java | 372 ++++++++++++
.../api/serialization/RecordDeserializer.java | 4 +-
.../network/api/serialization/SpanningWrapper.java | 314 ++++++++++
...SpillingAdaptiveSpanningRecordDeserializer.java | 661 ++-------------------
.../partition/consumer/RemoteInputChannel.java | 3 +-
.../ChannelStateWriteRequestDispatcherTest.java | 10 +-
.../ChannelStateWriteRequestExecutorImplTest.java | 1 -
.../channel/ChannelStateWriterImplTest.java | 13 +-
.../channel/CheckpointInProgressRequestTest.java | 7 +-
.../checkpoint/channel/MockChannelStateWriter.java | 11 +-
.../channel/RecordingChannelStateWriter.java | 12 +-
.../SpanningRecordSerializationTest.java | 9 +-
.../api/serialization/SpanningWrapperTest.java | 115 ++++
.../partition/consumer/SingleInputGateTest.java | 14 +-
.../runtime/state/ChannelPersistenceITCase.java | 3 +-
.../runtime/io/CheckpointBarrierUnaligner.java | 3 +-
.../runtime/io/StreamTaskNetworkInput.java | 11 +-
45 files changed, 1428 insertions(+), 937 deletions(-)
copy
{flink-filesystems/flink-s3-fs-base/src/main/java/org/apache/flink/fs/s3/common/utils
=> flink-core/src/main/java/org/apache/flink/core/fs}/RefCountedFile.java (60%)
rename
{flink-filesystems/flink-s3-fs-base/src/main/java/org/apache/flink/fs/s3/common/utils
=> flink-core/src/main/java/org/apache/flink/util}/RefCounted.java (96%)
copy
{flink-filesystems/flink-s3-fs-base/src/test/java/org/apache/flink/fs/s3/common/utils
=> flink-core/src/test/java/org/apache/flink/core/fs}/RefCountedFileTest.java
(56%)
create mode 100644
flink-core/src/test/java/org/apache/flink/core/memory/MemorySegmentFactoryTest.java
create mode 100644
flink-core/src/test/java/org/apache/flink/util/CloseableIteratorTest.java
rename
flink-filesystems/flink-s3-fs-base/src/main/java/org/apache/flink/fs/s3/common/utils/{RefCountedFile.java
=> RefCountedFileWithStream.java} (60%)
rename
flink-filesystems/flink-s3-fs-base/src/test/java/org/apache/flink/fs/s3/common/utils/{RefCountedFileTest.java
=> RefCountedFileWithStreamTest.java} (56%)
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/io/disk/FileBasedBufferIterator.java
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/io/network/api/serialization/NonSpanningWrapper.java
create mode 100644
flink-runtime/src/main/java/org/apache/flink/runtime/io/network/api/serialization/SpanningWrapper.java
create mode 100644
flink-runtime/src/test/java/org/apache/flink/runtime/io/network/api/serialization/SpanningWrapperTest.java