Class BeamKafkaCSVTable
- 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
-
- org.apache.beam.sdk.extensions.sql.meta.provider.kafka.BeamKafkaCSVTable
-
- All Implemented Interfaces:
java.io.Serializable,BeamSqlTable
public class BeamKafkaCSVTable extends BeamKafkaTable
A Kafka topic that saves records as CSV format.- See Also:
- Serialized Form
-
-
Field Summary
-
Fields inherited from class org.apache.beam.sdk.extensions.sql.meta.provider.kafka.BeamKafkaTable
numberOfRecordsForRate
-
Fields inherited from class org.apache.beam.sdk.extensions.sql.meta.SchemaBaseBeamTable
schema
-
-
Constructor Summary
Constructors Constructor Description BeamKafkaCSVTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics)BeamKafkaCSVTable(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)BeamKafkaCSVTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics, org.apache.commons.csv.CSVFormat format, org.apache.beam.sdk.io.kafka.TimestampPolicyFactory timestampPolicyFactory)
-
Method Summary
All Methods Instance Methods Concrete Methods Modifier and Type Method Description protected 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 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()-
Methods inherited from class org.apache.beam.sdk.extensions.sql.meta.provider.kafka.BeamKafkaTable
buildIOReader, buildIOWriter, createKafkaRead, getBootstrapServers, getConfigUpdates, getTableStatistics, getTimestampPolicyFactory, getTopics, isBounded, updateConsumerProperties
-
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
-
BeamKafkaCSVTable
public BeamKafkaCSVTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics)
-
BeamKafkaCSVTable
public BeamKafkaCSVTable(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)
-
BeamKafkaCSVTable
public BeamKafkaCSVTable(org.apache.beam.sdk.schemas.Schema beamSchema, java.lang.String bootstrapServers, java.util.List<java.lang.String> topics, org.apache.commons.csv.CSVFormat format, org.apache.beam.sdk.io.kafka.TimestampPolicyFactory timestampPolicyFactory)
-
-
Method Detail
-
getPTransformForInput
protected 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()
- Specified by:
getPTransformForInputin classBeamKafkaTable
-
getPTransformForOutput
protected 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()
- Specified by:
getPTransformForOutputin classBeamKafkaTable
-
-