Package | Description |
---|---|
org.apache.flink.streaming.connectors.kinesis.internals |
Modifier and Type | Method and Description |
---|---|
List<KinesisStreamShardState> |
KinesisDataFetcher.getSubscribedShardsState() |
Modifier and Type | Method and Description |
---|---|
int |
KinesisDataFetcher.registerNewSubscribedShardState(KinesisStreamShardState newSubscribedShardState)
Register a new subscribed shard state.
|
Constructor and Description |
---|
KinesisDataFetcher(List<String> streams,
SourceFunction.SourceContext<T> sourceContext,
Object checkpointLock,
RuntimeContext runtimeContext,
Properties configProps,
KinesisDeserializationSchema<T> deserializationSchema,
KinesisShardAssigner shardAssigner,
AssignerWithPeriodicWatermarks<T> periodicWatermarkAssigner,
WatermarkTracker watermarkTracker,
AtomicReference<Throwable> error,
List<KinesisStreamShardState> subscribedShardsState,
HashMap<String,String> subscribedStreamsToLastDiscoveredShardIds,
KinesisDataFetcher.FlinkKinesisProxyFactory kinesisProxyFactory,
KinesisDataFetcher.FlinkKinesisProxyV2Factory kinesisProxyV2Factory) |
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.