Class PubSubSource

java.lang.Object
io.github.flink.gcp.connector.pubsub.source.PubSubSource

@Public public final class PubSubSource extends Object
Entry point for building Pub/Sub sources.

 PubSubDeserializationSchema<String> deserializer =
         PubSubDeserializationSchema.payload(new SimpleStringSchema());
 Source<String, ?, ?> source =
         PubSubSource.<String>builder()
                 .subscription(SubscriptionDestination.of("my-project", "my-subscription"))
                 .deserializer(deserializer)
                 .build();
 
  • Method Details

    • builder

      public static <T> PubSubSourceBuilder<T> builder()
      Returns a new builder. A deserializer and at least one subscription are required.
      Type Parameters:
      T - type of the records produced by the source
      Returns:
      the builder