org.apache.beam.sdk.transforms.SerializableFunction<InputT,OutputT> schemaRegistryClientProviderFn
java.lang.String schemaRegistryUrl
java.lang.String subject
java.lang.Integer version
java.lang.String topic
int partition
long nextOffset
long watermarkMillis
KafkaIO.ReadSourceDescriptors<K,V> readSourceDescriptors
KafkaIO.Read<K,V> read
org.apache.beam.sdk.transforms.SerializableFunction<InputT,OutputT> valueMapper
org.apache.beam.sdk.metrics.Counter errorCounter
java.lang.Long errorsInBundle
boolean handleErrors
org.apache.beam.sdk.schemas.Schema errorSchema
org.apache.beam.sdk.coders.KvCoder<K,V> kvCoder
org.apache.kafka.common.TopicPartition topicPartition
org.apache.beam.sdk.transforms.SerializableFunction<InputT,OutputT> toBytesFn
org.apache.beam.sdk.metrics.Counter errorCounter
java.lang.Long errorsInBundle
boolean handleErrors
org.apache.beam.sdk.schemas.Schema errorSchema
org.apache.beam.sdk.coders.KvCoder<K,V> kvCoder