apache/seatunnel · error · GooglePubSubConnectorException
WRITE_FAILED
WRITE_FAILED
Error message
Failed to publish message to Google Pub/Sub
What it means
GooglePubSubSinkWriter.write serializes each row and calls publisher.publish(message). A synchronous RuntimeException from the Pub/Sub client (publisher shut down, immediate auth/transport failure) is rethrown as a SeaTunnelRuntimeException with code WRITE_FAILED: 'Failed to publish message to Google Pub/Sub'. Async failures on the ApiFuture are surfaced separately via callback.
Solutions
- Verify the service account has pubsub.publisher on the topic and credentials are present on all worker nodes
- Confirm the topic path (project/topic) in config matches an existing topic
- Check whether the writer is being closed/reused across checkpoint boundaries incorrectly; keep publish calls before close
- Retry the job for transient gRPC unavailability; enable retry settings on the Pub/Sub publisher
Example fix
// before
publisher = Publisher.newBuilder("projects/wrong-project/topics/my-topic").build(); // permission/topic error
// after
publisher = Publisher.newBuilder("projects/my-project/topics/my-topic").build(); Defensive patterns
Strategy: try-catch
Validate before calling
// preflight: publish a test message and delete it
try (Publisher p = Publisher.newBuilder(topicName).build()) {
p.publish(PubsubMessage.newBuilder().setData(ByteString.EMPTY).build()).get();
} Try / catch
try {
publishFuture = publisher.publish(message);
} catch (SeaTunnelRuntimeException e) {
// check credentials/topic permissions; do not swallow
if (publisher.isShutdown()) { /* writer reused after close — fix lifecycle */ }
throw e;
} Prevention
- Grant roles/pubsub.publisher to the service account on the exact topic
- Verify topic project/name in config and that the publisher is closed only after all writes
- Run a preflight publish at job start to fail fast on auth/permission issues
- Enable Pub/Sub client retry settings for transient gRPC errors
When it happens
Trigger: publisher.publish throwing RuntimeException at call time — publisher already closed/shut down (e.g. after prepareCommit/close races), invalid topic state, immediate credential or transport errors; also async publish failures surfaced through shouldSurfaceAsynchronousPublishFailure handling.
Common situations: Missing or invalid Google service-account credentials on worker nodes, topic deleted or permission revoked (roles/pubsub.publisher missing), project id mismatch, or checkpoint/flush racing with writer close.
Understand the failure class
Background: 'Something went wrong' / 'Request failed (500)' / 'HTTP error! status: 404' — what failed HTTP requests actually mean and how to find the real cause — this error's family across 28 libraries.
Related errors
- FLUSH_DATA_FAILED
- FLUSH_DATA_FAILED
- Option 'field_delimiter' cannot be empty
- Options 'credentials_path' and 'emulator_host' cannot be…
- SNMP agent returned error status
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/414a3fa0e802a519.
Report an issue: GitHub.
Appendix: source
Thrown at seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/sink/GooglePubSubSinkWriter.java:74
}
public GooglePubSubSinkWriter(SeaTunnelRowType rowType, GooglePubSubSinkConfig config) {
this(rowType, config, GooglePubSubPublisher.create(config));
}
@Override
public void write(SeaTunnelRow row) throws IOException {
checkPublishError();
PubsubMessage message =
PubsubMessage.newBuilder()
.setData(ByteString.copyFrom(serializationSchema.serialize(row)))
.build();
ApiFuture<String> publishFuture;
try {
publishFuture = publisher.publish(message);
} catch (RuntimeException e) {
throw writeFailure(e);
}
pendingPublishes.add(publishFuture);
ApiFutures.addCallback(
publishFuture,
new ApiFutureCallback<String>() {
@Override
public void onSuccess(String messageId) {
pendingPublishes.remove(publishFuture);
}
@Override
public void onFailure(Throwable throwable) {
publishError.compareAndSet(null, throwable);
pendingPublishes.remove(publishFuture);
}
},
Runnable::run);View on GitHub (pinned to cf67b549a7)