public class PubsubUnboundedSource extends PTransform<PBegin,PCollection<PubsubMessage>>
PubsubIO#read instead.
A PTransform which streams messages from Pubsub.
UnboundedSource which receives messages in
batches and hands them out one at a time.
UnboundedSource.UnboundedReader instances to execute concurrently and thus hide latency.
name, resourceHints| Constructor and Description |
|---|
PubsubUnboundedSource(com.google.api.client.util.Clock clock,
PubsubClient.PubsubClientFactory pubsubFactory,
@Nullable ValueProvider<PubsubClient.ProjectPath> project,
@Nullable ValueProvider<PubsubClient.TopicPath> topic,
@Nullable ValueProvider<PubsubClient.SubscriptionPath> subscription,
@Nullable java.lang.String timestampAttribute,
@Nullable java.lang.String idAttribute,
boolean needsAttributes)
Construct an unbounded source to consume from the Pubsub
subscription. |
PubsubUnboundedSource(PubsubClient.PubsubClientFactory pubsubFactory,
@Nullable ValueProvider<PubsubClient.ProjectPath> project,
@Nullable ValueProvider<PubsubClient.TopicPath> topic,
@Nullable ValueProvider<PubsubClient.SubscriptionPath> subscription,
@Nullable java.lang.String timestampAttribute,
@Nullable java.lang.String idAttribute,
boolean needsAttributes)
Construct an unbounded source to consume from the Pubsub
subscription. |
PubsubUnboundedSource(PubsubClient.PubsubClientFactory pubsubFactory,
@Nullable ValueProvider<PubsubClient.ProjectPath> project,
@Nullable ValueProvider<PubsubClient.TopicPath> topic,
@Nullable ValueProvider<PubsubClient.SubscriptionPath> subscription,
@Nullable java.lang.String timestampAttribute,
@Nullable java.lang.String idAttribute,
boolean needsAttributes,
boolean needsMessageId)
Construct an unbounded source to consume from the Pubsub
subscription. |
| Modifier and Type | Method and Description |
|---|---|
PCollection<PubsubMessage> |
expand(PBegin input)
Override this method to specify how this
PTransform should be expanded on the given
InputT. |
@Nullable java.lang.String |
getIdAttribute()
Get the id attribute.
|
boolean |
getNeedsAttributes() |
boolean |
getNeedsMessageId() |
boolean |
getNeedsOrderingKey() |
@Nullable PubsubClient.ProjectPath |
getProject()
Get the project path.
|
@Nullable PubsubClient.SubscriptionPath |
getSubscription()
Get the subscription being read from.
|
@Nullable ValueProvider<PubsubClient.SubscriptionPath> |
getSubscriptionProvider()
Get the
ValueProvider for the subscription being read from. |
@Nullable java.lang.String |
getTimestampAttribute()
Get the timestamp attribute.
|
@Nullable PubsubClient.TopicPath |
getTopic()
Get the topic being read from.
|
@Nullable ValueProvider<PubsubClient.TopicPath> |
getTopicProvider()
Get the
ValueProvider for the topic being read from. |
compose, compose, getAdditionalInputs, getDefaultOutputCoder, getDefaultOutputCoder, getDefaultOutputCoder, getKindString, getName, getResourceHints, populateDisplayData, setResourceHints, toString, validate, validatepublic PubsubUnboundedSource(PubsubClient.PubsubClientFactory pubsubFactory, @Nullable ValueProvider<PubsubClient.ProjectPath> project, @Nullable ValueProvider<PubsubClient.TopicPath> topic, @Nullable ValueProvider<PubsubClient.SubscriptionPath> subscription, @Nullable java.lang.String timestampAttribute, @Nullable java.lang.String idAttribute, boolean needsAttributes)
subscription.public PubsubUnboundedSource(com.google.api.client.util.Clock clock,
PubsubClient.PubsubClientFactory pubsubFactory,
@Nullable ValueProvider<PubsubClient.ProjectPath> project,
@Nullable ValueProvider<PubsubClient.TopicPath> topic,
@Nullable ValueProvider<PubsubClient.SubscriptionPath> subscription,
@Nullable java.lang.String timestampAttribute,
@Nullable java.lang.String idAttribute,
boolean needsAttributes)
subscription.public PubsubUnboundedSource(PubsubClient.PubsubClientFactory pubsubFactory, @Nullable ValueProvider<PubsubClient.ProjectPath> project, @Nullable ValueProvider<PubsubClient.TopicPath> topic, @Nullable ValueProvider<PubsubClient.SubscriptionPath> subscription, @Nullable java.lang.String timestampAttribute, @Nullable java.lang.String idAttribute, boolean needsAttributes, boolean needsMessageId)
subscription.public @Nullable PubsubClient.ProjectPath getProject()
public @Nullable PubsubClient.TopicPath getTopic()
public @Nullable ValueProvider<PubsubClient.TopicPath> getTopicProvider()
ValueProvider for the topic being read from.public @Nullable PubsubClient.SubscriptionPath getSubscription()
public @Nullable ValueProvider<PubsubClient.SubscriptionPath> getSubscriptionProvider()
ValueProvider for the subscription being read from.public @Nullable java.lang.String getTimestampAttribute()
public @Nullable java.lang.String getIdAttribute()
public boolean getNeedsAttributes()
public boolean getNeedsMessageId()
public boolean getNeedsOrderingKey()
public PCollection<PubsubMessage> expand(PBegin input)
PTransformPTransform should be expanded on the given
InputT.
NOTE: This method should not be called directly. Instead apply the PTransform should
be applied to the InputT using the apply method.
Composite transforms, which are defined in terms of other transforms, should return the output of one of the composed transforms. Non-composite transforms, which do not apply any transforms internally, should return a new unbound output and register evaluators (via backend-specific registration methods).
expand in class PTransform<PBegin,PCollection<PubsubMessage>>