Class BeamKafkaCSVTable

    • 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:
        getPTransformForInput in class BeamKafkaTable
      • 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:
        getPTransformForOutput in class BeamKafkaTable