public abstract class PubsubClient
extends java.lang.Object
implements java.io.Closeable
Modifier and Type | Class and Description |
---|---|
static class |
PubsubClient.ProjectPath
Path representing a cloud project id.
|
static interface |
PubsubClient.PubsubClientFactory
Factory for creating clients.
|
static class |
PubsubClient.SubscriptionPath
Path representing a Pubsub subscription.
|
static class |
PubsubClient.TopicPath
Path representing a Pubsub topic.
|
Constructor and Description |
---|
PubsubClient() |
Modifier and Type | Method and Description |
---|---|
abstract int |
ackDeadlineSeconds(PubsubClient.SubscriptionPath subscription)
Return the ack deadline, in seconds, for
subscription . |
abstract void |
acknowledge(PubsubClient.SubscriptionPath subscription,
java.util.List<java.lang.String> ackIds)
Acknowldege messages from
subscription with ackIds . |
PubsubClient.SubscriptionPath |
createRandomSubscription(PubsubClient.ProjectPath project,
PubsubClient.TopicPath topic,
int ackDeadlineSeconds)
Create a random subscription for
topic . |
abstract void |
createSubscription(PubsubClient.TopicPath topic,
PubsubClient.SubscriptionPath subscription,
int ackDeadlineSeconds)
Create
subscription to topic . |
abstract void |
createTopic(PubsubClient.TopicPath topic)
Create
topic . |
abstract void |
deleteSubscription(PubsubClient.SubscriptionPath subscription)
Delete
subscription . |
abstract void |
deleteTopic(PubsubClient.TopicPath topic) |
protected static long |
extractTimestamp(java.lang.String timestampAttribute,
java.lang.String pubsubTimestamp,
java.util.Map<java.lang.String,java.lang.String> attributes)
Return the timestamp (in ms since unix epoch) to use for a Pubsub message with
attributes and pubsubTimestamp . |
abstract boolean |
isEOF()
Return true if
pull(long, org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.SubscriptionPath, int, boolean) will always return empty list. |
abstract java.util.List<PubsubClient.SubscriptionPath> |
listSubscriptions(PubsubClient.ProjectPath project,
PubsubClient.TopicPath topic)
Return a list of subscriptions for
topic in project . |
abstract java.util.List<PubsubClient.TopicPath> |
listTopics(PubsubClient.ProjectPath project)
Return a list of topics for
project . |
abstract void |
modifyAckDeadline(PubsubClient.SubscriptionPath subscription,
java.util.List<java.lang.String> ackIds,
int deadlineSeconds)
Modify the ack deadline for messages from
subscription with ackIds to
be deadlineSeconds from now. |
static PubsubClient.ProjectPath |
projectPathFromId(java.lang.String projectId) |
static PubsubClient.ProjectPath |
projectPathFromPath(java.lang.String path) |
abstract int |
publish(PubsubClient.TopicPath topic,
java.util.List<org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.OutgoingMessage> outgoingMessages)
Publish
outgoingMessages to Pubsub topic . |
abstract java.util.List<org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.IncomingMessage> |
pull(long requestTimeMsSinceEpoch,
PubsubClient.SubscriptionPath subscription,
int batchSize,
boolean returnImmediately)
Request the next batch of up to
batchSize messages from subscription . |
static PubsubClient.SubscriptionPath |
subscriptionPathFromName(java.lang.String projectId,
java.lang.String subscriptionName) |
static PubsubClient.SubscriptionPath |
subscriptionPathFromPath(java.lang.String path) |
static PubsubClient.TopicPath |
topicPathFromName(java.lang.String projectId,
java.lang.String topicName) |
static PubsubClient.TopicPath |
topicPathFromPath(java.lang.String path) |
protected static long extractTimestamp(@Nullable java.lang.String timestampAttribute, @Nullable java.lang.String pubsubTimestamp, @Nullable java.util.Map<java.lang.String,java.lang.String> attributes)
attributes
and pubsubTimestamp
.
If timestampAttribute
is non-null then the message attributes must
contain that attribute, and the value of that attribute will be taken as the timestamp.
Otherwise the timestamp will be taken from the Pubsub publish timestamp pubsubTimestamp
.
java.lang.IllegalArgumentException
- if the timestamp cannot be recognized as a ms-since-unix-epoch
or RFC3339 time.public static PubsubClient.ProjectPath projectPathFromPath(java.lang.String path)
public static PubsubClient.ProjectPath projectPathFromId(java.lang.String projectId)
public static PubsubClient.SubscriptionPath subscriptionPathFromPath(java.lang.String path)
public static PubsubClient.SubscriptionPath subscriptionPathFromName(java.lang.String projectId, java.lang.String subscriptionName)
public static PubsubClient.TopicPath topicPathFromPath(java.lang.String path)
public static PubsubClient.TopicPath topicPathFromName(java.lang.String projectId, java.lang.String topicName)
public abstract int publish(PubsubClient.TopicPath topic, java.util.List<org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.OutgoingMessage> outgoingMessages) throws java.io.IOException
outgoingMessages
to Pubsub topic
. Return number of messages
published.java.io.IOException
public abstract java.util.List<org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.IncomingMessage> pull(long requestTimeMsSinceEpoch, PubsubClient.SubscriptionPath subscription, int batchSize, boolean returnImmediately) throws java.io.IOException
batchSize
messages from subscription
.
Return the received messages, or empty collection if none were available. Does not
wait for messages to arrive if returnImmediately
is true.
Returned messages will record their request time as requestTimeMsSinceEpoch
.java.io.IOException
public abstract void acknowledge(PubsubClient.SubscriptionPath subscription, java.util.List<java.lang.String> ackIds) throws java.io.IOException
subscription
with ackIds
.java.io.IOException
public abstract void modifyAckDeadline(PubsubClient.SubscriptionPath subscription, java.util.List<java.lang.String> ackIds, int deadlineSeconds) throws java.io.IOException
subscription
with ackIds
to
be deadlineSeconds
from now.java.io.IOException
public abstract void createTopic(PubsubClient.TopicPath topic) throws java.io.IOException
topic
.java.io.IOException
public abstract void deleteTopic(PubsubClient.TopicPath topic) throws java.io.IOException
java.io.IOException
public abstract java.util.List<PubsubClient.TopicPath> listTopics(PubsubClient.ProjectPath project) throws java.io.IOException
project
.java.io.IOException
public abstract void createSubscription(PubsubClient.TopicPath topic, PubsubClient.SubscriptionPath subscription, int ackDeadlineSeconds) throws java.io.IOException
subscription
to topic
.java.io.IOException
public PubsubClient.SubscriptionPath createRandomSubscription(PubsubClient.ProjectPath project, PubsubClient.TopicPath topic, int ackDeadlineSeconds) throws java.io.IOException
topic
. Return the PubsubClient.SubscriptionPath
. It
is the responsibility of the caller to later delete the subscription.java.io.IOException
public abstract void deleteSubscription(PubsubClient.SubscriptionPath subscription) throws java.io.IOException
subscription
.java.io.IOException
public abstract java.util.List<PubsubClient.SubscriptionPath> listSubscriptions(PubsubClient.ProjectPath project, PubsubClient.TopicPath topic) throws java.io.IOException
topic
in project
.java.io.IOException
public abstract int ackDeadlineSeconds(PubsubClient.SubscriptionPath subscription) throws java.io.IOException
subscription
.java.io.IOException
public abstract boolean isEOF()
pull(long, org.apache.beam.sdk.io.gcp.pubsub.PubsubClient.SubscriptionPath, int, boolean)
will always return empty list. Actual clients
will return false. Test clients may return true to signal that all
expected messages have been pulled and the test may complete.