apache/seatunnel · error · GooglePubSubConnectorException

ACKNOWLEDGE_FAILED

ACKNOWLEDGE_FAILED

Error message

Google Pub/Sub returned ${response} while acknowledging checkpoint ${checkpointId}

What it means

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.

Solutions

  1. Shorten checkpoint.interval so messages are acknowledged within the Pub/Sub ack deadline (default 10 minutes), or raise the subscription's ackDeadline.
  2. Verify the service account still has pubsub.subscriber (and pubsub.viewer) on the subscription.
  3. Check GCP status/incidents for Pub/Sub at the failure time and rely on redelivery for the affected messages.
  4. Ensure the subscription was not deleted/recreated between receive and ack; recreate and restart the job if it was.

Example fix

// before
env {
  checkpoint.interval = 3600000 # 1h, exceeds ack deadline
}
// after
env {
  checkpoint.interval = 60000
}
Defensive patterns

Strategy: retry

Validate before calling

// Pre-submit check
long ackDeadlineSeconds = subscriptionAckDeadline; // fetch via SubscriptionAdminClient
if (checkpointIntervalMs > ackDeadlineSeconds * 1000) {
    throw new IllegalStateException("checkpoint.interval exceeds subscription ackDeadline");
}

Try / catch

try {
    notifyCheckpointComplete(checkpointId);
} catch (GooglePubSubConnectorException e) {
    if (e.getErrorCode() == GooglePubSubConnectorErrorCode.ACKNOWLEDGE_FAILED) {
        logger.warn("Ack failed at checkpoint {} ({}); messages will be redelivered", checkpointId, e.getMessage());
        // rely on Pub/Sub redelivery; optionally retry ack for transient responses
    }
}

Prevention

When it happens

Trigger: 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.

Common situations: 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.

Understand the failure class

Background: "invalid response format", "malformed payload", "missing data field": when an API returns 200 but the response shape is wrong — this error's family across 23 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/62fa266abd381be7. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReader.java:203

    @Override
    public void close() throws IOException {
        if (subscriber != null) {
            subscriber.close();
        }
    }

    private void acknowledge(
            List<AckReplyConsumerWithResponse> acknowledgements, long checkpointId) {
        List<ApiFuture<AckResponse>> futures = new ArrayList<>(acknowledgements.size());
        for (AckReplyConsumerWithResponse acknowledgement : acknowledgements) {
            futures.add(acknowledgement.ack());
        }

        try {
            for (AckResponse response : ApiFutures.allAsList(futures).get()) {
                if (response != AckResponse.SUCCESSFUL) {
                    throw new GooglePubSubConnectorException(
                            GooglePubSubConnectorErrorCode.ACKNOWLEDGE_FAILED,
                            "Google Pub/Sub returned "
                                    + response
                                    + " while acknowledging checkpoint "
                                    + checkpointId);
                }
            }
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw acknowledgeFailure(checkpointId, e);
        } catch (ExecutionException e) {
            throw acknowledgeFailure(checkpointId, e.getCause());
        }
    }

    private void checkSubscriberFailure() {
        Throwable failure = subscriberFailure.get();
        if (failure != null) {

View on GitHub (pinned to cf67b549a7)