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
- 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.
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
- 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.
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
- CONFIGURATION_FAILED
- CONNECTION_FAILED
- Failed to close Google Pub/Sub publisher
- Google Pub/Sub source expects exactly one source split
- Option 'field_delimiter' cannot be empty
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)