Modifier and Type | Method and Description |
---|---|
static <IN,OUT,E extends Throwable> |
ConfigUtils.decodeListFromConfig(ReadableConfig configuration,
ConfigOption<List<IN>> key,
FunctionWithException<IN,OUT,E> mapper)
Gets a
List of values of type IN from a ReadableConfig and transforms
it to a List of type OUT based on the provided mapper function. |
Modifier and Type | Method and Description |
---|---|
static <T> List<T> |
PulsarSerdeUtils.deserializeList(DataInputStream in,
FunctionWithException<DataInputStream,T,IOException> deserializer) |
static <K,V> Map<K,V> |
PulsarSerdeUtils.deserializeMap(DataInputStream in,
FunctionWithException<DataInputStream,K,IOException> keyDeserializer,
FunctionWithException<DataInputStream,V,IOException> valueDeserializer) |
static <K,V> Map<K,V> |
PulsarSerdeUtils.deserializeMap(DataInputStream in,
FunctionWithException<DataInputStream,K,IOException> keyDeserializer,
FunctionWithException<DataInputStream,V,IOException> valueDeserializer) |
static <T> Set<T> |
PulsarSerdeUtils.deserializeSet(DataInputStream in,
FunctionWithException<DataInputStream,T,IOException> deserializer) |
Modifier and Type | Class and Description |
---|---|
class |
RefCountedTmpFileCreator
A utility class that creates local
reference counted files that
serve as temporary files. |
Modifier and Type | Method and Description |
---|---|
static RefCountedBufferingFileStream |
RefCountedBufferingFileStream.openNew(FunctionWithException<File,RefCountedFileWithStream,IOException> tmpFileProvider) |
static RefCountedBufferingFileStream |
RefCountedBufferingFileStream.restore(FunctionWithException<File,RefCountedFileWithStream,IOException> tmpFileProvider,
File initialTmpFile) |
Constructor and Description |
---|
OSSRecoverableFsDataOutputStream(long ossUploadPartSize,
FunctionWithException<File,RefCountedFileWithStream,IOException> cachedFileCreator,
OSSRecoverableMultipartUpload upload,
long sizeBeforeCurrentPart) |
OSSRecoverableWriter(OSSAccessor ossAccessor,
long ossUploadPartSize,
int streamConcurrentUploads,
Executor executor,
FunctionWithException<File,RefCountedFileWithStream,IOException> cachedFileCreator) |
Modifier and Type | Method and Description |
---|---|
static S3RecoverableFsDataOutputStream |
S3RecoverableFsDataOutputStream.newStream(org.apache.flink.fs.s3.common.writer.RecoverableMultiPartUpload upload,
FunctionWithException<File,RefCountedFileWithStream,IOException> tmpFileCreator,
long userDefinedMinPartSize) |
static S3RecoverableFsDataOutputStream |
S3RecoverableFsDataOutputStream.recoverStream(org.apache.flink.fs.s3.common.writer.RecoverableMultiPartUpload upload,
FunctionWithException<File,RefCountedFileWithStream,IOException> tmpFileCreator,
long userDefinedMinPartSize,
long bytesBeforeCurrentPart) |
static S3RecoverableWriter |
S3RecoverableWriter.writer(org.apache.hadoop.fs.FileSystem fs,
FunctionWithException<File,RefCountedFileWithStream,IOException> tempFileCreator,
S3AccessHelper s3AccessHelper,
Executor uploadThreadPool,
long userDefinedMinPartSize,
int maxConcurrentUploadsPerStream) |
Modifier and Type | Method and Description |
---|---|
<V,E extends Exception> |
DeterminismEnvelope.map(FunctionWithException<? super T,? extends V,E> mapper) |
Modifier and Type | Interface and Description |
---|---|
static interface |
ChangelogBackendRestoreOperation.BaseBackendBuilder<K>
Builds base backend for
ChangelogKeyedStateBackend from state. |
Constructor and Description |
---|
BackendRestorerProcedure(FunctionWithException<Collection<S>,T,Exception> instanceSupplier,
CloseableRegistry backendCloseableRegistry,
String logDescription)
Creates a new backend restorer using the given backend supplier and the closeable registry.
|
SourceOperator(FunctionWithException<SourceReaderContext,SourceReader<OUT,SplitT>,Exception> readerFactory,
OperatorEventGateway operatorEventGateway,
SimpleVersionedSerializer<SplitT> splitSerializer,
WatermarkStrategy<OUT> watermarkStrategy,
ProcessingTimeService timeService,
Configuration configuration,
String localHostname,
boolean emitProgressiveWatermarks) |
Modifier and Type | Method and Description |
---|---|
protected void |
SingleCheckpointBarrierHandler.markCheckpointAlignedAndTransformState(InputChannelInfo alignedChannel,
CheckpointBarrier barrier,
FunctionWithException<org.apache.flink.streaming.runtime.io.checkpointing.BarrierHandlerState,org.apache.flink.streaming.runtime.io.checkpointing.BarrierHandlerState,Exception> stateTransformer) |
Modifier and Type | Method and Description |
---|---|
static <A,B> java.util.function.Function<A,B> |
FunctionUtils.uncheckedFunction(FunctionWithException<A,B,?> functionWithException)
Convert at
FunctionWithException into a Function . |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.