Class BeamKafkaTable
- java.lang.Object
-
- org.apache.beam.sdk.extensions.sql.meta.BaseBeamTable
-
- org.apache.beam.sdk.extensions.sql.meta.SchemaBaseBeamTable
-
- org.apache.beam.sdk.extensions.sql.meta.provider.kafka.BeamKafkaTable
-
- All Implemented Interfaces:
java.io.Serializable,BeamSqlTable
- Direct Known Subclasses:
BeamKafkaCSVTable,PayloadSerializerKafkaTable
public abstract class BeamKafkaTable extends SchemaBaseBeamTable
BeamKafkaTablerepresent a Kafka topic, as source or target. Need to extend to convert betweenBeamSqlRowandKV<byte[], byte[]>.- See Also:
- Serialized Form
-
-
Field Summary
Fields Modifier and Type Field Description protected intnumberOfRecordsForRate-
Fields inherited from class org.apache.beam.sdk.extensions.sql.meta.SchemaBaseBeamTable
schema
-
-
Constructor Summary
Constructors Modifier Constructor Description protectedBeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema)BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics)BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics, org.apache.beam.sdk.io.kafka.TimestampPolicyFactory timestampPolicyFactory)BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.util.List<org.apache.kafka.common.TopicPartition> topicPartitions, java.lang.String bootstrapServers)BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.util.List<org.apache.kafka.common.TopicPartition> topicPartitions, java.lang.String bootstrapServers, org.apache.beam.sdk.io.kafka.TimestampPolicyFactory timestampPolicyFactory)
-
Method Summary
All Methods Instance Methods Abstract Methods Concrete Methods Modifier and Type Method Description org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>buildIOReader(org.apache.beam.sdk.values.PBegin begin)create aPCollection<Row>from source.org.apache.beam.sdk.values.POutputbuildIOWriter(org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> input)create aIO.write()instance to write to target.protected org.apache.beam.sdk.io.kafka.KafkaIO.Read<byte[],byte[]>createKafkaRead()java.lang.StringgetBootstrapServers()java.util.Map<java.lang.String,java.lang.Object>getConfigUpdates()protected abstract org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.io.kafka.KafkaRecord<byte[],byte[]>>,org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>>getPTransformForInput()protected abstract org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,org.apache.beam.sdk.values.PCollection<org.apache.kafka.clients.producer.ProducerRecord<byte[],byte[]>>>getPTransformForOutput()BeamTableStatisticsgetTableStatistics(org.apache.beam.sdk.options.PipelineOptions options)Estimates the number of rows or the rate for unbounded Tables.org.apache.beam.sdk.io.kafka.TimestampPolicyFactorygetTimestampPolicyFactory()java.util.List<java.lang.String>getTopics()org.apache.beam.sdk.values.PCollection.IsBoundedisBounded()Whether this table is bounded (known to be finite) or unbounded (may or may not be finite).BeamKafkaTableupdateConsumerProperties(java.util.Map<java.lang.String,java.lang.Object> configUpdates)-
Methods inherited from class org.apache.beam.sdk.extensions.sql.meta.SchemaBaseBeamTable
getSchema
-
Methods inherited from class org.apache.beam.sdk.extensions.sql.meta.BaseBeamTable
buildIOReader, constructFilter, supportsProjects
-
-
-
-
Constructor Detail
-
BeamKafkaTable
protected BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema)
-
BeamKafkaTable
public BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics)
-
BeamKafkaTable
public BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics, org.apache.beam.sdk.io.kafka.TimestampPolicyFactory timestampPolicyFactory)
-
BeamKafkaTable
public BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.util.List<org.apache.kafka.common.TopicPartition> topicPartitions, java.lang.String bootstrapServers)
-
BeamKafkaTable
public BeamKafkaTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.util.List<org.apache.kafka.common.TopicPartition> topicPartitions, java.lang.String bootstrapServers, org.apache.beam.sdk.io.kafka.TimestampPolicyFactory timestampPolicyFactory)
-
-
Method Detail
-
updateConsumerProperties
public BeamKafkaTable updateConsumerProperties(java.util.Map<java.lang.String,java.lang.Object> configUpdates)
-
isBounded
public org.apache.beam.sdk.values.PCollection.IsBounded isBounded()
Description copied from interface:BeamSqlTableWhether this table is bounded (known to be finite) or unbounded (may or may not be finite).
-
getPTransformForInput
protected abstract org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.io.kafka.KafkaRecord<byte[],byte[]>>,org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>> getPTransformForInput()
-
getPTransformForOutput
protected abstract org.apache.beam.sdk.transforms.PTransform<org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row>,org.apache.beam.sdk.values.PCollection<org.apache.kafka.clients.producer.ProducerRecord<byte[],byte[]>>> getPTransformForOutput()
-
buildIOReader
public org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> buildIOReader(org.apache.beam.sdk.values.PBegin begin)
Description copied from interface:BeamSqlTablecreate aPCollection<Row>from source.
-
createKafkaRead
protected org.apache.beam.sdk.io.kafka.KafkaIO.Read<byte[],byte[]> createKafkaRead()
-
buildIOWriter
public org.apache.beam.sdk.values.POutput buildIOWriter(org.apache.beam.sdk.values.PCollection<org.apache.beam.sdk.values.Row> input)
Description copied from interface:BeamSqlTablecreate aIO.write()instance to write to target.
-
getBootstrapServers
public java.lang.String getBootstrapServers()
-
getTopics
public java.util.List<java.lang.String> getTopics()
-
getConfigUpdates
public java.util.Map<java.lang.String,java.lang.Object> getConfigUpdates()
-
getTimestampPolicyFactory
public org.apache.beam.sdk.io.kafka.TimestampPolicyFactory getTimestampPolicyFactory()
-
getTableStatistics
public BeamTableStatistics getTableStatistics(org.apache.beam.sdk.options.PipelineOptions options)
Description copied from interface:BeamSqlTableEstimates the number of rows or the rate for unbounded Tables. If it is not possible to estimate the row count or rate it will return BeamTableStatistics.BOUNDED_UNKNOWN.- Specified by:
getTableStatisticsin interfaceBeamSqlTable- Overrides:
getTableStatisticsin classBaseBeamTable
-
-