{"record":{"id":"62fa266abd381be7","repo":"apache/seatunnel","slug":"acknowledge-failed-62fa26","errorCode":"ACKNOWLEDGE_FAILED","errorMessage":"Google Pub/Sub returned ${response} while acknowledging checkpoint ${checkpointId}","messagePattern":"Google Pub/Sub returned (.+?) while acknowledging checkpoint (.+?)","errorType":"error_code","errorClass":"GooglePubSubConnectorException","httpStatus":null,"severity":"error","filePath":"seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReader.java","lineNumber":203,"sourceCode":"\n    @Override\n    public void close() throws IOException {\n        if (subscriber != null) {\n            subscriber.close();\n        }\n    }\n\n    private void acknowledge(\n            List<AckReplyConsumerWithResponse> acknowledgements, long checkpointId) {\n        List<ApiFuture<AckResponse>> futures = new ArrayList<>(acknowledgements.size());\n        for (AckReplyConsumerWithResponse acknowledgement : acknowledgements) {\n            futures.add(acknowledgement.ack());\n        }\n\n        try {\n            for (AckResponse response : ApiFutures.allAsList(futures).get()) {\n                if (response != AckResponse.SUCCESSFUL) {\n                    throw new GooglePubSubConnectorException(\n                            GooglePubSubConnectorErrorCode.ACKNOWLEDGE_FAILED,\n                            \"Google Pub/Sub returned \"\n                                    + response\n                                    + \" while acknowledging checkpoint \"\n                                    + checkpointId);\n                }\n            }\n        } catch (InterruptedException e) {\n            Thread.currentThread().interrupt();\n            throw acknowledgeFailure(checkpointId, e);\n        } catch (ExecutionException e) {\n            throw acknowledgeFailure(checkpointId, e.getCause());\n        }\n    }\n\n    private void checkSubscriberFailure() {\n        Throwable failure = subscriberFailure.get();\n        if (failure != null) {","sourceCodeStart":185,"sourceCodeEnd":221,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReader.java#L185-L221","documentation":"Thrown in GooglePubSubSourceReader.acknowledge (invoked from notifyCheckpointComplete) after awaiting all ack futures: if any AckResponse is not SUCCESSFUL, the connector raises ACKNOWLEDGE_FAILED with the actual response and the checkpoint id. This means Pub/Sub did not confirm acknowledgement for some messages, so they may be redelivered after ack deadline expiry.","triggerScenarios":"notifyCheckpointComplete triggers acknowledge(checkpointId); the accumulated ack() futures resolve to a non-SUCCESSFUL AckResponse (e.g. FAILED, INVALID, PERMISSION_DENIED) for at least one message.","commonSituations":"Messages whose ack IDs expired (ack deadline passed before checkpoint completed — very long checkpoints); permission changes removing pubsub.subscriber on the service account; temporary Pub/Sub outages during checkpoint completion; acknowledging after subscription recreation.","solutions":["Shorten checkpoint.interval so messages are acknowledged within the Pub/Sub ack deadline (default 10 minutes), or raise the subscription's ackDeadline.","Verify the service account still has pubsub.subscriber (and pubsub.viewer) on the subscription.","Check GCP status/incidents for Pub/Sub at the failure time and rely on redelivery for the affected messages.","Ensure the subscription was not deleted/recreated between receive and ack; recreate and restart the job if it was."],"exampleFix":"// before\nenv {\n  checkpoint.interval = 3600000 # 1h, exceeds ack deadline\n}\n// after\nenv {\n  checkpoint.interval = 60000\n}","handlingStrategy":"retry","validationCode":"// Pre-submit check\nlong ackDeadlineSeconds = subscriptionAckDeadline; // fetch via SubscriptionAdminClient\nif (checkpointIntervalMs > ackDeadlineSeconds * 1000) {\n    throw new IllegalStateException(\"checkpoint.interval exceeds subscription ackDeadline\");\n}","typeGuard":null,"tryCatchPattern":"try {\n    notifyCheckpointComplete(checkpointId);\n} catch (GooglePubSubConnectorException e) {\n    if (e.getErrorCode() == GooglePubSubConnectorErrorCode.ACKNOWLEDGE_FAILED) {\n        logger.warn(\"Ack failed at checkpoint {} ({}); messages will be redelivered\", checkpointId, e.getMessage());\n        // rely on Pub/Sub redelivery; optionally retry ack for transient responses\n    }\n}","preventionTips":["Keep checkpoint.interval well below the subscription ackDeadline.","Grant and periodically audit pubsub.subscriber permissions for the service account.","Monitor Pub/Sub unacknowledged-message metrics to catch ack failures early.","Don't delete/recreate subscriptions while jobs are running."],"tags":["google-pubsub","acknowledgement","checkpoint","java"],"backgroundTag":"unexpected-api-response-shape","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}