Class PayloadSerializerKafkaTable

    • 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