Serialized Form
-
Package io.github.flink.gcp.connector.base.failure
-
Package io.github.flink.gcp.connector.base.lineage
-
Class io.github.flink.gcp.connector.base.lineage.PhysicalResourceFacet
class PhysicalResourceFacet extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
resources
List<ResourceIdentifier> resources
-
-
Class io.github.flink.gcp.connector.base.lineage.ResourceIdentifier
class ResourceIdentifier extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.base.rpc
-
Class io.github.flink.gcp.connector.base.rpc.EmulatorEndpoint
class EmulatorEndpoint extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
host
String host
-
port
int port
-
-
-
Package io.github.flink.gcp.connector.base.source
-
Class io.github.flink.gcp.connector.base.source.StartPosition
class StartPosition extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigquery.sink
-
Class io.github.flink.gcp.connector.bigquery.sink.BigQuerySinkConfig
class BigQuerySinkConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
cdcOptions
CdcOptions<? super T> cdcOptions
-
cdcTableOptionsProvider
CdcTableOptionsProvider cdcTableOptionsProvider
-
cdcTableReconciliationPolicy
CdcTableReconciliationPolicy cdcTableReconciliationPolicy
-
createDisposition
CreateDisposition createDisposition
-
destinationResolver
DestinationResolver<? super T> destinationResolver
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
emulatorRestEndpoint
EmulatorEndpoint emulatorRestEndpoint
-
failureHandler
FailureHandler<? super BigQueryFailure> failureHandler
-
location
String location
-
manageCdcTableCreation
boolean manageCdcTableCreation
-
rowAugmentingSerializer
ProtoRowAugmentingSerializer<T> rowAugmentingSerializer
-
schemaUpdateOptions
SchemaUpdateOptions schemaUpdateOptions
-
serializer
BigQueryProtoSerializationSchema<? super T> serializer
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
tableCreateOptionsProvider
TableCreateOptionsProvider tableCreateOptionsProvider
-
-
Class io.github.flink.gcp.connector.bigquery.sink.CdcTableOptions
class CdcTableOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigquery.sink.FixedDestinationResolver
class FixedDestinationResolver extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destination
TableDestination destination
-
-
Class io.github.flink.gcp.connector.bigquery.sink.SchemaUpdateOptions
class SchemaUpdateOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
allowFieldRelaxation
boolean allowFieldRelaxation
-
allowNewFields
boolean allowNewFields
-
-
Class io.github.flink.gcp.connector.bigquery.sink.TableCreateOptions
class TableCreateOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
clusteredFields
List<String> clusteredFields
-
timePartitioningExpirationMs
Long timePartitioningExpirationMs
-
timePartitioningField
String timePartitioningField
-
timePartitioningType
TableCreateOptions.TimePartitioningType timePartitioningType
-
-
Class io.github.flink.gcp.connector.bigquery.sink.TableDestination
class TableDestination extends DestinationResolution implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigquery.sink.cdc
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.CdcOptions
class CdcOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
changeTypeProvider
CdcChangeTypeProvider<? super T> changeTypeProvider
-
sequenceNumberProvider
CdcSequenceNumberProvider<? super T> sequenceNumberProvider
-
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.DebeziumMySqlCdcSequenceNumberEncoder
class DebeziumMySqlCdcSequenceNumberEncoder extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.DebeziumMySqlCdcSequenceNumberProvider
class DebeziumMySqlCdcSequenceNumberProvider extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
encoder
DebeziumMySqlCdcSequenceNumberEncoder encoder
-
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.DebeziumPostgreSqlCdcSequenceNumberProvider
class DebeziumPostgreSqlCdcSequenceNumberProvider extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.DebeziumSpannerCdcSequenceNumberProvider
class DebeziumSpannerCdcSequenceNumberProvider extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.TiCdcSequenceNumberEncoder
class TiCdcSequenceNumberEncoder extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
clusterId
String clusterId
-
-
Class io.github.flink.gcp.connector.bigquery.sink.cdc.TiCdcSequenceNumberProvider
class TiCdcSequenceNumberProvider extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
encoder
TiCdcSequenceNumberEncoder encoder
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.fileloads
-
Class io.github.flink.gcp.connector.bigquery.sink.fileloads.BigQueryFileLoadsSink
class BigQueryFileLoadsSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigQuerySinkConfig<T> config
-
lineageTableName
String lineageTableName
-
options
FileLoadsOptions options
-
storage
StagingStorage storage
-
-
Class io.github.flink.gcp.connector.bigquery.sink.fileloads.FileLoadsOptions
class FileLoadsOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destinationIdleTimeout
Duration destinationIdleTimeout
-
loadJobPollInitialBackoff
Duration loadJobPollInitialBackoff
-
loadJobPollMaxBackoff
Duration loadJobPollMaxBackoff
-
maxConcurrentCheckpointFinalizations
int maxConcurrentCheckpointFinalizations
-
maxConcurrentDestinations
int maxConcurrentDestinations
-
maxOpenDestinations
int maxOpenDestinations
-
maxPendingFiles
int maxPendingFiles
-
maxSerializedRowBytes
long maxSerializedRowBytes
-
maxStagingFileBytes
long maxStagingFileBytes
-
minCheckpointInterval
Duration minCheckpointInterval
-
parquetCompression
ParquetCompression parquetCompression
-
perDestinationMetrics
boolean perDestinationMetrics
-
schemaReconcileInitialBackoff
Duration schemaReconcileInitialBackoff
-
schemaReconcileMaxAttempts
int schemaReconcileMaxAttempts
-
schemaReconcileMaxBackoff
Duration schemaReconcileMaxBackoff
-
stagingFormat
StagingFormat stagingFormat
-
stagingPath
String stagingPath
-
tempDataset
String tempDataset
-
writeDisposition
WriteDisposition writeDisposition
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.fileloads.committer
-
Class io.github.flink.gcp.connector.bigquery.sink.fileloads.committer.FileLoadsCheckpointStamper
class FileLoadsCheckpointStamper extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigquery.sink.fileloads.writer
-
Class io.github.flink.gcp.connector.bigquery.sink.fileloads.writer.GcsStagingStorage
class GcsStagingStorage extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.serializer
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.AdditionalField
class AdditionalField extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
name
String name
-
nullPolicy
AdditionalFieldNullPolicy nullPolicy
-
type
AdditionalFieldType type
-
valueProvider
AdditionalFieldValueProvider<? super T> valueProvider
-
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.AdditionalFields
class AdditionalFields extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
fields
List<AdditionalField<? super T>> fields
-
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.BigQueryProtoSerializationSchema
class BigQueryProtoSerializationSchema extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.LazyDerivedState
class LazyDerivedState extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.ProtoRowAugmentationField
class ProtoRowAugmentationField extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
descriptorField
com.google.protobuf.DescriptorProtos.FieldDescriptorProto descriptorField
-
nullPolicy
AdditionalFieldNullPolicy nullPolicy
-
nullValueMessage
String nullValueMessage
-
providerFailureMessage
String providerFailureMessage
-
schemaOwnership
io.github.flink.gcp.connector.bigquery.sink.serializer.ProtoRowAugmentationField.SchemaOwnership schemaOwnership
-
tableField
com.google.cloud.bigquery.storage.v1.TableFieldSchema tableField
-
valueProvider
ProtoRowAugmentationField.ValueProvider<? super T> valueProvider
-
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.ProtoRowAugmentingSerializer
class ProtoRowAugmentingSerializer extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
delegate
BigQueryProtoSerializationSchema<? super T> delegate
-
descriptorFieldDescription
String descriptorFieldDescription
-
fields
List<ProtoRowAugmentationField<? super T>> fields
-
rowFailureMessage
String rowFailureMessage
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.serializer.avro
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.avro.AvroRecordSerializationSchema
class AvroRecordSerializationSchema extends BigQueryProtoSerializationSchema<org.apache.avro.generic.IndexedRecord> implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
avroSchemaJson
String avroSchemaJson
-
conversionState
LazyDerivedState<io.github.flink.gcp.connector.bigquery.sink.serializer.avro.AvroRecordSerializationSchema.ConversionState> conversionState
-
options
AvroSchemaOptions options
-
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.avro.AvroSchemaOptions
class AvroSchemaOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigquery.sink.serializer.json
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.json.JsonDocumentOptions
class JsonDocumentOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
ignoreUnknownFields
boolean ignoreUnknownFields
-
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.json.JsonDocumentSerializationSchema
class JsonDocumentSerializationSchema extends BigQueryProtoSerializationSchema<String> implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
options
JsonDocumentOptions options
-
rowDescriptor
LazyDerivedState<com.google.protobuf.Descriptors.Descriptor> rowDescriptor
-
tableSchema
com.google.cloud.bigquery.storage.v1.TableSchema tableSchema
The destination schema. Protobuf messages are Java-serializable, so this travels in the job graph as it is; the descriptor derived from it is not, and is rebuilt on the task manager.
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.serializer.proto
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.proto.ProtoMessageSerializationSchema
class ProtoMessageSerializationSchema extends BigQueryProtoSerializationSchema<T extends com.google.protobuf.Message> implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
conversionState
LazyDerivedState<io.github.flink.gcp.connector.bigquery.sink.serializer.proto.ProtoMessageSerializationSchema.ConversionState> conversionState
-
messageClass
Class<T extends com.google.protobuf.Message> messageClass
-
options
ProtoSchemaOptions options
-
-
Class io.github.flink.gcp.connector.bigquery.sink.serializer.proto.ProtoSchemaOptions
class ProtoSchemaOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
deriveRequiredColumns
boolean deriveRequiredColumns
-
geographyFieldOptions
Map<Integer,
String> geographyFieldOptions AsjsonFieldOptions, forGEOGRAPHYcolumns. -
geographyFieldPaths
Set<String> geographyFieldPaths
-
jsonFieldOptions
Map<Integer,
String> jsonFieldOptions Configured field options, keyed by extension number, valued by the option's full name ornullwhen it was configured by number alone. Keyed by number because two entries for one number would be contradictory — and because an unnamed entry sitting beside a named one would match anything at that number, defeating the name check the named entry exists for. -
jsonFieldPaths
Set<String> jsonFieldPaths
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.storage
-
Class io.github.flink.gcp.connector.bigquery.sink.storage.BigQueryBufferedStreamSink
class BigQueryBufferedStreamSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigQuerySinkConfig<T> config
-
lineageTableName
String lineageTableName
-
options
BufferedStreamOptions options
-
serviceFactory
BufferedStreamServiceFactory serviceFactory
-
-
Class io.github.flink.gcp.connector.bigquery.sink.storage.BigQueryDefaultStreamSink
class BigQueryDefaultStreamSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigQuerySinkConfig<T> config
-
lineageTableName
String lineageTableName
-
options
DefaultStreamOptions options
-
-
Class io.github.flink.gcp.connector.bigquery.sink.storage.BufferedStreamOptions
class BufferedStreamOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destinationIdleTimeout
Duration destinationIdleTimeout
-
maxAppendRequestBytes
long maxAppendRequestBytes
-
maxRetryDuration
Duration maxRetryDuration
-
recoveryInitialBackoff
Duration recoveryInitialBackoff
-
recoveryMaxAttempts
int recoveryMaxAttempts
-
recoveryMaxBackoff
Duration recoveryMaxBackoff
-
retryDelayMultiplier
double retryDelayMultiplier
-
retryInitialDelay
Duration retryInitialDelay
-
retryMaxAttempts
int retryMaxAttempts
-
retryMaxDelay
Duration retryMaxDelay
-
-
Class io.github.flink.gcp.connector.bigquery.sink.storage.DefaultStreamOptions
class DefaultStreamOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destinationIdleTimeout
Duration destinationIdleTimeout
-
flushInterval
Duration flushInterval
-
maxAppendRequestBytes
long maxAppendRequestBytes
-
maxConnectionsPerRegion
int maxConnectionsPerRegion
-
maxInflightBytes
long maxInflightBytes
-
maxInflightRequests
int maxInflightRequests
-
maxRetryDuration
Duration maxRetryDuration
-
minConnectionsPerRegion
int minConnectionsPerRegion
-
perDestinationMetrics
boolean perDestinationMetrics
-
recoveryInitialBackoff
Duration recoveryInitialBackoff
-
recoveryMaxAttempts
int recoveryMaxAttempts
-
recoveryMaxBackoff
Duration recoveryMaxBackoff
-
retryDelayMultiplier
double retryDelayMultiplier
-
retryInitialDelay
Duration retryInitialDelay
-
retryMaxAttempts
int retryMaxAttempts
-
retryMaxDelay
Duration retryMaxDelay
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.storage.writer
-
Class io.github.flink.gcp.connector.bigquery.sink.storage.writer.StreamWriterRowAppenderFactory
class StreamWriterRowAppenderFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
options
DefaultStreamOptions options
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
Class io.github.flink.gcp.connector.bigquery.sink.storage.writer.WriteClientBufferedStreamServiceFactory
class WriteClientBufferedStreamServiceFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.bigquery.sink.tables
-
Exception io.github.flink.gcp.connector.bigquery.sink.tables.RetriableTableAdminException
class RetriableTableAdminException extends TableAdminException implements Serializable- serialVersionUID:
- 1L
-
Exception io.github.flink.gcp.connector.bigquery.sink.tables.SchemaUnifier.SchemaUnionException
class SchemaUnionException extends IOException implements Serializable- serialVersionUID:
- 1L
-
Exception io.github.flink.gcp.connector.bigquery.sink.tables.TableAdminException
class TableAdminException extends IOException implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
creationRequested
boolean creationRequested
-
-
-
Package io.github.flink.gcp.connector.bigquery.source
-
Class io.github.flink.gcp.connector.bigquery.source.BigQuerySourceConfig
class BigQuerySourceConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
deserializer
BigQueryRowDeserializationSchema<T> deserializer
-
materializeViews
boolean materializeViews
-
maxBytesPerFetch
long maxBytesPerFetch
-
maxRecordsPerFetch
int maxRecordsPerFetch
-
maxStreamCount
int maxStreamCount
-
parentProject
String parentProject
-
preferredMinStreamCount
int preferredMinStreamCount
-
query
String query
-
queryLocation
String queryLocation
-
queryResultDataset
String queryResultDataset
-
queryRunner
QueryRunner queryRunner
-
reuseQueryResultWithin
Duration reuseQueryResultWithin
-
rowRestriction
String rowRestriction
-
rowStreamOpener
RowStreamOpener rowStreamOpener
-
selectedFields
List<String> selectedFields
-
sessionCreatorFactory
ReadSessionCreatorFactory sessionCreatorFactory
-
snapshotTime
Instant snapshotTime
-
table
TableDestination table
-
-
Class io.github.flink.gcp.connector.bigquery.source.BigQueryStorageReadSource
class BigQueryStorageReadSource extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigQuerySourceConfig<T> config
-
lineageTableName
String lineageTableName
-
-
-
Package io.github.flink.gcp.connector.bigquery.source.enumerator
-
Class io.github.flink.gcp.connector.bigquery.source.enumerator.DefaultReadSessionCreatorFactory
class DefaultReadSessionCreatorFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.bigquery.source.query
-
Class io.github.flink.gcp.connector.bigquery.source.query.BigQueryQueryRunner
class BigQueryQueryRunner extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.bigquery.source.reader
-
Class io.github.flink.gcp.connector.bigquery.source.reader.ReadClientRowStreamOpener
class ReadClientRowStreamOpener extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
retryMaxAttempts
int retryMaxAttempts
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.bigquery.source.serializer
-
Package io.github.flink.gcp.connector.bigquery.table.sink
-
Class io.github.flink.gcp.connector.bigquery.table.sink.RowDataSchemaOptions
class RowDataSchemaOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigtable
-
Class io.github.flink.gcp.connector.bigtable.LazyBigtableDataClient
class LazyBigtableDataClient extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
owner
String owner
-
-
Class io.github.flink.gcp.connector.bigtable.TableDestination
class TableDestination extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigtable.sink
-
Class io.github.flink.gcp.connector.bigtable.sink.BigtableMutateRowsSink
class BigtableMutateRowsSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigtableSinkConfig<T> config
-
expectedFamilies
TableCreateOptions expectedFamilies
-
initialDestination
TableDestination initialDestination
-
lineageTableName
String lineageTableName
-
-
Class io.github.flink.gcp.connector.bigtable.sink.BigtableSinkConfig
class BigtableSinkConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
createDisposition
CreateDisposition createDisposition
-
destinationResolver
DestinationResolver<? super T> destinationResolver
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
failedMutationHandler
FailureHandler<? super FailedMutation> failedMutationHandler
-
serializer
BigtableSerializationSchema<? super T> serializer
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
tableCreateOptions
TableCreateOptions tableCreateOptions
-
writerOptions
BigtableWriterOptions writerOptions
-
-
Class io.github.flink.gcp.connector.bigtable.sink.BigtableStagedOptions
class BigtableStagedOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
markerFamily
String markerFamily
-
maxStagedBytes
long maxStagedBytes
-
maxStagedEntries
int maxStagedEntries
-
requestOptions
BigtableRequestOptions requestOptions
-
-
Class io.github.flink.gcp.connector.bigtable.sink.BigtableWriterOptions
class BigtableWriterOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
batchElementCountThreshold
Long batchElementCountThreshold
-
batchRequestByteThreshold
Long batchRequestByteThreshold
-
destinationIdleTimeout
Duration destinationIdleTimeout
-
maxActiveInstances
int maxActiveInstances
-
maxConsecutiveRejections
int maxConsecutiveRejections
-
maxInFlightBytes
long maxInFlightBytes
-
maxInFlightEntries
int maxInFlightEntries
-
perDestinationMetrics
boolean perDestinationMetrics
-
recoveryInitialBackoff
Duration recoveryInitialBackoff
-
recoveryMaxAttempts
int recoveryMaxAttempts
-
recoveryMaxBackoff
Duration recoveryMaxBackoff
-
-
Class io.github.flink.gcp.connector.bigtable.sink.FixedDestinationResolver
class FixedDestinationResolver extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destination
TableDestination destination
-
-
Class io.github.flink.gcp.connector.bigtable.sink.GcRule
class GcRule extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
kind
GcRule.Kind kind
-
maxAge
Duration maxAge
-
maxVersions
Integer maxVersions
-
rules
List<GcRule> rules
-
-
Class io.github.flink.gcp.connector.bigtable.sink.TableCreateOptions
class TableCreateOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
columnFamilies
LinkedHashMap<String,
GcRule> columnFamilies Family name to garbage-collection rule; anullvalue is a family without one. -
columnFamilyTypes
Map<String,
ColumnFamilyType> columnFamilyTypes
-
-
-
Package io.github.flink.gcp.connector.bigtable.sink.conditional
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.AggregateValue
class AggregateValue extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
bytes
com.google.protobuf.ByteString bytes
-
integer
long integer
-
kind
io.github.flink.gcp.connector.bigtable.sink.conditional.AggregateValue.Kind kind
-
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.BigtableConditionalSink
class BigtableConditionalSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
SingleRowRequestConfig<T> config
-
lineageTableName
String lineageTableName
-
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalFilter
class ConditionalFilter extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
bytes
com.google.protobuf.ByteString bytes
-
children
List<ConditionalFilter> children
-
count
int count
-
end
long end
-
family
String family
-
kind
io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalFilter.Kind kind
-
start
long start
-
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalMutation
class ConditionalMutation extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
aggregate
AggregateValue aggregate
-
end
Long end
-
family
String family
-
kind
io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalMutation.Kind kind
-
qualifier
com.google.protobuf.ByteString qualifier
-
start
Long start
-
timestamp
long timestamp
-
value
com.google.protobuf.ByteString value
-
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalRequest
class ConditionalRequest extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
otherwiseMutations
List<ConditionalMutation> otherwiseMutations
-
predicate
ConditionalFilter predicate
-
rowKey
com.google.protobuf.ByteString rowKey
-
thenMutations
List<ConditionalMutation> thenMutations
-
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalResult
class ConditionalResult extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destination
TableDestination destination
-
predicateMatched
boolean predicateMatched
-
rowKey
com.google.protobuf.ByteString rowKey
-
selectedBranchHasMutations
boolean selectedBranchHasMutations
-
-
Class io.github.flink.gcp.connector.bigtable.sink.conditional.ConditionalResultSerializer
class ConditionalResultSerializer extends org.apache.flink.api.common.typeutils.TypeSerializer<ConditionalResult> implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigtable.sink.mutaterows.writer
-
Class io.github.flink.gcp.connector.bigtable.sink.mutaterows.writer.DefaultMutationBatcherFactory
class DefaultMutationBatcherFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
credentialsOverride
com.google.api.gax.core.CredentialsProvider credentialsOverride
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
writerOptions
BigtableWriterOptions writerOptions
-
-
-
Package io.github.flink.gcp.connector.bigtable.sink.readmodifywrite
-
Class io.github.flink.gcp.connector.bigtable.sink.readmodifywrite.BigtableReadModifyWriteSink
class BigtableReadModifyWriteSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
SingleRowRequestConfig<T> config
-
lineageTableName
String lineageTableName
-
-
Class io.github.flink.gcp.connector.bigtable.sink.readmodifywrite.ReadModifyWriteRequest
class ReadModifyWriteRequest extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
rowKey
com.google.protobuf.ByteString rowKey
-
rules
List<ReadModifyWriteRule> rules
-
-
Class io.github.flink.gcp.connector.bigtable.sink.readmodifywrite.ReadModifyWriteResult
class ReadModifyWriteResult extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destination
TableDestination destination
-
row
BigtableRow row
-
-
Class io.github.flink.gcp.connector.bigtable.sink.readmodifywrite.ReadModifyWriteResultSerializer
class ReadModifyWriteResultSerializer extends org.apache.flink.api.common.typeutils.TypeSerializer<ReadModifyWriteResult> implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
rowSerializer
org.apache.flink.api.common.typeutils.TypeSerializer<BigtableRow> rowSerializer
-
-
Class io.github.flink.gcp.connector.bigtable.sink.readmodifywrite.ReadModifyWriteRule
class ReadModifyWriteRule extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appendValue
com.google.protobuf.ByteString appendValue
-
family
String family
-
incrementAmount
long incrementAmount
-
qualifier
com.google.protobuf.ByteString qualifier
-
-
-
Package io.github.flink.gcp.connector.bigtable.sink.serializer
-
Package io.github.flink.gcp.connector.bigtable.sink.singlerow
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.BigtableRequestOptions
class BigtableRequestOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.BigtableRow
class BigtableRow extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
cells
List<BigtableRow.Cell> cells
-
key
com.google.protobuf.ByteString key
-
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.BigtableRow.Cell
class Cell extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.BigtableRowSerializer
class BigtableRowSerializer extends org.apache.flink.api.common.typeutils.TypeSerializer<BigtableRow> implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.BigtableStagedSink
class BigtableStagedSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigtableSinkConfig<T> config
-
expectedFamilies
Map<String,
ColumnFamilyType> expectedFamilies -
logicalTableName
String logicalTableName
-
options
BigtableStagedOptions options
-
-
-
Package io.github.flink.gcp.connector.bigtable.sink.singlerow.writer
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.writer.BigtableRequestFunction
class BigtableRequestFunction extends org.apache.flink.streaming.api.functions.async.RichAsyncFunction<IN,OUT> implements Serializable - serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
injectedFactory
SingleRowClientFactory injectedFactory
An injected factory, ornullfor the production one built inBigtableRequestFunction.open(org.apache.flink.api.common.functions.OpenContext). -
options
BigtableRequestOptions options
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.writer.DefaultSingleRowClientFactory
class DefaultSingleRowClientFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
credentialsOverride
com.google.api.gax.core.CredentialsProvider credentialsOverride
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
options
BigtableRequestOptions options
-
-
Class io.github.flink.gcp.connector.bigtable.sink.singlerow.writer.SingleRowRequestConfig
class SingleRowRequestConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
destinationResolver
DestinationResolver<? super T> destinationResolver
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
failedRequestHandler
FailureHandler<? super FailedRequest> failedRequestHandler
-
requestOptions
BigtableRequestOptions requestOptions
-
serializer
RowRequestSerializer<? super T> serializer
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.bigtable.source
-
Class io.github.flink.gcp.connector.bigtable.source.BigtableChangeStreamSource
class BigtableChangeStreamSource extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigtableChangeStreamSourceConfig<T> config
-
lineageTableName
String lineageTableName
-
-
Class io.github.flink.gcp.connector.bigtable.source.BigtableChangeStreamSourceConfig
class BigtableChangeStreamSourceConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
boundedTimestamp
Instant boundedTimestamp
-
coordinatorClientFactory
ChangeStreamCoordinatorClientFactory coordinatorClientFactory
-
deserializer
BigtableChangeStreamDeserializationSchema<T> deserializer
-
maxConcurrentStreamsPerSubtask
int maxConcurrentStreamsPerSubtask
-
mutationFilter
BigtableChangeStreamMutationFilter mutationFilter
-
opener
ChangeStreamOpener opener
-
restoreResolver
ChangeStreamRestoreResolver restoreResolver
-
resumeFallback
StartPosition resumeFallback
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
startPosition
StartPosition startPosition
-
table
TableDestination table
-
-
Class io.github.flink.gcp.connector.bigtable.source.BigtableSourceConfig
class BigtableSourceConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
deserializer
BigtableRowDeserializationSchema<T> deserializer
-
filter
com.google.cloud.bigtable.data.v2.models.Filters.Filter filter
-
maxBytesPerFetch
long maxBytesPerFetch
-
maxRowsPerFetch
int maxRowsPerFetch
-
opener
RowStreamOpener opener
-
ranges
List<com.google.cloud.bigtable.data.v2.models.Range.ByteStringRange> ranges
-
samplerFactory
RowKeySamplerFactory samplerFactory
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
table
TableDestination table
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.changestream
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation
class BigtableChangeStreamMutation extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
commitTime
Instant commitTime
-
entries
List<BigtableChangeStreamMutation.Entry> entries
-
estimatedLowWatermarkTime
Instant estimatedLowWatermarkTime
-
rowKey
com.google.protobuf.ByteString rowKey
-
sourceClusterId
String sourceClusterId
-
tieBreaker
int tieBreaker
-
token
String token
-
type
BigtableChangeStreamMutation.MutationType type
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.AddToCellEntry
class AddToCellEntry extends BigtableChangeStreamMutation.Entry implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
familyName
String familyName
-
input
BigtableChangeStreamMutation.Value input
-
qualifier
BigtableChangeStreamMutation.Value qualifier
-
timestamp
BigtableChangeStreamMutation.Value timestamp
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.DeleteCellsEntry
class DeleteCellsEntry extends BigtableChangeStreamMutation.Entry implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
familyName
String familyName
-
qualifier
com.google.protobuf.ByteString qualifier
-
timestampRange
BigtableChangeStreamMutation.TimestampRange timestampRange
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.DeleteFamilyEntry
class DeleteFamilyEntry extends BigtableChangeStreamMutation.Entry implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
familyName
String familyName
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.Entry
class Entry extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.Int64Value
class Int64Value extends BigtableChangeStreamMutation.Value implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
value
long value
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.MergeToCellEntry
class MergeToCellEntry extends BigtableChangeStreamMutation.Entry implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
familyName
String familyName
-
input
BigtableChangeStreamMutation.Value input
-
qualifier
BigtableChangeStreamMutation.Value qualifier
-
timestamp
BigtableChangeStreamMutation.Value timestamp
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.RawTimestamp
class RawTimestamp extends BigtableChangeStreamMutation.Value implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
value
long value
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.RawValue
class RawValue extends BigtableChangeStreamMutation.Value implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
value
com.google.protobuf.ByteString value
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.SetCellEntry
class SetCellEntry extends BigtableChangeStreamMutation.Entry implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
familyName
String familyName
-
qualifier
com.google.protobuf.ByteString qualifier
-
timestampMicros
long timestampMicros
-
value
com.google.protobuf.ByteString value
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.TimestampBound
class TimestampBound extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
timestampMicros
long timestampMicros
-
type
BigtableChangeStreamMutation.BoundType type
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.TimestampRange
class TimestampRange extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutation.Value
class Value extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutationFilter
class BigtableChangeStreamMutationFilter extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.BigtableChangeStreamMutationSerializer
class BigtableChangeStreamMutationSerializer extends org.apache.flink.api.common.typeutils.TypeSerializer<BigtableChangeStreamMutation> implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.PartitionProgressEvent
class PartitionProgressEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.PartitionTransitionEvent
class PartitionTransitionEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
finishedSplitId
String finishedSplitId
-
lowWatermark
Instant lowWatermark
-
successors
List<PartitionTransitionEvent.Successor> successors
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.PartitionTransitionEvent.Successor
class Successor extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
continuationToken
com.google.cloud.bigtable.data.v2.models.ChangeStreamContinuationToken continuationToken
-
partition
com.google.cloud.bigtable.data.v2.models.Range.ByteStringRange partition
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.ReaderCapacityEvent
class ReaderCapacityEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
freeSlots
int freeSlots
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.changestream.enumerator
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.enumerator.DefaultChangeStreamCoordinatorClientFactory
class DefaultChangeStreamCoordinatorClientFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
table
TableDestination table
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.changestream.reader
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.reader.DataClientChangeStreamOpener
class DataClientChangeStreamOpener extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
client
LazyBigtableDataClient client
The client this opener streams through, and everything around building it. Unlike the scan source's two seams, every call here runs on the reader's task thread — the change-stream reader has no split-fetcher pool — so the holder's thread guarding is inherited rather than demanded by this seam.
-
-
Class io.github.flink.gcp.connector.bigtable.source.changestream.reader.DefaultChangeStreamRestoreResolver
class DefaultChangeStreamRestoreResolver extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
table
TableDestination table
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.readrows
-
Class io.github.flink.gcp.connector.bigtable.source.readrows.BigtableScanSource
class BigtableScanSource extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
BigtableSourceConfig<T> config
-
lineageTableName
String lineageTableName
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.readrows.enumerator
-
Class io.github.flink.gcp.connector.bigtable.source.readrows.enumerator.DefaultRowKeySamplerFactory
class DefaultRowKeySamplerFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
appProfileId
String appProfileId
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.readrows.reader
-
Class io.github.flink.gcp.connector.bigtable.source.readrows.reader.DataClientRowStreamOpener
class DataClientRowStreamOpener extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
client
LazyBigtableDataClient client
The client this opener reads through, and everything around building it: a fetcher thread opens a stream through it whileDataClientRowStreamOpener.close()runs on the task thread once the fetchers are down, and a fetcher generation can start while the previous one is still finishing.
-
-
-
Package io.github.flink.gcp.connector.bigtable.source.serializer
-
Class io.github.flink.gcp.connector.bigtable.source.serializer.BigtableChangeStreamMutationDeserializationSchema
class BigtableChangeStreamMutationDeserializationSchema extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigtable.table
-
Class io.github.flink.gcp.connector.bigtable.table.BigtableLookupConfig
class BigtableLookupConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
async
boolean async
-
cacheType
org.apache.flink.table.connector.source.lookup.LookupOptions.LookupCacheType cacheType
-
fullPeriodicReloadInterval
Duration fullPeriodicReloadInterval
-
fullPeriodicScheduleMode
org.apache.flink.table.connector.source.lookup.cache.trigger.PeriodicCacheReloadTrigger.ScheduleMode fullPeriodicScheduleMode
-
fullReloadStrategy
org.apache.flink.table.connector.source.lookup.LookupOptions.ReloadStrategy fullReloadStrategy
-
fullTimedReloadIntervalDays
int fullTimedReloadIntervalDays
-
fullTimedReloadIsoTime
String fullTimedReloadIsoTime
-
maxRetries
int maxRetries
-
partialCacheMissingKey
boolean partialCacheMissingKey
-
partialExpireAfterAccess
Duration partialExpireAfterAccess
-
partialExpireAfterWrite
Duration partialExpireAfterWrite
-
partialMaxRows
Long partialMaxRows
-
-
-
Package io.github.flink.gcp.connector.bigtable.table.function
-
Class io.github.flink.gcp.connector.bigtable.table.function.BigtableCheckAndMutateFunction
class BigtableCheckAndMutateFunction extends io.github.flink.gcp.connector.bigtable.table.function.AbstractBigtableWriteFunction<Boolean> implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.bigtable.table.function.BigtableReadModifyWriteFunction
class BigtableReadModifyWriteFunction extends io.github.flink.gcp.connector.bigtable.table.function.AbstractBigtableWriteFunction<org.apache.flink.types.Row> implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.bigtable.table.source
-
Class io.github.flink.gcp.connector.bigtable.table.source.BigtableRowDataAsyncLookupFunction
class BigtableRowDataAsyncLookupFunction extends org.apache.flink.table.functions.AsyncLookupFunction implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
lookup
io.github.flink.gcp.connector.bigtable.table.source.BigtableRowDataLookup lookup
-
maxRetries
int maxRetries
-
-
Class io.github.flink.gcp.connector.bigtable.table.source.BigtableRowDataLookupFunction
class BigtableRowDataLookupFunction extends org.apache.flink.table.functions.LookupFunction implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
lookup
io.github.flink.gcp.connector.bigtable.table.source.BigtableRowDataLookup lookup
-
maxRetries
int maxRetries
-
-
-
Package io.github.flink.gcp.connector.cloudtasks.sink
-
Class io.github.flink.gcp.connector.cloudtasks.sink.CloudTasksCreateTaskSink
class CloudTasksCreateTaskSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
CloudTasksSinkConfig<T> config
-
logicalTableName
String logicalTableName
-
-
Class io.github.flink.gcp.connector.cloudtasks.sink.CloudTasksSinkConfig
class CloudTasksSinkConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destinationResolver
DestinationResolver<? super T> destinationResolver
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
failedTaskHandler
FailureHandler<? super FailedTask> failedTaskHandler
-
serializer
CloudTasksSerializationSchema<? super T> serializer
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
taskIdExtractor
TaskIdExtractor<? super T> taskIdExtractor
-
writerOptions
CloudTasksWriterOptions writerOptions
-
-
Class io.github.flink.gcp.connector.cloudtasks.sink.CloudTasksStagedCreateTaskSink
class CloudTasksStagedCreateTaskSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
CloudTasksSinkConfig<T> config
-
logicalTableName
String logicalTableName
-
options
CloudTasksStagedOptions options
-
staging
CloudTasksStagingConfig staging
-
-
Class io.github.flink.gcp.connector.cloudtasks.sink.CloudTasksStagedOptions
class CloudTasksStagedOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
clockSkewAllowance
Duration clockSkewAllowance
-
expiredEnvelopePolicy
CloudTasksStagedOptions.ExpiredEnvelopePolicy expiredEnvelopePolicy
-
maxStagedBytes
long maxStagedBytes
-
maxStagedTasks
int maxStagedTasks
-
nameRetention
Duration nameRetention
-
requestTimeout
Duration requestTimeout
-
verifyQueueRetention
boolean verifyQueueRetention
-
-
Class io.github.flink.gcp.connector.cloudtasks.sink.CloudTasksStagingConfig
class CloudTasksStagingConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.cloudtasks.sink.CloudTasksWriterOptions
class CloudTasksWriterOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
channelPoolSize
Integer channelPoolSize
-
maxInFlightTasks
int maxInFlightTasks
-
notFoundRecoveryInitialBackoff
Duration notFoundRecoveryInitialBackoff
-
notFoundRecoveryMaxAttempts
int notFoundRecoveryMaxAttempts
-
notFoundRecoveryMaxBackoff
Duration notFoundRecoveryMaxBackoff
-
perDestinationMetrics
boolean perDestinationMetrics
-
recoveryInitialBackoff
Duration recoveryInitialBackoff
-
recoveryMaxAttempts
int recoveryMaxAttempts
-
recoveryMaxBackoff
Duration recoveryMaxBackoff
-
-
Class io.github.flink.gcp.connector.cloudtasks.sink.FixedDestinationResolver
class FixedDestinationResolver extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destination
QueueDestination destination
-
-
Class io.github.flink.gcp.connector.cloudtasks.sink.QueueDestination
class QueueDestination extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.cloudtasks.sink.serializer
-
Class io.github.flink.gcp.connector.cloudtasks.sink.serializer.AppEngineTargetSerializationSchema
class AppEngineTargetSerializationSchema extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.cloudtasks.sink.serializer.HttpTargetSerializationSchema
class HttpTargetSerializationSchema extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
body
org.apache.flink.api.common.serialization.SerializationSchema<T> body
-
carriesBody
boolean carriesBody
Whether the method allows a request body. Cloud Tasks accepts one only onPOST,PUTandPATCH, and rejects a task that carries one under any other method. -
headersExtractor
HttpTargetSerializationSchema.HeadersExtractor<? super T> headersExtractor
-
method
com.google.cloud.tasks.v2.HttpMethod method
-
oauthScope
String oauthScope
-
oauthServiceAccount
String oauthServiceAccount
-
oidcAudience
String oidcAudience
-
oidcServiceAccount
String oidcServiceAccount
-
url
String url
-
urlExtractor
HttpTargetSerializationSchema.UrlExtractor<? super T> urlExtractor
-
-
-
Package io.github.flink.gcp.connector.cloudtasks.sink.writer
-
Class io.github.flink.gcp.connector.cloudtasks.sink.writer.DefaultTaskCreatorFactory
class DefaultTaskCreatorFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
channelPoolSize
Integer channelPoolSize
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.cloudtasks.table.sink
-
Class io.github.flink.gcp.connector.cloudtasks.table.sink.AppEngineTargetSpec
class AppEngineTargetSpec extends TargetSpec implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.cloudtasks.table.sink.HttpTargetSpec
class HttpTargetSpec extends TargetSpec implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
bodyContentType
String bodyContentType
-
headers
Map<String,
String> headers -
method
com.google.cloud.tasks.v2.HttpMethod method
-
oauthScope
String oauthScope
-
oauthServiceAccountEmail
String oauthServiceAccountEmail
-
oidcAudience
String oidcAudience
-
oidcServiceAccountEmail
String oidcServiceAccountEmail
-
url
String url
-
-
Class io.github.flink.gcp.connector.cloudtasks.table.sink.TargetSpec
class TargetSpec extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.pubsub.deadletter
-
Class io.github.flink.gcp.connector.pubsub.deadletter.PubSubDeadLetterQueue
class PubSubDeadLetterQueue extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
flushTimeout
Duration flushTimeout
-
maxInFlightMessages
int maxInFlightMessages
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
shutdownTimeout
Duration shutdownTimeout
-
topic
TopicDestination topic
-
-
-
Package io.github.flink.gcp.connector.pubsub.sink
-
Class io.github.flink.gcp.connector.pubsub.sink.FixedDestinationResolver
class FixedDestinationResolver extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
destination
TopicDestination destination
-
-
Class io.github.flink.gcp.connector.pubsub.sink.PubSubPublisherOptions
class PubSubPublisherOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
batchDelayThreshold
Duration batchDelayThreshold
-
batchElementCountThreshold
Long batchElementCountThreshold
-
batchRequestByteThreshold
Long batchRequestByteThreshold
-
destinationIdleTimeout
Duration destinationIdleTimeout
-
enableMessageOrdering
boolean enableMessageOrdering
-
maxActivePublishers
int maxActivePublishers
-
maxConsecutiveRejections
int maxConsecutiveRejections
-
maxInFlightBytes
long maxInFlightBytes
-
maxInFlightMessages
int maxInFlightMessages
-
perDestinationMetrics
boolean perDestinationMetrics
-
publishProgressTimeout
Duration publishProgressTimeout
-
recoveryInitialBackoff
Duration recoveryInitialBackoff
-
recoveryMaxAttempts
int recoveryMaxAttempts
-
recoveryMaxBackoff
Duration recoveryMaxBackoff
-
retryDelayMultiplier
Double retryDelayMultiplier
-
retryInitialDelay
Duration retryInitialDelay
-
retryInitialRpcTimeout
Duration retryInitialRpcTimeout
-
retryMaxAttempts
Integer retryMaxAttempts
-
retryMaxDelay
Duration retryMaxDelay
-
retryMaxRpcTimeout
Duration retryMaxRpcTimeout
-
retryRpcTimeoutMultiplier
Double retryRpcTimeoutMultiplier
-
retryTotalTimeout
Duration retryTotalTimeout
-
shutdownTimeout
Duration shutdownTimeout
-
-
Class io.github.flink.gcp.connector.pubsub.sink.PubSubPublisherSink
class PubSubPublisherSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
PubSubSinkConfig<T> config
-
lineageTableName
String lineageTableName
-
-
Class io.github.flink.gcp.connector.pubsub.sink.PubSubSinkConfig
class PubSubSinkConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
createDisposition
CreateDisposition createDisposition
-
destinationResolver
DestinationResolver<? super T> destinationResolver
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
failedMessageHandler
FailureHandler<? super FailedMessage> failedMessageHandler
-
publisherOptions
PubSubPublisherOptions publisherOptions
-
serializer
PubSubSerializationSchema<? super T> serializer
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
topicCreateOptions
TopicCreateOptions topicCreateOptions
-
-
Class io.github.flink.gcp.connector.pubsub.sink.TopicCreateOptions
class TopicCreateOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.pubsub.sink.TopicDestination
class TopicDestination extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.pubsub.sink.serializer
-
Package io.github.flink.gcp.connector.pubsub.sink.writer
-
Class io.github.flink.gcp.connector.pubsub.sink.writer.DefaultPublisherFactory
class DefaultPublisherFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
credentialsOverride
com.google.api.gax.core.CredentialsProvider credentialsOverride
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
options
PubSubPublisherOptions options
-
-
-
Package io.github.flink.gcp.connector.pubsub.source
-
Class io.github.flink.gcp.connector.pubsub.source.PubSubSourceConfig
class PubSubSourceConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
createOptions
Map<SubscriptionDestination,
SubscriptionCreateOptions> createOptions -
deserializationFailurePolicy
DeserializationFailurePolicy deserializationFailurePolicy
-
deserializationSchema
PubSubDeserializationSchema<T> deserializationSchema
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
orderingMode
OrderingMode orderingMode
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
startPosition
PubSubStartPosition startPosition
-
subscriberOptions
PubSubSubscriberOptions subscriberOptions
-
subscriptions
List<SubscriptionDestination> subscriptions
-
-
Class io.github.flink.gcp.connector.pubsub.source.PubSubStartPosition
class PubSubStartPosition extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
mode
PubSubStartPosition.Mode mode
-
timestamp
Instant timestamp
-
-
Class io.github.flink.gcp.connector.pubsub.source.PubSubSubscriberOptions
class PubSubSubscriberOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
awaitAckConfirmation
Duration awaitAckConfirmation
-
firstCheckpointTimeout
Duration firstCheckpointTimeout
-
flowControlMaxOutstandingElementCount
Long flowControlMaxOutstandingElementCount
-
flowControlMaxOutstandingRequestBytes
Long flowControlMaxOutstandingRequestBytes
-
maxAckExtensionPeriod
Duration maxAckExtensionPeriod
-
maxDurationPerAckExtension
Duration maxDurationPerAckExtension
-
maxRecordsPerFetch
int maxRecordsPerFetch
-
minDurationPerAckExtension
Duration minDurationPerAckExtension
-
parallelPullCount
Integer parallelPullCount
-
pausedSplitBufferMaxBytes
Long pausedSplitBufferMaxBytes
-
pausedSplitBufferMaxMessages
Long pausedSplitBufferMaxMessages
-
shutdownTimeout
Duration shutdownTimeout
-
subscriberBufferMaxBytes
Long subscriberBufferMaxBytes
-
subscriberBufferMaxMessages
Long subscriberBufferMaxMessages
-
-
Class io.github.flink.gcp.connector.pubsub.source.SubscriptionCreateOptions
class SubscriptionCreateOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
ackDeadline
Duration ackDeadline
-
deadLetterMaxDeliveryAttempts
int deadLetterMaxDeliveryAttempts
-
deadLetterTopic
TopicDestination deadLetterTopic
-
enableMessageOrdering
boolean enableMessageOrdering
-
expirationTtl
Duration expirationTtl
-
filter
String filter
-
messageRetention
Duration messageRetention
-
neverExpire
boolean neverExpire
-
retainAckedMessages
boolean retainAckedMessages
-
topic
TopicDestination topic
-
-
Class io.github.flink.gcp.connector.pubsub.source.SubscriptionDestination
class SubscriptionDestination extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.pubsub.source.serializer
-
Package io.github.flink.gcp.connector.pubsub.source.streamingpull
-
Class io.github.flink.gcp.connector.pubsub.source.streamingpull.PubSubStreamingPullSource
class PubSubStreamingPullSource extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
PubSubSourceConfig<T> config
-
lineageTableName
String lineageTableName
-
-
Class io.github.flink.gcp.connector.pubsub.source.streamingpull.SubscriberBufferLimitExceededEvent
class SubscriberBufferLimitExceededEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
attemptedBytes
long attemptedBytes
-
attemptedMessages
long attemptedMessages
-
maxBytes
long maxBytes
-
maxMessages
long maxMessages
-
splitId
String splitId
-
-
-
Package io.github.flink.gcp.connector.spanner
-
Class io.github.flink.gcp.connector.spanner.DatabaseDestination
class DatabaseDestination extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.SpannerTableName
class SpannerTableName extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.SpannerTableName.AccessPathName
class AccessPathName extends Object implements Serializable- serialVersionUID:
- 1L
-
-
Package io.github.flink.gcp.connector.spanner.sink
-
Class io.github.flink.gcp.connector.spanner.sink.SpannerMutationsSink
class SpannerMutationsSink extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
SpannerSinkConfig<T> config
-
-
Class io.github.flink.gcp.connector.spanner.sink.SpannerSinkConfig
class SpannerSinkConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
constraintViolationPolicy
ConstraintViolationPolicy constraintViolationPolicy
-
database
DatabaseDestination database
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
failedMutationHandler
FailureHandler<? super FailedMutation> failedMutationHandler
-
serializer
SpannerMutationSerializationSchema<? super T> serializer
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
writerOptions
SpannerWriterOptions writerOptions
-
-
Class io.github.flink.gcp.connector.spanner.sink.SpannerWriterOptions
class SpannerWriterOptions extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
batchWriteTimeout
Duration batchWriteTimeout
-
maxBatchBytes
long maxBatchBytes
-
maxBatchCells
int maxBatchCells
-
maxBatchMutations
int maxBatchMutations
-
maxCommitDelay
Duration maxCommitDelay
-
recoveryInitialBackoff
Duration recoveryInitialBackoff
-
recoveryMaxAttempts
int recoveryMaxAttempts
-
recoveryMaxBackoff
Duration recoveryMaxBackoff
-
rpcPriority
SpannerRpcPriority rpcPriority
-
-
-
Package io.github.flink.gcp.connector.spanner.sink.serializer
-
Package io.github.flink.gcp.connector.spanner.sink.writer
-
Class io.github.flink.gcp.connector.spanner.sink.writer.DefaultSpannerDatabaseAccessFactory
class DefaultSpannerDatabaseAccessFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
credentialsOverride
com.google.auth.Credentials credentialsOverride
-
database
DatabaseDestination database
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
writerOptions
SpannerWriterOptions writerOptions
-
-
-
Package io.github.flink.gcp.connector.spanner.source
-
Class io.github.flink.gcp.connector.spanner.source.SpannerChangeStreamSource
class SpannerChangeStreamSource extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
SpannerChangeStreamSourceConfig<T> config
-
-
Class io.github.flink.gcp.connector.spanner.source.SpannerChangeStreamSourceConfig
class SpannerChangeStreamSourceConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
absentRetentionFallback
Duration absentRetentionFallback
-
changeStreamName
String changeStreamName
-
coordinatorClientFactory
SpannerChangeStreamCoordinatorClientFactory coordinatorClientFactory
-
database
DatabaseDestination database
-
deserializer
SpannerChangeStreamDeserializationSchema<T> deserializer
-
endTimestamp
Instant endTimestamp
-
heartbeatMillis
long heartbeatMillis
-
maxConcurrentQueriesPerSubtask
int maxConcurrentQueriesPerSubtask
-
queryClientFactory
SpannerChangeStreamQueryClientFactory queryClientFactory
-
recordFilter
SpannerChangeStreamRecordFilter recordFilter
-
resumeFallback
StartPosition resumeFallback
-
rpcPriority
SpannerRpcPriority rpcPriority
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
startPosition
StartPosition startPosition
-
-
Class io.github.flink.gcp.connector.spanner.source.SpannerReadOperation
class SpannerReadOperation extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
columns
List<String> columns
-
index
String index
-
keys
com.google.cloud.spanner.KeySet keys
-
resolver
SpannerReadOperationResolver resolver
-
statement
com.google.cloud.spanner.Statement statement
-
table
String table
-
-
Class io.github.flink.gcp.connector.spanner.source.SpannerSourceConfig
class SpannerSourceConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
database
DatabaseDestination database
-
dataBoostEnabled
boolean dataBoostEnabled
-
deserializer
SpannerStructDeserializationSchema<T> deserializer
-
maxBytesPerFetch
long maxBytesPerFetch
-
maxRowsPerFetch
int maxRowsPerFetch
-
opener
StructStreamOpener opener
-
partitionOptions
com.google.cloud.spanner.PartitionOptions partitionOptions
-
plannerFactory
PartitionPlannerFactory plannerFactory
-
readOperation
SpannerReadOperation readOperation
-
rpcPriority
SpannerRpcPriority rpcPriority
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
timestampBound
com.google.cloud.spanner.TimestampBound timestampBound
-
-
-
Package io.github.flink.gcp.connector.spanner.source.batch
-
Class io.github.flink.gcp.connector.spanner.source.batch.SpannerBatchReadSource
class SpannerBatchReadSource extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
config
SpannerSourceConfig<T> config
-
-
-
Package io.github.flink.gcp.connector.spanner.source.batch.enumerator
-
Class io.github.flink.gcp.connector.spanner.source.batch.enumerator.DefaultPartitionPlannerFactory
class DefaultPartitionPlannerFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
database
DatabaseDestination database
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
-
-
Package io.github.flink.gcp.connector.spanner.source.batch.reader
-
Class io.github.flink.gcp.connector.spanner.source.batch.reader.BatchClientStructStreamOpener
class BatchClientStructStreamOpener extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
database
DatabaseDestination database
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
-
-
Package io.github.flink.gcp.connector.spanner.source.changestream
-
Class io.github.flink.gcp.connector.spanner.source.changestream.ChildPartitionsEvent
class ChildPartitionsEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
children
List<ChildPartitionsEvent.ChildPartition> children
-
parentSplitId
String parentSplitId
-
startTimestamp
Instant startTimestamp
-
-
Class io.github.flink.gcp.connector.spanner.source.changestream.ChildPartitionsEvent.ChildPartition
class ChildPartition extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.DataChangeRecord
class DataChangeRecord extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
columnTypes
List<DataChangeRecord.ColumnType> columnTypes
-
commitTimestamp
Instant commitTimestamp
-
lastRecordInTransactionInPartition
boolean lastRecordInTransactionInPartition
-
mods
List<Mod> mods
-
modType
ModType modType
-
numberOfPartitionsInTransaction
long numberOfPartitionsInTransaction
-
numberOfRecordsInTransaction
long numberOfRecordsInTransaction
-
recordSequence
String recordSequence
-
serverTransactionId
String serverTransactionId
-
systemTransaction
boolean systemTransaction
-
tableName
String tableName
-
transactionTag
String transactionTag
-
valueCaptureType
ValueCaptureType valueCaptureType
-
-
Class io.github.flink.gcp.connector.spanner.source.changestream.DataChangeRecord.ColumnType
class ColumnType extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.DataChangeRecordSerializer
class DataChangeRecordSerializer extends org.apache.flink.api.common.typeutils.TypeSerializer<DataChangeRecord> implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.Mod
class Mod extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.PartitionFinishedEvent
class PartitionFinishedEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.PartitionProgressEvent
class PartitionProgressEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.SpannerChangeStreamInitializationEvent
class SpannerChangeStreamInitializationEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
discardRestoredSplits
boolean discardRestoredSplits
-
sourceWatermark
long sourceWatermark
-
-
Class io.github.flink.gcp.connector.spanner.source.changestream.SpannerChangeStreamRecordFilter
class SpannerChangeStreamRecordFilter extends Object implements Serializable- serialVersionUID:
- 1L
-
Class io.github.flink.gcp.connector.spanner.source.changestream.SpannerChangeStreamWatermarkEvent
class SpannerChangeStreamWatermarkEvent extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
sourceWatermark
long sourceWatermark
-
-
-
Package io.github.flink.gcp.connector.spanner.source.changestream.enumerator
-
Class io.github.flink.gcp.connector.spanner.source.changestream.enumerator.DefaultSpannerChangeStreamCoordinatorClientFactory
class DefaultSpannerChangeStreamCoordinatorClientFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
absentRetentionFallback
Duration absentRetentionFallback
-
changeStreamName
String changeStreamName
-
database
DatabaseDestination database
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.spanner.source.changestream.reader
-
Class io.github.flink.gcp.connector.spanner.source.changestream.reader.DefaultSpannerChangeStreamQueryClientFactory
class DefaultSpannerChangeStreamQueryClientFactory extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
changeStreamName
String changeStreamName
-
database
DatabaseDestination database
-
emulatorEndpoint
EmulatorEndpoint emulatorEndpoint
-
maxConcurrentQueries
int maxConcurrentQueries
-
rpcPriority
SpannerRpcPriority rpcPriority
-
serviceAccountKeyFile
String serviceAccountKeyFile
-
-
-
Package io.github.flink.gcp.connector.spanner.source.serializer
-
Package io.github.flink.gcp.connector.spanner.table
-
Class io.github.flink.gcp.connector.spanner.table.SpannerLookupConfig
class SpannerLookupConfig extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
Class io.github.flink.gcp.connector.spanner.table.SpannerTableLineage
class SpannerTableLineage extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
logicalName
String logicalName
-
namespace
String namespace
-
resources
List<ResourceIdentifier> resources
-
-
Class io.github.flink.gcp.connector.spanner.table.SpannerTableSchemaConverter
class SpannerTableSchemaConverter extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
columns
List<SpannerTableSchemaConverter.Column> columns
-
primaryKeyIndexes
int[] primaryKeyIndexes
-
-
Class io.github.flink.gcp.connector.spanner.table.SpannerTableSchemaConverter.Column
class Column extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
index
int index
-
logicalType
org.apache.flink.table.types.logical.LogicalType logicalType
-
name
String name
-
spannerType
com.google.cloud.spanner.Type spannerType
-
-
-
Package io.github.flink.gcp.connector.spanner.table.sink
-
Class io.github.flink.gcp.connector.spanner.table.sink.RowDataSerializationSchema
class RowDataSerializationSchema extends Object implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
schema
SpannerTableSchemaConverter schema
-
table
String table
-
-
-
Package io.github.flink.gcp.connector.spanner.table.source
-
Class io.github.flink.gcp.connector.spanner.table.source.SpannerRowDataAsyncLookupFunction
class SpannerRowDataAsyncLookupFunction extends org.apache.flink.table.functions.AsyncLookupFunction implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
converter
io.github.flink.gcp.connector.spanner.table.source.StructToRowDataConverter converter
-
filters
io.github.flink.gcp.connector.spanner.table.source.SpannerFilterPushDown.RuntimeState filters
-
keyEncoder
io.github.flink.gcp.connector.spanner.table.source.SpannerLookupKeyEncoder keyEncoder
-
lookup
io.github.flink.gcp.connector.spanner.table.source.SpannerRowLookup lookup
-
maxRetries
int maxRetries
-
-
Class io.github.flink.gcp.connector.spanner.table.source.SpannerRowDataLookupFunction
class SpannerRowDataLookupFunction extends org.apache.flink.table.functions.LookupFunction implements Serializable- serialVersionUID:
- 1L
-
Serialized Fields
-
converter
io.github.flink.gcp.connector.spanner.table.source.StructToRowDataConverter converter
-
filters
io.github.flink.gcp.connector.spanner.table.source.SpannerFilterPushDown.RuntimeState filters
-
keyEncoder
io.github.flink.gcp.connector.spanner.table.source.SpannerLookupKeyEncoder keyEncoder
-
lookup
io.github.flink.gcp.connector.spanner.table.source.SpannerRowLookup lookup
-
maxRetries
int maxRetries
-
-