Skip to content

Pub/Sub streamingPull subscriber: large number of duplicate messages, modifyAckDeadline calls observed #2465

Description

@kir-titievsky

Large number of duplicate messages is observed with Pub/Sub streamingPull client library using code that pulls from a Pub/Sub subscription and inserts messages into BigQuery with a synchronous blocking operation, with flowControl set to max 500 outstanding messages. See [1] for code.

For the same code, we also observe an excessive number of modifyAckDeadline operations (>> streamingPull message operations). And tracing a single message, we see modifyAcks and Acks in alternating order for the same message (modAck, modAck, Ack, modAck, Ack) [2]. This suggest that the implementation might fail to remove Ack'ed messages from a queue of messages to process and keep re-processing messages already on the client. This also suggests that ack requests may not actually be sent.

[2] https://docs.google.com/spreadsheets/d/1mqtxxm0guZcOcRy8ORG0ri787XLQZNF_FLBujAiayFI/edit?ts=59cac0f0#gid=2139642597

[1]


package kir.pubsub;

import com.google.api.gax.batching.FlowControlSettings;
import com.google.cloud.bigquery.*;
import com.google.cloud.pubsub.v1.Subscriber;
import com.google.pubsub.v1.PubsubMessage;
import com.google.pubsub.v1.SubscriptionName;
import com.google.cloud.pubsub.v1.MessageReceiver;
import com.google.cloud.pubsub.v1.AckReplyConsumer;

import java.text.DateFormat;
import java.text.SimpleDateFormat;
import java.util.*;
import java.util.concurrent.atomic.AtomicInteger;

public class Sub {

// Instantiate an asynchronous message receiver

public static void main(String... args) throws Exception {
    final String projectId = args[0];
    final String subscriptionId = args[1];

    final BigQuery bq = BigQueryOptions.getDefaultInstance().getService();
    final String datasetName = "pubsub_debug";
    final String tableName = "gke_subscriber";
    final AtomicInteger messageCount = new AtomicInteger(0);
    final DateFormat dateFormat = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS");
    MessageReceiver receiver = new MessageReceiver() {
                public void receiveMessage(PubsubMessage message, AckReplyConsumer consumer) {
                    // handle incoming message, then ack/nack the received message
                    System.out.printf("%s,\t%s,\t%d\n",message.getData().toStringUtf8()
                            , message.getMessageId()
                            , messageCount.incrementAndGet());
                    Map<String,Object> row = new HashMap<String,Object>();
                    long timestampMs = message.getPublishTime().getSeconds()*1000 + message.getPublishTime().getNanos() / 1000000;
                    Date timestampDate = new Date(timestampMs);

                    row.put("messageId", message.getMessageId());
                    row.put("messageData", message.getData().toStringUtf8());
                    row.put("messagePublishTime", dateFormat.format(timestampDate));
                    // a version of this code without the bq.insert was ran, where consumer.ack() 
                    // was called immediately. The results were the same.
                    InsertAllResponse response = bq.insertAll(
                                InsertAllRequest.newBuilder(TableId.of(projectId, datasetName, tableName)).addRow(row).build()
                    );
                    if (response.hasErrors()) {
                        System.err.println("Error inserting into BigQuery " + response.toString() );
                        consumer.nack();
                    } else{
                        consumer.ack();
                    }
                }

            };

    SubscriptionName subscriptionName = SubscriptionName.create(projectId, subscriptionId);
    Subscriber subscriber = Subscriber.defaultBuilder(subscriptionName, receiver).setFlowControlSettings(
            FlowControlSettings.newBuilder().setMaxOutstandingElementCount(1000L).build()
    ).build();
    subscriber.startAsync();
    subscriber.awaitRunning();
    System.out.println("Started async subscriber.");
    subscriber.awaitTerminated();
}

}

