{"record":{"id":"414a3fa0e802a519","repo":"apache/seatunnel","slug":"write-failed-414a3f","errorCode":"WRITE_FAILED","errorMessage":"Failed to publish message to Google Pub/Sub","messagePattern":"Failed to publish message to Google Pub/Sub","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/sink/GooglePubSubSinkWriter.java","lineNumber":74,"sourceCode":"    }\n\n    public GooglePubSubSinkWriter(SeaTunnelRowType rowType, GooglePubSubSinkConfig config) {\n        this(rowType, config, GooglePubSubPublisher.create(config));\n    }\n\n    @Override\n    public void write(SeaTunnelRow row) throws IOException {\n        checkPublishError();\n\n        PubsubMessage message =\n                PubsubMessage.newBuilder()\n                        .setData(ByteString.copyFrom(serializationSchema.serialize(row)))\n                        .build();\n        ApiFuture<String> publishFuture;\n        try {\n            publishFuture = publisher.publish(message);\n        } catch (RuntimeException e) {\n            throw writeFailure(e);\n        }\n\n        pendingPublishes.add(publishFuture);\n        ApiFutures.addCallback(\n                publishFuture,\n                new ApiFutureCallback<String>() {\n                    @Override\n                    public void onSuccess(String messageId) {\n                        pendingPublishes.remove(publishFuture);\n                    }\n\n                    @Override\n                    public void onFailure(Throwable throwable) {\n                        publishError.compareAndSet(null, throwable);\n                        pendingPublishes.remove(publishFuture);\n                    }\n                },\n                Runnable::run);","sourceCodeStart":56,"sourceCodeEnd":92,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/sink/GooglePubSubSinkWriter.java#L56-L92","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before\npublisher = Publisher.newBuilder(\"projects/wrong-project/topics/my-topic\").build(); // permission/topic error\n// after\npublisher = Publisher.newBuilder(\"projects/my-project/topics/my-topic\").build();","handlingStrategy":"try-catch","validationCode":"// preflight: publish a test message and delete it\ntry (Publisher p = Publisher.newBuilder(topicName).build()) {\n    p.publish(PubsubMessage.newBuilder().setData(ByteString.EMPTY).build()).get();\n}","typeGuard":null,"tryCatchPattern":"try {\n    publishFuture = publisher.publish(message);\n} catch (SeaTunnelRuntimeException e) {\n    // check credentials/topic permissions; do not swallow\n    if (publisher.isShutdown()) { /* writer reused after close — fix lifecycle */ }\n    throw e;\n}","preventionTips":["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"],"tags":["pubsub","google-cloud","sink","write"],"backgroundTag":"http-request-failed","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"}