Class KinesisReadSchemaTransformProvider.KinesisReadSchemaTransformConfiguration
java.lang.Object
org.apache.beam.sdk.io.aws2.kinesis.KinesisReadSchemaTransformProvider.KinesisReadSchemaTransformConfiguration
- All Implemented Interfaces:
Serializable
- Enclosing class:
KinesisReadSchemaTransformProvider
@DefaultSchema(AutoValueSchema.class)
public abstract static class KinesisReadSchemaTransformProvider.KinesisReadSchemaTransformConfiguration
extends Object
implements Serializable
- See Also:
-
Constructor Details
-
KinesisReadSchemaTransformConfiguration
public KinesisReadSchemaTransformConfiguration()
-
-
Method Details
-
builder
public static KinesisReadSchemaTransformProvider.KinesisReadSchemaTransformConfiguration.Builder builder() -
getStreamName
-
getAwsAccessKey
-
getAwsSecretKey
-
getRegion
-
getServiceEndpoint
@SchemaFieldDescription("Optional custom Kinesis service endpoint URI.") public abstract @Nullable String getServiceEndpoint() -
getVerifyCertificate
@SchemaFieldDescription("Whether to verify TLS certificates. Defaults to true. Never disable in production.") public abstract @Nullable Boolean getVerifyCertificate() -
getMaxNumRecords
@SchemaFieldDescription("Maximum number of records to read. When set, the resulting PCollection is bounded.") public abstract @Nullable Long getMaxNumRecords() -
getMaxReadTime
@SchemaFieldDescription("Maximum read time in milliseconds. When set, the resulting PCollection is bounded.") public abstract @Nullable Long getMaxReadTime() -
getInitialPositionInStream
@SchemaFieldDescription("Where to start reading in the stream: LATEST, TRIM_HORIZON, or AT_TIMESTAMP.") public abstract @Nullable String getInitialPositionInStream() -
getInitialTimestampInStream
@SchemaFieldDescription("Epoch millis timestamp used when initial_position_in_stream is AT_TIMESTAMP.") public abstract @Nullable Long getInitialTimestampInStream() -
getRequestRecordsLimit
@SchemaFieldDescription("Max records returned by a single GetRecords call (1-10000).") public abstract @Nullable Long getRequestRecordsLimit() -
getUpToDateThreshold
@SchemaFieldDescription("Threshold duration in milliseconds after which a shard is considered up to date.") public abstract @Nullable Long getUpToDateThreshold() -
getMaxCapacityPerShard
@SchemaFieldDescription("Maximum number of records to hold in memory per shard.") public abstract @Nullable Long getMaxCapacityPerShard() -
getWatermarkPolicy
@SchemaFieldDescription("Watermark policy: ARRIVAL_TIME or PROCESSING_TIME.") public abstract @Nullable String getWatermarkPolicy() -
getWatermarkIdleDurationThreshold
@SchemaFieldDescription("Idle duration threshold in milliseconds for ARRIVAL_TIME watermark policy.") public abstract @Nullable Long getWatermarkIdleDurationThreshold() -
getRateLimit
@SchemaFieldDescription("Fixed delay between GetRecords calls in milliseconds.") public abstract @Nullable Long getRateLimit()
-