Activity

  1. added
    api: pubsubIssues related to the Pub/Sub API.
    priority: p1Important issue which blocks shipping the next release. Will be fixed prior to next release.
    type: bugError or flaw in code with unintended results or allowing sub-optimal usage patterns.
    on Sep 26, 2017
  2. pongad commented on Sep 27, 2017

    @pongad
    Contributor

    @kir-titievsky Do you have a workload that can reliably reproduce this? If you do, could you try removing the flow control (hopefully the workload isn't so large that it crashes your machine)? If the problem goes away, this is probably a dup of #2452; the symptoms are nearly identical.

    If this doesn't fix the problem, could you share the repro workload with me?

    EDIT: If the workload is too large to remove flow control, the fix for the linked issue is already in master, so we can to test with that version. Slightly less convenient as we'll need to compile from source.

  3. robertsaxby commented on Sep 27, 2017

    @robertsaxby

    The behaviour observed with #2452 was seen with the older non streaming pull implementation (0.21.1-beta). This issue came about when using the newer client library. It might also be worth noting that with the Flow Control set to 1000 max messages on the older library duplicates where not seen, just the stuck messages.

  4. pongad commented on Sep 28, 2017

    @pongad
    Contributor

    @robertsaxby That makes sense. I can reproduce the messages getting stuck, but not redelivery.

    @kir-titievsky I need a little help understanding the spreadsheet. Are all rows for the same message? I assume that stream_id uniquely identifies a StreamingPull stream? Ie, one physical computer can have multiple IDs by opening many streams, but one stream ID is unique to one computer?

    FWIW, I have a PR opened to address a potential race condition in streaming. It's concievable that the race condition causes this problem.

    If you could set up a reproduction, please let me know.

  5. kir-titievsky commented on Sep 28, 2017

    @kir-titievsky
    Author

    @pongad You are right on all counts about the spreadsheet. All rows are for the same message, including traffic from several bidi streams, that had been opened at different times.

  6. kir-titievsky commented on Sep 28, 2017

    @kir-titievsky
    Author

    @pongad Did a couple experiments:

    • No flow control, single machine: everything is acked within seconds ~10 modify ack deadline operations for 10K messages.
    • Added a synchronous 50ms sleep to every operation. Got a peak of 7K modifyAckDeadline operations over a minute, about 10K total. modifyAckDeadline ops peaked around the same time as Ack operations. Overall, ~2-3 minutes of activity. This is very much unexpected as the subscription has a 10 minute ack deadline (at least that's the setting).
  7. pongad commented on Oct 18, 2017

    @pongad
    Contributor

    To summarize: This is an issue on the server side. Currently the client library does not have enough information to properly handle ack deadlines.

    In the immediate term, consider using v0.21.1. That version uses a different (slower) pubsub endpoint, that isn't affected by this problem.

  8. pongad commented on Nov 1, 2017

    @pongad
    Contributor

    Pubsub team is working on a server-side mitigation for this. The client lib will need to be updated to take advantage of it. Fortunately, this new feature "piggybacks" on an already existing one, so the work on client lib can progress right away.

    I hope to create a PR for this soon.

  9. pongad commented on Nov 10, 2017

    @pongad
    Contributor

    Update: the server side release should happen this week. The feature should be enabled next week. The client (in master) has already been modified to take advantage of this new feature.

    When the server feature is enabled, we'll test to see how this helps.

  10. pongad commented on Nov 22, 2017

    @pongad
    Contributor

    The server-side fix has landed. If you are affected, could you try again and see if you observe fewer duplicates?

    While we expect the fix to help reduce duplication on older client libs, I'd encourage moving to latest release (v0.30.0) since more fixes has landed during that time.

  11. ericmartineau commented on Dec 11, 2017

    @ericmartineau

    We are seeing this behavior currently. We are using v0.30.0-beta of the pubsub library, and our subscriptions are all set to a 60s ack deadline. We have a subscription that is currently extended the deadline for over 750K unacked messages:

    Screenshot of stackdriver below:
    https://screencast.com/t/KogMzx0q7f5

    Our receiver always performs either ack() or nack(), and the occasional message that sneaks through when symptoms look like this complete within 1s, usually faster.

    I deleted the subscription and recreated it, and saw the modack calls drop to zero, only to climb back almost instantly to where they were before.

    Is there anything else that will help you/us troubleshoot this issue?

  12. kir-titievsky commented on Dec 11, 2017

    @kir-titievsky
    Author
  13. 24 remaining items

  14. kir-titievsky commented on Jun 14, 2018

    @kir-titievsky
    Author

    That sounds like a bug somewhere. The re-delivery after an hour, in particular, makes me suspicious. If you could, might you file a separate bug with a reproduction? If not,

    1. Might you open a support case with GCP support, if you have a support plan, detailing when you observed this on what project and subscription?
    2. If not, could you send the same to cloud-pubsub@google.com? No guarantees, but I might be able to take a look.

    An alternative explanation for this behavior is that the acks never succeed. Which might make this a client-side bug. But hard to tell.

  15. luann-trend commented on Jun 21, 2018

    @luann-trend

    We have experienced similar kind of issue here with the java google-cloud-pubsub lib GA version 1.31.0. We don't get the duplicate messages but the messages seem stuck in the queue even though we send ack back. After restarts the clients, the stuck messages got cleared up.

  16. pongad commented on Jun 21, 2018

    @pongad
    Contributor

    @luann-trend What does "stuck" here mean? Are new messages not being processed?

  17. luann-trend commented on Jun 22, 2018

    @luann-trend

    @pongad The new messages still being processed, only there couple hundreds message keep get redelivered but not able to process for some reason. We have experienced the same issues on 2 different cluster environments after about 4-5 days start using the new Java Google-pubsub client GA version.

  18. pongad commented on Jun 22, 2018

    @pongad
    Contributor

    @luann-trend This is interesting. Is it possible that the messages are causing you to throw exception? We catch exception and nack messages automatically, assuming that user code failed to process it. Do you know the duration of time between redelivers?

  19. chingor13 commented on Dec 4, 2018

    @chingor13
    Contributor

    Should have been fixed in #3743

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

🚨 criticalP0 critical issue. Requires immediate fixapi: pubsubIssues related to the Pub/Sub API.priority: p2Moderately-important priority. Fix may not be included in next release.status: blockedResolving the issue is dependent on other work.triaged for GAtype: bugError or flaw in code with unintended results or allowing sub-optimal usage patterns.

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions