如何串行拉取Pub/Sub消息?解决未ACK时重复拉取问题
Root Cause
The default Pub/Sub Java subscriber uses flow control settings that allow multiple unacked messages to be held in memory. Additionally, if you don’t acknowledge a message within the default 10-second acknowledgment deadline, Pub/Sub automatically redelivers it—this is why you’re seeing repeated pulls. The subscriber.isRunning() method only checks if the subscriber service is active; it doesn’t control message delivery flow.
Step-by-Step Fixes
1. Restrict Flow Control to One Unacked Message
Configure the subscriber to only fetch and hold one message at a time. This ensures no new messages are pulled until you explicitly ack or nack the current one. Use FlowControlSettings to limit outstanding messages:
import com.google.cloud.pubsub.v1.Subscriber; import com.google.pubsub.v1.ProjectSubscriptionName; import com.google.cloud.pubsub.v1.MessageReceiver; import com.google.cloud.pubsub.v1.FlowControlSettings; // Build flow control to limit 1 unacked message FlowControlSettings flowControl = FlowControlSettings.newBuilder() .setMaxOutstandingMessages(1L) // Only 1 message in flight at a time .setMaxOutstandingBytes(Long.MAX_VALUE) // No byte limit (adjust if needed) .build(); // Initialize subscriber with custom flow control ProjectSubscriptionName subscriptionName = ProjectSubscriptionName.of(projectId, subscriptionId); subscriber = Subscriber.newBuilder(subscriptionName, messageReceiver) .setFlowControlSettings(flowControl) .build(); subscriber.startAsync().awaitRunning(); subscriber.awaitTerminated();
2. Extend Ack Deadline for Long-Running Tasks
If your message processing takes longer than 10 seconds, extend the acknowledgment deadline periodically to prevent premature redelivery. Use modifyAckDeadline on the AckReplyConsumer:
MessageReceiver messageReceiver = (message, consumer) -> { try { // Simulate long-running processing (e.g., 30 seconds) Thread.sleep(30000); // Extend deadline to 60 seconds (call repeatedly for longer tasks) consumer.modifyAckDeadline(60); // Acknowledge once processing is complete consumer.ack(); } catch (InterruptedException e) { // Nack to trigger immediate redelivery on failure consumer.nack(); Thread.currentThread().interrupt(); } };
3. Critical Notes
- Always call
ack()ornack()after processing: Failing to do so will trigger redelivery once the deadline expires. - Adjust flow control limits: If your messages are large, set a reasonable
maxOutstandingBytesto avoid memory overload. - Clean shutdown: Call
subscriber.stopAsync().awaitTerminated()when closing the application to handle in-flight messages properly.
内容的提问来源于stack exchange,提问作者ankita ojha

