This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new c5b58851ba [Feature] [Connector-V2] [GooglePubSub] Add Google Pub/Sub 
source connector (#11989)
c5b58851ba is described below

commit c5b58851ba9e5aed7ea5d7cee872952616df814c
Author: Goutam Adwant <[email protected]>
AuthorDate: Sun Aug 30 04:24:39 2026 +0000

    [Feature] [Connector-V2] [GooglePubSub] Add Google Pub/Sub source connector 
(#11989)
    
    Signed-off-by: goutamadwant <[email protected]>
---
 .../changelog/connector-google-pubsub.md           |   1 +
 docs/en/connectors/source/GooglePubSub.md          | 166 +++++++++++++
 .../changelog/connector-google-pubsub.md           |   1 +
 docs/zh/connectors/source/GooglePubSub.md          | 166 +++++++++++++
 plugin-mapping.properties                          |   1 +
 .../pubsub/config/GooglePubSubSourceConfig.java    | 100 ++++++++
 .../pubsub/config/GooglePubSubSourceOptions.java   |  91 +++++++
 .../exception/GooglePubSubConnectorErrorCode.java  |   7 +-
 .../google/pubsub/source/GooglePubSubSource.java   | 110 +++++++++
 .../pubsub/source/GooglePubSubSourceFactory.java   |  98 ++++++++
 .../pubsub/source/GooglePubSubSourceReader.java    | 252 +++++++++++++++++++
 .../pubsub/source/GooglePubSubSubscriber.java      | 184 ++++++++++++++
 .../PubSubSubscriber.java}                         |  28 +--
 .../config/GooglePubSubSourceConfigTest.java       | 106 ++++++++
 .../source/GooglePubSubSourceFactoryTest.java}     |  31 +--
 .../source/GooglePubSubSourceReaderTest.java       | 273 +++++++++++++++++++++
 .../pubsub/source/GooglePubSubSourceTest.java      |  72 ++++++
 .../connector/google/pubsub/GooglePubSubIT.java    |  97 ++++++++
 .../resources/pubsub/google_pubsub_to_console.conf |  41 ++++
 19 files changed, 1779 insertions(+), 46 deletions(-)

diff --git a/docs/en/connectors/changelog/connector-google-pubsub.md 
b/docs/en/connectors/changelog/connector-google-pubsub.md
index 67f31cd76e..663a6ed8fd 100644
--- a/docs/en/connectors/changelog/connector-google-pubsub.md
+++ b/docs/en/connectors/changelog/connector-google-pubsub.md
@@ -2,6 +2,7 @@
 
 | Change | Commit | Version |
 | --- | --- | --- |
+|[Feature][Connector-V2] Add Google Pub/Sub source connector|-|Next|
 |[Feature][Connector-V2] Add Google Pub/Sub sink connector|-|Next|
 
 </details>
diff --git a/docs/en/connectors/source/GooglePubSub.md 
b/docs/en/connectors/source/GooglePubSub.md
new file mode 100644
index 0000000000..5637bffeba
--- /dev/null
+++ b/docs/en/connectors/source/GooglePubSub.md
@@ -0,0 +1,166 @@
+import ChangeLog from '../changelog/connector-google-pubsub.md';
+
+# GooglePubSub
+
+> Google Pub/Sub source connector
+
+## Description
+
+Reads messages from an existing Google Pub/Sub subscription and converts each 
message payload to a SeaTunnel row.
+
+## Support Those Engines
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## Key Features
+
+- [ ] [batch](../../introduction/concepts/connector-v2-features.md)
+- [x] [stream](../../introduction/concepts/connector-v2-features.md)
+- [ ] [exactly-once](../../introduction/concepts/connector-v2-features.md)
+- [ ] [column projection](../../introduction/concepts/connector-v2-features.md)
+- [ ] [parallelism](../../introduction/concepts/connector-v2-features.md)
+- [ ] [support user-defined 
split](../../introduction/concepts/connector-v2-features.md)
+
+## Options
+
+| name | type | required | default value |
+| --- | --- | --- | --- |
+| project_id | string | yes | - |
+| subscription | string | yes | - |
+| credentials_path | string | no | - |
+| emulator_host | string | no | - |
+| format | enum | no | json |
+| field_delimiter | string | no | , |
+| max_outstanding_messages | long | no | Google client default |
+| max_outstanding_bytes | long | no | Google client default |
+| parallel_pull_count | int | no | Google client default |
+| schema | config | yes | - |
+| common-options | | no | - |
+
+### project_id [string]
+
+Google Cloud project ID that owns the subscription.
+
+### subscription [string]
+
+Pub/Sub subscription ID. The subscription and its topic must exist before the 
job starts.
+
+### credentials_path [string]
+
+Path to a Google Cloud service account JSON key file. If this option is not 
set, the connector uses [Application Default 
Credentials](https://cloud.google.com/docs/authentication/application-default-credentials).
+
+### emulator_host [string]
+
+Pub/Sub emulator host and port, for example `pubsub-emulator:8085`. When set, 
the connector uses a plaintext connection without credentials. Do not use this 
option for a production Pub/Sub endpoint.
+
+### format [enum]
+
+Message payload format. Supported values:
+
+- `json`: converts a JSON object to a row using the configured schema.
+- `text`: splits the payload into fields using `field_delimiter`.
+
+### field_delimiter [string]
+
+Field delimiter used when `format = text`. The default is `,`.
+
+### max_outstanding_messages [long]
+
+Maximum number of messages the subscriber can hold before applying flow 
control. The value must be greater than `0`. When omitted, the Google client 
default is used.
+
+### max_outstanding_bytes [long]
+
+Maximum total message bytes the subscriber can hold before applying flow 
control. The value must be greater than `0`. When omitted, the Google client 
default is used.
+
+### parallel_pull_count [int]
+
+Number of streaming pull connections opened by each source reader. The value 
must be greater than `0`. When omitted, the Google client default is used.
+
+### schema [config]
+
+Schema used to deserialize message payloads. See [Schema 
Feature](../../introduction/concepts/schema-feature.md) for details.
+
+### common options
+
+Source plugin common parameters, please refer to [Source Common 
Options](../common-options/source-common-options.md) for details.
+
+## Delivery Semantics
+
+The connector uses one logical Pub/Sub subscription split. Messages are 
acknowledged only after the SeaTunnel checkpoint containing their rows 
completes. If the task fails before that checkpoint completes, Pub/Sub can 
redeliver the unacknowledged messages.
+
+This provides at-least-once delivery. Consumers must tolerate duplicate rows 
after recovery. Periodic SeaTunnel checkpoints must be enabled so the connector 
can acknowledge processed messages. The source currently does not expose 
Pub/Sub message attributes, ordering keys, or publish timestamps as metadata 
fields.
+
+If a message cannot be deserialized, the connector negatively acknowledges it 
and fails the source task. Pub/Sub can redeliver the same message after 
recovery, so a permanently invalid message can repeatedly restart the job. 
Configure a Pub/Sub dead-letter topic or remove the invalid message when this 
behavior is not acceptable.
+
+## Task Example
+
+### Application Default Credentials
+
+```hocon
+env {
+  parallelism = 1
+  job.mode = "STREAMING"
+  checkpoint.interval = 5000
+}
+
+source {
+  GooglePubSub {
+    project_id = "my-gcp-project"
+    subscription = "events-subscription"
+    format = json
+    schema = {
+      fields {
+        event_id = string
+        event_type = string
+      }
+    }
+  }
+}
+
+sink {
+  Console {}
+}
+```
+
+### Service Account Key File
+
+```hocon
+source {
+  GooglePubSub {
+    project_id = "my-gcp-project"
+    subscription = "events-subscription"
+    credentials_path = "/secrets/service-account.json"
+    format = text
+    field_delimiter = "|"
+    schema = {
+      fields {
+        event_id = string
+        event_type = string
+      }
+    }
+  }
+}
+```
+
+### Pub/Sub Emulator
+
+```hocon
+source {
+  GooglePubSub {
+    project_id = "local-project"
+    subscription = "events-subscription"
+    emulator_host = "pubsub-emulator:8085"
+    schema = {
+      fields {
+        event_id = string
+      }
+    }
+  }
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/docs/zh/connectors/changelog/connector-google-pubsub.md 
b/docs/zh/connectors/changelog/connector-google-pubsub.md
index 67f31cd76e..663a6ed8fd 100644
--- a/docs/zh/connectors/changelog/connector-google-pubsub.md
+++ b/docs/zh/connectors/changelog/connector-google-pubsub.md
@@ -2,6 +2,7 @@
 
 | Change | Commit | Version |
 | --- | --- | --- |
+|[Feature][Connector-V2] Add Google Pub/Sub source connector|-|Next|
 |[Feature][Connector-V2] Add Google Pub/Sub sink connector|-|Next|
 
 </details>
diff --git a/docs/zh/connectors/source/GooglePubSub.md 
b/docs/zh/connectors/source/GooglePubSub.md
new file mode 100644
index 0000000000..23d66a198d
--- /dev/null
+++ b/docs/zh/connectors/source/GooglePubSub.md
@@ -0,0 +1,166 @@
+import ChangeLog from '../changelog/connector-google-pubsub.md';
+
+# GooglePubSub
+
+> Google Pub/Sub Source 连接器
+
+## 描述
+
+从已有的 Google Pub/Sub 订阅读取消息,并将每条消息负载转换为 SeaTunnel 行。
+
+## 支持这些引擎
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## 主要特性
+
+- [ ] [批处理](../../introduction/concepts/connector-v2-features.md)
+- [x] [流处理](../../introduction/concepts/connector-v2-features.md)
+- [ ] [精确一次](../../introduction/concepts/connector-v2-features.md)
+- [ ] [列投影](../../introduction/concepts/connector-v2-features.md)
+- [ ] [并行度](../../introduction/concepts/connector-v2-features.md)
+- [ ] [支持用户自定义 Split](../../introduction/concepts/connector-v2-features.md)
+
+## 参数
+
+| 参数名 | 类型 | 是否必填 | 默认值 |
+| --- | --- | --- | --- |
+| project_id | string | 是 | - |
+| subscription | string | 是 | - |
+| credentials_path | string | 否 | - |
+| emulator_host | string | 否 | - |
+| format | enum | 否 | json |
+| field_delimiter | string | 否 | , |
+| max_outstanding_messages | long | 否 | Google 客户端默认值 |
+| max_outstanding_bytes | long | 否 | Google 客户端默认值 |
+| parallel_pull_count | int | 否 | Google 客户端默认值 |
+| schema | config | 是 | - |
+| common-options | | 否 | - |
+
+### project_id [string]
+
+订阅所属的 Google Cloud 项目 ID。
+
+### subscription [string]
+
+Pub/Sub 订阅 ID。启动作业前必须先创建该订阅及其关联主题。
+
+### credentials_path [string]
+
+Google Cloud 服务账号 JSON 密钥文件的路径。未配置时,连接器使用 [Application Default 
Credentials](https://cloud.google.com/docs/authentication/application-default-credentials)。
+
+### emulator_host [string]
+
+Pub/Sub 模拟器的主机和端口,例如 `pubsub-emulator:8085`。配置后,连接器使用无凭证的明文连接。生产环境中不要使用该选项。
+
+### format [enum]
+
+消息负载格式。支持以下值:
+
+- `json`:按照配置的 Schema 将 JSON 对象转换为行。
+- `text`:使用 `field_delimiter` 将负载拆分为字段。
+
+### field_delimiter [string]
+
+`format = text` 时使用的字段分隔符。默认值为 `,`。
+
+### max_outstanding_messages [long]
+
+订阅客户端在触发流量控制前最多保留的消息数。该值必须大于 `0`。未配置时使用 Google 客户端默认值。
+
+### max_outstanding_bytes [long]
+
+订阅客户端在触发流量控制前最多保留的消息总字节数。该值必须大于 `0`。未配置时使用 Google 客户端默认值。
+
+### parallel_pull_count [int]
+
+每个 Source Reader 建立的流式拉取连接数。该值必须大于 `0`。未配置时使用 Google 客户端默认值。
+
+### schema [config]
+
+反序列化消息负载使用的 Schema。详情请参阅 [Schema 
特性](../../introduction/concepts/schema-feature.md)。
+
+### common options
+
+Source 插件通用参数,请参考 [Source 通用选项](../common-options/source-common-options.md)。
+
+## 交付语义
+
+连接器使用一个逻辑 Pub/Sub 订阅 Split。只有当包含消息对应行的 SeaTunnel 
检查点完成后,连接器才确认这些消息。如果任务在检查点完成前失败,Pub/Sub 可以重新投递未确认的消息。
+
+该机制提供至少一次交付语义。恢复后可能出现重复行,使用方需要具备去重能力。必须启用周期性 SeaTunnel 
检查点,连接器才能确认已处理的消息。当前版本不将 Pub/Sub 消息属性、排序键或发布时间公开为元数据字段。
+
+如果消息无法反序列化,连接器会对该消息进行否定确认并使 Source 任务失败。Pub/Sub 
可以在任务恢复后再次投递同一条消息,因此永久无效的消息可能导致作业反复重启。如不能接受该行为,请配置 Pub/Sub 死信主题或移除无效消息。
+
+## 任务示例
+
+### Application Default Credentials
+
+```hocon
+env {
+  parallelism = 1
+  job.mode = "STREAMING"
+  checkpoint.interval = 5000
+}
+
+source {
+  GooglePubSub {
+    project_id = "my-gcp-project"
+    subscription = "events-subscription"
+    format = json
+    schema = {
+      fields {
+        event_id = string
+        event_type = string
+      }
+    }
+  }
+}
+
+sink {
+  Console {}
+}
+```
+
+### 服务账号密钥文件
+
+```hocon
+source {
+  GooglePubSub {
+    project_id = "my-gcp-project"
+    subscription = "events-subscription"
+    credentials_path = "/secrets/service-account.json"
+    format = text
+    field_delimiter = "|"
+    schema = {
+      fields {
+        event_id = string
+        event_type = string
+      }
+    }
+  }
+}
+```
+
+### Pub/Sub 模拟器
+
+```hocon
+source {
+  GooglePubSub {
+    project_id = "local-project"
+    subscription = "events-subscription"
+    emulator_host = "pubsub-emulator:8085"
+    schema = {
+      fields {
+        event_id = string
+      }
+    }
+  }
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/plugin-mapping.properties b/plugin-mapping.properties
index 83ce19f65f..066de0c6ce 100644
--- a/plugin-mapping.properties
+++ b/plugin-mapping.properties
@@ -92,6 +92,7 @@ seatunnel.sink.InfluxDB = connector-influxdb
 seatunnel.source.GoogleSheets = connector-google-sheets
 seatunnel.sink.GoogleFirestore = connector-google-firestore
 seatunnel.sink.GooglePubSub = connector-google-pubsub
+seatunnel.source.GooglePubSub = connector-google-pubsub
 seatunnel.sink.AzureQueueStorage = connector-azure-queue-storage
 seatunnel.source.GoogleBigtable = connector-google-bigtable
 seatunnel.sink.GoogleBigtable = connector-google-bigtable
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceConfig.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceConfig.java
new file mode 100644
index 0000000000..aff187a7d5
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceConfig.java
@@ -0,0 +1,100 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+
+import lombok.Builder;
+import lombok.Getter;
+
+import java.io.Serializable;
+
+/** Immutable runtime configuration for the Google Pub/Sub source. */
+@Getter
+@Builder
+public class GooglePubSubSourceConfig implements Serializable {
+
+    private final String projectId;
+    private final String subscription;
+    private final String credentialsPath;
+    private final String emulatorHost;
+    private final MessageFormat format;
+    private final String fieldDelimiter;
+    private final Long maxOutstandingMessages;
+    private final Long maxOutstandingBytes;
+    private final Integer parallelPullCount;
+
+    public static GooglePubSubSourceConfig from(ReadonlyConfig config) {
+        String projectId = config.get(GooglePubSubSourceOptions.PROJECT_ID);
+        String subscription = 
config.get(GooglePubSubSourceOptions.SUBSCRIPTION);
+        String credentialsPath = 
config.get(GooglePubSubSourceOptions.CREDENTIALS_PATH);
+        String emulatorHost = 
config.get(GooglePubSubSourceOptions.EMULATOR_HOST);
+        MessageFormat format = config.get(GooglePubSubSourceOptions.FORMAT);
+        String fieldDelimiter = 
config.get(GooglePubSubSourceOptions.FIELD_DELIMITER);
+        Long maxOutstandingMessages =
+                config.get(GooglePubSubSourceOptions.MAX_OUTSTANDING_MESSAGES);
+        Long maxOutstandingBytes = 
config.get(GooglePubSubSourceOptions.MAX_OUTSTANDING_BYTES);
+        Integer parallelPullCount = 
config.get(GooglePubSubSourceOptions.PARALLEL_PULL_COUNT);
+
+        requireNonBlank(projectId, GooglePubSubSourceOptions.PROJECT_ID.key());
+        requireNonBlank(subscription, 
GooglePubSubSourceOptions.SUBSCRIPTION.key());
+        requireNonBlankIfPresent(credentialsPath, 
GooglePubSubSourceOptions.CREDENTIALS_PATH.key());
+        requireNonBlankIfPresent(emulatorHost, 
GooglePubSubSourceOptions.EMULATOR_HOST.key());
+        if (credentialsPath != null && emulatorHost != null) {
+            throw new IllegalArgumentException(
+                    "Options 'credentials_path' and 'emulator_host' cannot be 
configured together");
+        }
+        if (format == MessageFormat.TEXT && fieldDelimiter.isEmpty()) {
+            throw new IllegalArgumentException("Option 'field_delimiter' 
cannot be empty");
+        }
+        requirePositive(
+                maxOutstandingMessages, 
GooglePubSubSourceOptions.MAX_OUTSTANDING_MESSAGES.key());
+        requirePositive(maxOutstandingBytes, 
GooglePubSubSourceOptions.MAX_OUTSTANDING_BYTES.key());
+        requirePositive(parallelPullCount, 
GooglePubSubSourceOptions.PARALLEL_PULL_COUNT.key());
+
+        return GooglePubSubSourceConfig.builder()
+                .projectId(projectId)
+                .subscription(subscription)
+                .credentialsPath(credentialsPath)
+                .emulatorHost(emulatorHost)
+                .format(format)
+                .fieldDelimiter(fieldDelimiter)
+                .maxOutstandingMessages(maxOutstandingMessages)
+                .maxOutstandingBytes(maxOutstandingBytes)
+                .parallelPullCount(parallelPullCount)
+                .build();
+    }
+
+    private static void requireNonBlank(String value, String option) {
+        if (value == null || value.trim().isEmpty()) {
+            throw new IllegalArgumentException("Option '" + option + "' cannot 
be blank");
+        }
+    }
+
+    private static void requireNonBlankIfPresent(String value, String option) {
+        if (value != null) {
+            requireNonBlank(value, option);
+        }
+    }
+
+    private static void requirePositive(Number value, String option) {
+        if (value != null && value.longValue() <= 0) {
+            throw new IllegalArgumentException("Option '" + option + "' must 
be greater than 0");
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceOptions.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceOptions.java
new file mode 100644
index 0000000000..c1e7e58269
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceOptions.java
@@ -0,0 +1,91 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.config;
+
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.Options;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+
+public class GooglePubSubSourceOptions extends ConnectorCommonOptions {
+
+    public static final Option<String> PROJECT_ID =
+            Options.key("project_id")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("Google Cloud project ID.");
+
+    public static final Option<String> SUBSCRIPTION =
+            Options.key("subscription")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("Google Pub/Sub subscription ID.");
+
+    public static final Option<String> CREDENTIALS_PATH =
+            Options.key("credentials_path")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "Path to a Google Cloud service account JSON key 
file. "
+                                    + "Application Default Credentials are 
used when this option is not set.");
+
+    public static final Option<String> EMULATOR_HOST =
+            Options.key("emulator_host")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "Pub/Sub emulator host and port, for example 
pubsub-emulator:8085. "
+                                    + "Authentication and TLS are disabled 
when this option is set.");
+
+    public static final Option<MessageFormat> FORMAT =
+            Options.key("format")
+                    .enumType(MessageFormat.class)
+                    .defaultValue(MessageFormat.JSON)
+                    .withDescription("Message payload format. Supported values 
are json and text.");
+
+    public static final Option<String> FIELD_DELIMITER =
+            Options.key("field_delimiter")
+                    .stringType()
+                    .defaultValue(",")
+                    .withDescription("Field delimiter used when format is 
text.");
+
+    public static final Option<Long> MAX_OUTSTANDING_MESSAGES =
+            Options.key("max_outstanding_messages")
+                    .longType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "Maximum number of messages held by the Pub/Sub 
subscriber. "
+                                    + "The Google client default is used when 
this option is not set.");
+
+    public static final Option<Long> MAX_OUTSTANDING_BYTES =
+            Options.key("max_outstanding_bytes")
+                    .longType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "Maximum total bytes held by the Pub/Sub 
subscriber. "
+                                    + "The Google client default is used when 
this option is not set.");
+
+    public static final Option<Integer> PARALLEL_PULL_COUNT =
+            Options.key("parallel_pull_count")
+                    .intType()
+                    .noDefaultValue()
+                    .withDescription(
+                            "Number of streaming pull connections used by the 
Pub/Sub subscriber. "
+                                    + "The Google client default is used when 
this option is not set.");
+
+    private GooglePubSubSourceOptions() {}
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
index f4e8fcbc8a..46bae412d9 100644
--- 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
@@ -20,9 +20,12 @@ package 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception;
 import org.apache.seatunnel.common.exception.SeaTunnelErrorCode;
 
 public enum GooglePubSubConnectorErrorCode implements SeaTunnelErrorCode {
-    CONNECTION_FAILED("GooglePubSub-01", "Create Google Pub/Sub publisher 
failed"),
+    CONNECTION_FAILED("GooglePubSub-01", "Create Google Pub/Sub client 
failed"),
     WRITE_FAILED("GooglePubSub-02", "Publish Google Pub/Sub message failed"),
-    CLOSE_FAILED("GooglePubSub-03", "Close Google Pub/Sub publisher failed");
+    CLOSE_FAILED("GooglePubSub-03", "Close Google Pub/Sub client failed"),
+    READ_FAILED("GooglePubSub-04", "Read Google Pub/Sub message failed"),
+    ACKNOWLEDGE_FAILED("GooglePubSub-05", "Acknowledge Google Pub/Sub message 
failed"),
+    CONFIGURATION_FAILED("GooglePubSub-06", "Validate Google Pub/Sub 
configuration failed");
 
     private final String code;
     private final String description;
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSource.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSource.java
new file mode 100644
index 0000000000..e326e2d6f5
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSource.java
@@ -0,0 +1,110 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.common.constants.JobMode;
+import org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplit;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitEnumerator;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitEnumeratorState;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.GooglePubSubSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorErrorCode;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorException;
+
+import java.util.Collections;
+import java.util.List;
+
+/** Unbounded source for one Google Pub/Sub subscription. */
+public class GooglePubSubSource
+        implements SeaTunnelSource<SeaTunnelRow, SingleSplit, 
SingleSplitEnumeratorState> {
+
+    public static final String PLUGIN_NAME = "GooglePubSub";
+
+    private final GooglePubSubSourceConfig config;
+    private final CatalogTable catalogTable;
+    private final DeserializationSchema<SeaTunnelRow> deserializationSchema;
+
+    private JobContext jobContext;
+
+    public GooglePubSubSource(
+            GooglePubSubSourceConfig config,
+            CatalogTable catalogTable,
+            DeserializationSchema<SeaTunnelRow> deserializationSchema) {
+        this.config = config;
+        this.catalogTable = catalogTable;
+        this.deserializationSchema = deserializationSchema;
+    }
+
+    @Override
+    public Boundedness getBoundedness() {
+        if (jobContext != null) {
+            if (!JobMode.STREAMING.equals(jobContext.getJobMode())) {
+                throw new GooglePubSubConnectorException(
+                        GooglePubSubConnectorErrorCode.CONFIGURATION_FAILED,
+                        "Google Pub/Sub source supports streaming jobs only");
+            }
+            if (!jobContext.isEnableCheckpoint()) {
+                throw new GooglePubSubConnectorException(
+                        GooglePubSubConnectorErrorCode.CONFIGURATION_FAILED,
+                        "Google Pub/Sub source requires checkpointing to 
acknowledge messages");
+            }
+        }
+        return Boundedness.UNBOUNDED;
+    }
+
+    @Override
+    public String getPluginName() {
+        return PLUGIN_NAME;
+    }
+
+    @Override
+    public List<CatalogTable> getProducedCatalogTables() {
+        return Collections.singletonList(catalogTable);
+    }
+
+    @Override
+    public SourceReader<SeaTunnelRow, SingleSplit> createReader(
+            SourceReader.Context readerContext) {
+        return new GooglePubSubSourceReader(config, deserializationSchema);
+    }
+
+    @Override
+    public SourceSplitEnumerator<SingleSplit, SingleSplitEnumeratorState> 
createEnumerator(
+            SourceSplitEnumerator.Context<SingleSplit> enumeratorContext) {
+        return new SingleSplitEnumerator(enumeratorContext);
+    }
+
+    @Override
+    public SourceSplitEnumerator<SingleSplit, SingleSplitEnumeratorState> 
restoreEnumerator(
+            SourceSplitEnumerator.Context<SingleSplit> enumeratorContext,
+            SingleSplitEnumeratorState checkpointState) {
+        return new SingleSplitEnumerator(enumeratorContext);
+    }
+
+    @Override
+    public void setJobContext(JobContext jobContext) {
+        this.jobContext = jobContext;
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceFactory.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceFactory.java
new file mode 100644
index 0000000000..ac2336f15a
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceFactory.java
@@ -0,0 +1,98 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.connector.TableSource;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.GooglePubSubSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.GooglePubSubSourceOptions;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.MessageFormat;
+import org.apache.seatunnel.format.json.JsonDeserializationSchema;
+import org.apache.seatunnel.format.text.TextDeserializationSchema;
+
+import com.google.auto.service.AutoService;
+
+import java.io.Serializable;
+
+import static org.apache.seatunnel.api.options.ConnectorCommonOptions.SCHEMA;
+
+/** Creates Google Pub/Sub sources and their payload deserializers. */
+@AutoService(Factory.class)
+public class GooglePubSubSourceFactory implements TableSourceFactory {
+
+    @Override
+    public String factoryIdentifier() {
+        return GooglePubSubSource.PLUGIN_NAME;
+    }
+
+    @Override
+    public OptionRule optionRule() {
+        return OptionRule.builder()
+                .required(
+                        GooglePubSubSourceOptions.PROJECT_ID,
+                        GooglePubSubSourceOptions.SUBSCRIPTION,
+                        SCHEMA)
+                .optional(
+                        GooglePubSubSourceOptions.CREDENTIALS_PATH,
+                        GooglePubSubSourceOptions.EMULATOR_HOST,
+                        GooglePubSubSourceOptions.FORMAT,
+                        GooglePubSubSourceOptions.FIELD_DELIMITER,
+                        GooglePubSubSourceOptions.MAX_OUTSTANDING_MESSAGES,
+                        GooglePubSubSourceOptions.MAX_OUTSTANDING_BYTES,
+                        GooglePubSubSourceOptions.PARALLEL_PULL_COUNT)
+                .build();
+    }
+
+    @Override
+    public <T, SplitT extends SourceSplit, StateT extends Serializable>
+            TableSource<T, SplitT, StateT> 
createSource(TableSourceFactoryContext context) {
+        GooglePubSubSourceConfig config = 
GooglePubSubSourceConfig.from(context.getOptions());
+        CatalogTable catalogTable = 
CatalogTableUtil.buildWithConfig(context.getOptions());
+        DeserializationSchema<SeaTunnelRow> deserializationSchema =
+                createDeserializationSchema(catalogTable, config);
+
+        return () ->
+                (SeaTunnelSource<T, SplitT, StateT>)
+                        new GooglePubSubSource(config, catalogTable, 
deserializationSchema);
+    }
+
+    @Override
+    public Class<? extends SeaTunnelSource> getSourceClass() {
+        return GooglePubSubSource.class;
+    }
+
+    private DeserializationSchema<SeaTunnelRow> createDeserializationSchema(
+            CatalogTable catalogTable, GooglePubSubSourceConfig config) {
+        if (config.getFormat() == MessageFormat.JSON) {
+            return new JsonDeserializationSchema(catalogTable, false, false);
+        }
+        return TextDeserializationSchema.builder()
+                .seaTunnelRowType(catalogTable.getSeaTunnelRowType())
+                .delimiter(config.getFieldDelimiter())
+                .build();
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReader.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReader.java
new file mode 100644
index 0000000000..6fdb907d3d
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReader.java
@@ -0,0 +1,252 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
+
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.source.SourceReader;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplit;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.GooglePubSubSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorErrorCode;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorException;
+
+import com.google.api.core.ApiFuture;
+import com.google.api.core.ApiFutures;
+import com.google.cloud.pubsub.v1.AckReplyConsumerWithResponse;
+import com.google.cloud.pubsub.v1.AckResponse;
+import com.google.cloud.pubsub.v1.MessageReceiverWithAckResponse;
+import com.google.pubsub.v1.PubsubMessage;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.LinkedHashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.NavigableMap;
+import java.util.Set;
+import java.util.TreeMap;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Consumer;
+
+/** Reads Pub/Sub messages and acknowledges them only after their SeaTunnel 
checkpoint completes. */
+public class GooglePubSubSourceReader implements SourceReader<SeaTunnelRow, 
SingleSplit> {
+
+    private static final long POLL_TIMEOUT_MILLIS = 500;
+
+    private final DeserializationSchema<SeaTunnelRow> deserializationSchema;
+    private final SubscriberFactory subscriberFactory;
+    // Guards all message acknowledgement state shared with checkpoint 
callbacks.
+    private final Object acknowledgementLock = new Object();
+    private final BlockingQueue<ReceivedMessage> receivedMessages = new 
LinkedBlockingQueue<>();
+    // Messages emitted since the last completed checkpoint.
+    private final Set<AckReplyConsumerWithResponse> unacknowledgedMessages = 
new LinkedHashSet<>();
+    // Immutable acknowledgement snapshots keyed by SeaTunnel checkpoint ID.
+    private final NavigableMap<Long, List<AckReplyConsumerWithResponse>> 
pendingAcknowledgements =
+            new TreeMap<>();
+    // First asynchronous subscriber failure observed by the polling thread.
+    private final AtomicReference<Throwable> subscriberFailure = new 
AtomicReference<>();
+
+    private PubSubSubscriber subscriber;
+    private boolean splitAssigned;
+
+    public GooglePubSubSourceReader(
+            GooglePubSubSourceConfig config,
+            DeserializationSchema<SeaTunnelRow> deserializationSchema) {
+        this(
+                deserializationSchema,
+                (receiver, failureHandler) ->
+                        GooglePubSubSubscriber.create(config, receiver, 
failureHandler));
+    }
+
+    GooglePubSubSourceReader(
+            DeserializationSchema<SeaTunnelRow> deserializationSchema,
+            SubscriberFactory subscriberFactory) {
+        this.deserializationSchema = deserializationSchema;
+        this.subscriberFactory = subscriberFactory;
+    }
+
+    @Override
+    public void open() {
+        // The subscriber starts after the source split has been assigned.
+    }
+
+    @Override
+    public void pollNext(Collector<SeaTunnelRow> output) throws Exception {
+        checkSubscriberFailure();
+        ReceivedMessage receivedMessage =
+                receivedMessages.poll(POLL_TIMEOUT_MILLIS, 
TimeUnit.MILLISECONDS);
+        if (receivedMessage == null) {
+            checkSubscriberFailure();
+            return;
+        }
+
+        synchronized (output.getCheckpointLock()) {
+            try {
+                deserializationSchema.deserialize(
+                        receivedMessage.message.getData().toByteArray(), 
output);
+                synchronized (acknowledgementLock) {
+                    
unacknowledgedMessages.add(receivedMessage.acknowledgement);
+                }
+            } catch (Exception e) {
+                receivedMessage.acknowledgement.nack();
+                throw new GooglePubSubConnectorException(
+                        GooglePubSubConnectorErrorCode.READ_FAILED,
+                        "Failed to deserialize Google Pub/Sub message "
+                                + receivedMessage.message.getMessageId(),
+                        e);
+            }
+        }
+    }
+
+    @Override
+    public List<SingleSplit> snapshotState(long checkpointId) {
+        synchronized (acknowledgementLock) {
+            pendingAcknowledgements.put(checkpointId, new 
ArrayList<>(unacknowledgedMessages));
+        }
+        return Collections.singletonList(new SingleSplit(null));
+    }
+
+    @Override
+    public void addSplits(List<SingleSplit> splits) {
+        if (splits.size() != 1) {
+            throw new IllegalArgumentException(
+                    "Google Pub/Sub source expects exactly one source split");
+        }
+        if (splitAssigned) {
+            return;
+        }
+
+        subscriber =
+                subscriberFactory.create(
+                        (message, acknowledgement) ->
+                                receivedMessages.add(new 
ReceivedMessage(message, acknowledgement)),
+                        failure -> subscriberFailure.compareAndSet(null, 
failure));
+        subscriber.start();
+        splitAssigned = true;
+    }
+
+    @Override
+    public void handleNoMoreSplits() {
+        // The single subscription split remains active for the lifetime of 
the streaming job.
+    }
+
+    @Override
+    public void notifyCheckpointComplete(long checkpointId) {
+        List<AckReplyConsumerWithResponse> acknowledgements;
+        synchronized (acknowledgementLock) {
+            Map.Entry<Long, List<AckReplyConsumerWithResponse>> checkpoint =
+                    pendingAcknowledgements.floorEntry(checkpointId);
+            if (checkpoint == null) {
+                return;
+            }
+            acknowledgements = new ArrayList<>(checkpoint.getValue());
+        }
+
+        // Advance local state only after every acknowledgement in the 
selected checkpoint succeeds.
+        // If Pub/Sub accepts only part of the batch, failing the callback 
leaves the remaining
+        // messages eligible for redelivery instead of silently losing them 
from checkpoint state.
+        acknowledge(acknowledgements, checkpointId);
+        synchronized (acknowledgementLock) {
+            unacknowledgedMessages.removeAll(acknowledgements);
+            pendingAcknowledgements.headMap(checkpointId, true).clear();
+            for (List<AckReplyConsumerWithResponse> pending : 
pendingAcknowledgements.values()) {
+                pending.removeAll(acknowledgements);
+            }
+        }
+    }
+
+    @Override
+    public void notifyCheckpointAborted(long checkpointId) {
+        synchronized (acknowledgementLock) {
+            pendingAcknowledgements.remove(checkpointId);
+        }
+    }
+
+    @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) {
+            throw new GooglePubSubConnectorException(
+                    GooglePubSubConnectorErrorCode.READ_FAILED,
+                    "Google Pub/Sub subscriber stopped unexpectedly",
+                    failure);
+        }
+    }
+
+    private GooglePubSubConnectorException acknowledgeFailure(long 
checkpointId, Throwable cause) {
+        return new GooglePubSubConnectorException(
+                GooglePubSubConnectorErrorCode.ACKNOWLEDGE_FAILED,
+                "Failed to acknowledge Google Pub/Sub messages for checkpoint 
" + checkpointId,
+                cause);
+    }
+
+    @FunctionalInterface
+    interface SubscriberFactory {
+        PubSubSubscriber create(
+                MessageReceiverWithAckResponse receiver, Consumer<Throwable> 
failureHandler);
+    }
+
+    private static final class ReceivedMessage {
+        private final PubsubMessage message;
+        private final AckReplyConsumerWithResponse acknowledgement;
+
+        private ReceivedMessage(
+                PubsubMessage message, AckReplyConsumerWithResponse 
acknowledgement) {
+            this.message = message;
+            this.acknowledgement = acknowledgement;
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSubscriber.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSubscriber.java
new file mode 100644
index 0000000000..036f4d8097
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSubscriber.java
@@ -0,0 +1,184 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
+
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.GooglePubSubSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorErrorCode;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorException;
+
+import com.google.api.core.ApiService;
+import com.google.api.gax.batching.FlowControlSettings;
+import com.google.api.gax.core.FixedCredentialsProvider;
+import com.google.api.gax.core.NoCredentialsProvider;
+import com.google.api.gax.grpc.GrpcTransportChannel;
+import com.google.api.gax.rpc.FixedTransportChannelProvider;
+import com.google.auth.oauth2.GoogleCredentials;
+import com.google.cloud.pubsub.v1.MessageReceiverWithAckResponse;
+import com.google.cloud.pubsub.v1.Subscriber;
+import com.google.cloud.pubsub.v1.stub.SubscriberStubSettings;
+import com.google.pubsub.v1.ProjectSubscriptionName;
+import io.grpc.ManagedChannel;
+import io.grpc.ManagedChannelBuilder;
+
+import java.io.FileInputStream;
+import java.io.IOException;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.function.Consumer;
+
+/** Google Pub/Sub client lifecycle and transport configuration for a source 
reader. */
+class GooglePubSubSubscriber implements PubSubSubscriber {
+
+    private static final long CLOSE_TIMEOUT_SECONDS = 60;
+
+    private final Subscriber subscriber;
+    private final ManagedChannel emulatorChannel;
+
+    private GooglePubSubSubscriber(Subscriber subscriber, ManagedChannel 
emulatorChannel) {
+        this.subscriber = subscriber;
+        this.emulatorChannel = emulatorChannel;
+    }
+
+    static PubSubSubscriber create(
+            GooglePubSubSourceConfig config,
+            MessageReceiverWithAckResponse receiver,
+            Consumer<Throwable> failureHandler) {
+        ManagedChannel emulatorChannel = null;
+        try {
+            Subscriber.Builder subscriberBuilder =
+                    Subscriber.newBuilder(
+                            ProjectSubscriptionName.of(
+                                    config.getProjectId(), 
config.getSubscription()),
+                            receiver);
+
+            configureFlowControl(subscriberBuilder, config);
+
+            if (config.getEmulatorHost() != null) {
+                emulatorChannel =
+                        
ManagedChannelBuilder.forTarget(config.getEmulatorHost())
+                                .usePlaintext()
+                                .build();
+                subscriberBuilder
+                        .setChannelProvider(
+                                FixedTransportChannelProvider.create(
+                                        
GrpcTransportChannel.create(emulatorChannel)))
+                        
.setCredentialsProvider(NoCredentialsProvider.create());
+            } else if (config.getCredentialsPath() != null) {
+                try (FileInputStream credentialsStream =
+                        new FileInputStream(config.getCredentialsPath())) {
+                    subscriberBuilder.setCredentialsProvider(
+                            FixedCredentialsProvider.create(
+                                    
GoogleCredentials.fromStream(credentialsStream)
+                                            .createScoped(
+                                                    SubscriberStubSettings
+                                                            
.getDefaultServiceScopes())));
+                }
+            }
+
+            Subscriber subscriber = subscriberBuilder.build();
+            subscriber.addListener(
+                    new ApiService.Listener() {
+                        @Override
+                        public void failed(ApiService.State from, Throwable 
failure) {
+                            failureHandler.accept(failure);
+                        }
+                    },
+                    Runnable::run);
+            return new GooglePubSubSubscriber(subscriber, emulatorChannel);
+        } catch (Exception e) {
+            if (emulatorChannel != null) {
+                emulatorChannel.shutdownNow();
+            }
+            throw new GooglePubSubConnectorException(
+                    GooglePubSubConnectorErrorCode.CONNECTION_FAILED,
+                    "Failed to create Google Pub/Sub subscriber for 
subscription "
+                            + config.getSubscription(),
+                    e);
+        }
+    }
+
+    private static void configureFlowControl(
+            Subscriber.Builder subscriberBuilder, GooglePubSubSourceConfig 
config) {
+        if (config.getMaxOutstandingMessages() != null || 
config.getMaxOutstandingBytes() != null) {
+            FlowControlSettings.Builder flowControlSettings =
+                    
Subscriber.Builder.getDefaultFlowControlSettings().toBuilder();
+            if (config.getMaxOutstandingMessages() != null) {
+                flowControlSettings.setMaxOutstandingElementCount(
+                        config.getMaxOutstandingMessages());
+            }
+            if (config.getMaxOutstandingBytes() != null) {
+                
flowControlSettings.setMaxOutstandingRequestBytes(config.getMaxOutstandingBytes());
+            }
+            
subscriberBuilder.setFlowControlSettings(flowControlSettings.build());
+        }
+        if (config.getParallelPullCount() != null) {
+            
subscriberBuilder.setParallelPullCount(config.getParallelPullCount());
+        }
+    }
+
+    @Override
+    public void start() {
+        try {
+            subscriber.startAsync().awaitRunning();
+        } catch (RuntimeException e) {
+            throw new GooglePubSubConnectorException(
+                    GooglePubSubConnectorErrorCode.CONNECTION_FAILED,
+                    "Failed to start Google Pub/Sub subscriber",
+                    e);
+        }
+    }
+
+    @Override
+    public void close() throws IOException {
+        Throwable failure = null;
+        subscriber.stopAsync();
+        try {
+            subscriber.awaitTerminated(CLOSE_TIMEOUT_SECONDS, 
TimeUnit.SECONDS);
+        } catch (TimeoutException e) {
+            failure = e;
+        }
+
+        if (emulatorChannel != null) {
+            emulatorChannel.shutdown();
+            try {
+                if (!emulatorChannel.awaitTermination(CLOSE_TIMEOUT_SECONDS, 
TimeUnit.SECONDS)) {
+                    emulatorChannel.shutdownNow();
+                    IOException timeout =
+                            new IOException("Timed out while closing the 
Pub/Sub emulator channel");
+                    if (failure == null) {
+                        failure = timeout;
+                    } else {
+                        failure.addSuppressed(timeout);
+                    }
+                }
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                emulatorChannel.shutdownNow();
+                if (failure == null) {
+                    failure = e;
+                } else {
+                    failure.addSuppressed(e);
+                }
+            }
+        }
+
+        if (failure != null) {
+            throw new IOException("Failed to close Google Pub/Sub subscriber", 
failure);
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/PubSubSubscriber.java
similarity index 53%
copy from 
seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
copy to 
seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/PubSubSubscriber.java
index f4e8fcbc8a..b6720eb183 100644
--- 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/PubSubSubscriber.java
@@ -15,30 +15,12 @@
  * limitations under the License.
  */
 
-package org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception;
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
 
-import org.apache.seatunnel.common.exception.SeaTunnelErrorCode;
+import java.io.Closeable;
 
-public enum GooglePubSubConnectorErrorCode implements SeaTunnelErrorCode {
-    CONNECTION_FAILED("GooglePubSub-01", "Create Google Pub/Sub publisher 
failed"),
-    WRITE_FAILED("GooglePubSub-02", "Publish Google Pub/Sub message failed"),
-    CLOSE_FAILED("GooglePubSub-03", "Close Google Pub/Sub publisher failed");
+/** Lifecycle abstraction for the Pub/Sub client used by a source reader. */
+interface PubSubSubscriber extends Closeable {
 
-    private final String code;
-    private final String description;
-
-    GooglePubSubConnectorErrorCode(String code, String description) {
-        this.code = code;
-        this.description = description;
-    }
-
-    @Override
-    public String getCode() {
-        return code;
-    }
-
-    @Override
-    public String getDescription() {
-        return description;
-    }
+    void start();
 }
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceConfigTest.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceConfigTest.java
new file mode 100644
index 0000000000..5f6985b233
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/config/GooglePubSubSourceConfigTest.java
@@ -0,0 +1,106 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+class GooglePubSubSourceConfigTest {
+
+    @Test
+    void shouldRejectCredentialsForEmulator() {
+        Map<String, Object> options = requiredOptions();
+        options.put("credentials_path", "service-account.json");
+        options.put("emulator_host", "localhost:8085");
+
+        IllegalArgumentException exception =
+                Assertions.assertThrows(
+                        IllegalArgumentException.class,
+                        () -> 
GooglePubSubSourceConfig.from(ReadonlyConfig.fromMap(options)));
+        Assertions.assertTrue(exception.getMessage().contains("cannot be 
configured together"));
+    }
+
+    @Test
+    void shouldRejectBlankSubscription() {
+        Map<String, Object> options = requiredOptions();
+        options.put("subscription", " ");
+
+        IllegalArgumentException exception =
+                Assertions.assertThrows(
+                        IllegalArgumentException.class,
+                        () -> 
GooglePubSubSourceConfig.from(ReadonlyConfig.fromMap(options)));
+        Assertions.assertTrue(exception.getMessage().contains("subscription"));
+    }
+
+    @Test
+    void shouldRejectEmptyTextDelimiter() {
+        Map<String, Object> options = requiredOptions();
+        options.put("format", "text");
+        options.put("field_delimiter", "");
+
+        IllegalArgumentException exception =
+                Assertions.assertThrows(
+                        IllegalArgumentException.class,
+                        () -> 
GooglePubSubSourceConfig.from(ReadonlyConfig.fromMap(options)));
+        
Assertions.assertTrue(exception.getMessage().contains("field_delimiter"));
+    }
+
+    @Test
+    void shouldRejectNonPositiveFlowControlOptions() {
+        for (String option :
+                new String[] {
+                    "max_outstanding_messages", "max_outstanding_bytes", 
"parallel_pull_count"
+                }) {
+            Map<String, Object> options = requiredOptions();
+            options.put(option, 0);
+
+            IllegalArgumentException exception =
+                    Assertions.assertThrows(
+                            IllegalArgumentException.class,
+                            () -> 
GooglePubSubSourceConfig.from(ReadonlyConfig.fromMap(options)));
+            Assertions.assertTrue(exception.getMessage().contains(option));
+        }
+    }
+
+    @Test
+    void shouldReadFlowControlOptions() {
+        Map<String, Object> options = requiredOptions();
+        options.put("max_outstanding_messages", 100L);
+        options.put("max_outstanding_bytes", 1024L);
+        options.put("parallel_pull_count", 2);
+
+        GooglePubSubSourceConfig config =
+                GooglePubSubSourceConfig.from(ReadonlyConfig.fromMap(options));
+
+        Assertions.assertEquals(100L, config.getMaxOutstandingMessages());
+        Assertions.assertEquals(1024L, config.getMaxOutstandingBytes());
+        Assertions.assertEquals(2, config.getParallelPullCount());
+    }
+
+    private Map<String, Object> requiredOptions() {
+        Map<String, Object> options = new HashMap<>();
+        options.put("project_id", "test-project");
+        options.put("subscription", "test-subscription");
+        return options;
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceFactoryTest.java
similarity index 53%
copy from 
seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
copy to 
seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceFactoryTest.java
index f4e8fcbc8a..d9644f24ef 100644
--- 
a/seatunnel-connectors-v2/connector-google-pubsub/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/exception/GooglePubSubConnectorErrorCode.java
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceFactoryTest.java
@@ -15,30 +15,19 @@
  * limitations under the License.
  */
 
-package org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception;
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
 
-import org.apache.seatunnel.common.exception.SeaTunnelErrorCode;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
 
-public enum GooglePubSubConnectorErrorCode implements SeaTunnelErrorCode {
-    CONNECTION_FAILED("GooglePubSub-01", "Create Google Pub/Sub publisher 
failed"),
-    WRITE_FAILED("GooglePubSub-02", "Publish Google Pub/Sub message failed"),
-    CLOSE_FAILED("GooglePubSub-03", "Close Google Pub/Sub publisher failed");
+class GooglePubSubSourceFactoryTest {
 
-    private final String code;
-    private final String description;
+    @Test
+    void shouldExposeGooglePubSubIdentifierAndOptions() {
+        GooglePubSubSourceFactory factory = new GooglePubSubSourceFactory();
 
-    GooglePubSubConnectorErrorCode(String code, String description) {
-        this.code = code;
-        this.description = description;
-    }
-
-    @Override
-    public String getCode() {
-        return code;
-    }
-
-    @Override
-    public String getDescription() {
-        return description;
+        Assertions.assertEquals("GooglePubSub", factory.factoryIdentifier());
+        Assertions.assertNotNull(factory.optionRule());
+        Assertions.assertEquals(GooglePubSubSource.class, 
factory.getSourceClass());
     }
 }
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReaderTest.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReaderTest.java
new file mode 100644
index 0000000000..109e7bd143
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceReaderTest.java
@@ -0,0 +1,273 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
+
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplit;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorException;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import com.google.api.core.ApiFuture;
+import com.google.api.core.ApiFutures;
+import com.google.cloud.pubsub.v1.AckReplyConsumerWithResponse;
+import com.google.cloud.pubsub.v1.AckResponse;
+import com.google.cloud.pubsub.v1.MessageReceiverWithAckResponse;
+import com.google.protobuf.ByteString;
+import com.google.pubsub.v1.PubsubMessage;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.Collections;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.function.Consumer;
+
+class GooglePubSubSourceReaderTest {
+
+    @Test
+    void shouldAcknowledgeOnlyAfterCheckpointCompletes() throws Exception {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader = createReader(subscriberFactory);
+        TestAcknowledgement acknowledgement = new TestAcknowledgement();
+
+        assignSplit(reader);
+        subscriberFactory.emit("first", acknowledgement);
+        reader.pollNext(new TestCollector());
+
+        Assertions.assertEquals(0, acknowledgement.ackCount.get());
+        reader.snapshotState(1L);
+        Assertions.assertEquals(0, acknowledgement.ackCount.get());
+
+        reader.notifyCheckpointComplete(1L);
+        Assertions.assertEquals(1, acknowledgement.ackCount.get());
+    }
+
+    // An aborted checkpoint must not release messages. A later successful 
checkpoint owns them.
+    @Test
+    void shouldKeepAcknowledgementAfterCheckpointIsAborted() throws Exception {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader = createReader(subscriberFactory);
+        TestAcknowledgement acknowledgement = new TestAcknowledgement();
+
+        assignSplit(reader);
+        subscriberFactory.emit("first", acknowledgement);
+        reader.pollNext(new TestCollector());
+        reader.snapshotState(1L);
+        reader.notifyCheckpointAborted(1L);
+
+        reader.snapshotState(2L);
+        reader.notifyCheckpointComplete(2L);
+
+        Assertions.assertEquals(1, acknowledgement.ackCount.get());
+    }
+
+    // Completing a newer overlapping checkpoint acknowledges every message 
contained in it once.
+    @Test
+    void shouldAcknowledgeMessagesFromOverlappingCheckpointsOnce() throws 
Exception {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader = createReader(subscriberFactory);
+        TestAcknowledgement first = new TestAcknowledgement();
+        TestAcknowledgement second = new TestAcknowledgement();
+
+        assignSplit(reader);
+        subscriberFactory.emit("first", first);
+        reader.pollNext(new TestCollector());
+        reader.snapshotState(1L);
+
+        subscriberFactory.emit("second", second);
+        reader.pollNext(new TestCollector());
+        reader.snapshotState(2L);
+        reader.notifyCheckpointComplete(2L);
+        reader.notifyCheckpointComplete(1L);
+
+        Assertions.assertEquals(1, first.ackCount.get());
+        Assertions.assertEquals(1, second.ackCount.get());
+    }
+
+    @Test
+    void shouldFailCheckpointWhenPubSubRejectsAcknowledgement() throws 
Exception {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader = createReader(subscriberFactory);
+        TestAcknowledgement acknowledgement = new 
TestAcknowledgement(AckResponse.INVALID);
+
+        assignSplit(reader);
+        subscriberFactory.emit("first", acknowledgement);
+        reader.pollNext(new TestCollector());
+        reader.snapshotState(1L);
+
+        GooglePubSubConnectorException exception =
+                Assertions.assertThrows(
+                        GooglePubSubConnectorException.class,
+                        () -> reader.notifyCheckpointComplete(1L));
+        Assertions.assertTrue(exception.getMessage().contains("INVALID"));
+    }
+
+    @Test
+    void shouldNackMessageWhenDeserializationFails() throws Exception {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader =
+                new GooglePubSubSourceReader(new 
FailingDeserializationSchema(), subscriberFactory);
+        TestAcknowledgement acknowledgement = new TestAcknowledgement();
+
+        assignSplit(reader);
+        subscriberFactory.emit("invalid", acknowledgement);
+
+        Assertions.assertThrows(
+                GooglePubSubConnectorException.class, () -> 
reader.pollNext(new TestCollector()));
+        Assertions.assertEquals(1, acknowledgement.nackCount.get());
+    }
+
+    @Test
+    void shouldPropagateSubscriberFailure() {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader = createReader(subscriberFactory);
+
+        assignSplit(reader);
+        subscriberFactory.fail(new IOException("stream stopped"));
+
+        GooglePubSubConnectorException exception =
+                Assertions.assertThrows(
+                        GooglePubSubConnectorException.class,
+                        () -> reader.pollNext(new TestCollector()));
+        Assertions.assertTrue(exception.getMessage().contains("stopped 
unexpectedly"));
+    }
+
+    @Test
+    void shouldDeserializeMessagePayload() throws Exception {
+        TestSubscriberFactory subscriberFactory = new TestSubscriberFactory();
+        GooglePubSubSourceReader reader = createReader(subscriberFactory);
+        TestCollector collector = new TestCollector();
+
+        assignSplit(reader);
+        subscriberFactory.emit("hello", new TestAcknowledgement());
+        reader.pollNext(collector);
+
+        Assertions.assertEquals("hello", collector.value);
+    }
+
+    private GooglePubSubSourceReader createReader(TestSubscriberFactory 
subscriberFactory) {
+        return new GooglePubSubSourceReader(new TestDeserializationSchema(), 
subscriberFactory);
+    }
+
+    private void assignSplit(GooglePubSubSourceReader reader) {
+        reader.addSplits(Collections.singletonList(new SingleSplit(null)));
+    }
+
+    private static final class TestSubscriberFactory
+            implements GooglePubSubSourceReader.SubscriberFactory {
+        private MessageReceiverWithAckResponse receiver;
+        private Consumer<Throwable> failureHandler;
+
+        @Override
+        public PubSubSubscriber create(
+                MessageReceiverWithAckResponse receiver, Consumer<Throwable> 
failureHandler) {
+            this.receiver = receiver;
+            this.failureHandler = failureHandler;
+            return new PubSubSubscriber() {
+                @Override
+                public void start() {}
+
+                @Override
+                public void close() {}
+            };
+        }
+
+        private void emit(String value, AckReplyConsumerWithResponse 
acknowledgement) {
+            receiver.receiveMessage(
+                    PubsubMessage.newBuilder()
+                            .setMessageId(value)
+                            .setData(ByteString.copyFromUtf8(value))
+                            .build(),
+                    acknowledgement);
+        }
+
+        private void fail(Throwable failure) {
+            failureHandler.accept(failure);
+        }
+    }
+
+    private static final class TestAcknowledgement implements 
AckReplyConsumerWithResponse {
+        private final AtomicInteger ackCount = new AtomicInteger();
+        private final AtomicInteger nackCount = new AtomicInteger();
+        private final AckResponse response;
+
+        private TestAcknowledgement() {
+            this(AckResponse.SUCCESSFUL);
+        }
+
+        private TestAcknowledgement(AckResponse response) {
+            this.response = response;
+        }
+
+        @Override
+        public ApiFuture<AckResponse> ack() {
+            ackCount.incrementAndGet();
+            return ApiFutures.immediateFuture(response);
+        }
+
+        @Override
+        public ApiFuture<AckResponse> nack() {
+            nackCount.incrementAndGet();
+            return ApiFutures.immediateFuture(AckResponse.SUCCESSFUL);
+        }
+    }
+
+    private static class TestDeserializationSchema implements 
DeserializationSchema<SeaTunnelRow> {
+        private static final SeaTunnelRowType ROW_TYPE =
+                new SeaTunnelRowType(
+                        new String[] {"value"}, new SeaTunnelDataType[] 
{BasicType.STRING_TYPE});
+
+        @Override
+        public SeaTunnelRow deserialize(byte[] message) throws IOException {
+            return new SeaTunnelRow(new Object[] {new String(message, 
StandardCharsets.UTF_8)});
+        }
+
+        @Override
+        public SeaTunnelDataType<SeaTunnelRow> getProducedType() {
+            return ROW_TYPE;
+        }
+    }
+
+    private static final class FailingDeserializationSchema extends 
TestDeserializationSchema {
+        @Override
+        public SeaTunnelRow deserialize(byte[] message) throws IOException {
+            throw new IOException("invalid payload");
+        }
+    }
+
+    private static final class TestCollector implements 
Collector<SeaTunnelRow> {
+        private final Object checkpointLock = new Object();
+        private String value;
+
+        @Override
+        public void collect(SeaTunnelRow record) {
+            value = record.getField(0).toString();
+        }
+
+        @Override
+        public Object getCheckpointLock() {
+            return checkpointLock;
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceTest.java
 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceTest.java
new file mode 100644
index 0000000000..80ec871ef4
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-google-pubsub/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/pubsub/source/GooglePubSubSourceTest.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.google.pubsub.source;
+
+import org.apache.seatunnel.api.common.JobContext;
+import org.apache.seatunnel.api.serialization.DeserializationSchema;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.common.constants.JobMode;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.GooglePubSubSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.config.MessageFormat;
+import 
org.apache.seatunnel.connectors.seatunnel.google.pubsub.exception.GooglePubSubConnectorException;
+import org.apache.seatunnel.format.json.JsonDeserializationSchema;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+class GooglePubSubSourceTest {
+
+    @Test
+    void shouldRequireStreamingModeWithCheckpointing() {
+        GooglePubSubSource source = createSource();
+
+        source.setJobContext(
+                new 
JobContext().setJobMode(JobMode.STREAMING).setEnableCheckpoint(true));
+        Assertions.assertEquals(Boundedness.UNBOUNDED, 
source.getBoundedness());
+
+        source.setJobContext(
+                new 
JobContext().setJobMode(JobMode.STREAMING).setEnableCheckpoint(false));
+        GooglePubSubConnectorException checkpointException =
+                Assertions.assertThrows(
+                        GooglePubSubConnectorException.class, 
source::getBoundedness);
+        
Assertions.assertTrue(checkpointException.getMessage().contains("requires 
checkpointing"));
+
+        source.setJobContext(new 
JobContext().setJobMode(JobMode.BATCH).setEnableCheckpoint(true));
+        GooglePubSubConnectorException batchException =
+                Assertions.assertThrows(
+                        GooglePubSubConnectorException.class, 
source::getBoundedness);
+        Assertions.assertTrue(batchException.getMessage().contains("streaming 
jobs only"));
+    }
+
+    private GooglePubSubSource createSource() {
+        CatalogTable catalogTable = CatalogTableUtil.buildSimpleTextTable();
+        DeserializationSchema<SeaTunnelRow> deserializationSchema =
+                new JsonDeserializationSchema(catalogTable, false, false);
+        GooglePubSubSourceConfig config =
+                GooglePubSubSourceConfig.builder()
+                        .projectId("project")
+                        .subscription("subscription")
+                        .format(MessageFormat.JSON)
+                        .fieldDelimiter(",")
+                        .build();
+        return new GooglePubSubSource(config, catalogTable, 
deserializationSchema);
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/java/org/apache/seatunnel/e2e/connector/google/pubsub/GooglePubSubIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/java/org/apache/seatunnel/e2e/connector/google/pubsub/GooglePubSubIT.java
index 2da4794714..544c5fbf08 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/java/org/apache/seatunnel/e2e/connector/google/pubsub/GooglePubSubIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/java/org/apache/seatunnel/e2e/connector/google/pubsub/GooglePubSubIT.java
@@ -23,7 +23,10 @@ import 
org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ObjectNode
 import org.apache.seatunnel.common.utils.JsonUtils;
 import org.apache.seatunnel.e2e.common.TestResource;
 import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.EngineType;
 import org.apache.seatunnel.e2e.common.container.TestContainer;
+import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer;
+import org.apache.seatunnel.e2e.common.util.JobIdGenerator;
 
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.Assertions;
@@ -46,6 +49,7 @@ import java.net.HttpURLConnection;
 import java.net.URL;
 import java.nio.charset.StandardCharsets;
 import java.util.Base64;
+import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.stream.Stream;
@@ -60,10 +64,16 @@ public class GooglePubSubIT extends TestSuiteBase 
implements TestResource {
     private static final String PROJECT_ID = "seatunnel-test";
     private static final String TOPIC_ID = "events";
     private static final String SUBSCRIPTION_ID = "events-test";
+    private static final String SOURCE_TOPIC_ID = "source-events";
+    private static final String SOURCE_SUBSCRIPTION_ID = "source-events-test";
     private static final int EMULATOR_PORT = 8085;
     private static final String EMULATOR_HOST = "pubsub-emulator";
     private static final String JOB_CONFIG = 
"/pubsub/fake_to_google_pubsub.conf";
+    private static final String SOURCE_JOB_CONFIG = 
"/pubsub/google_pubsub_to_console.conf";
     private static final String EXPECTED_MESSAGE = 
"{\"name\":\"alice\",\"age\":30}";
+    private static final String EXPECTED_SOURCE_EVENT_ID = 
"pubsub-source-checkpoint-event";
+    private static final String SOURCE_MESSAGE =
+            "{\"event_id\":\"" + EXPECTED_SOURCE_EVENT_ID + 
"\",\"event_type\":\"created\"}";
 
     private GenericContainer<?> emulator;
     private String emulatorEndpoint;
@@ -94,6 +104,7 @@ public class GooglePubSubIT extends TestSuiteBase implements 
TestResource {
         emulatorEndpoint =
                 "http://"; + emulator.getHost() + ":" + 
emulator.getMappedPort(EMULATOR_PORT);
         createResource("/v1/projects/" + PROJECT_ID + "/topics/" + TOPIC_ID, 
"{}");
+        createResource("/v1/projects/" + PROJECT_ID + "/topics/" + 
SOURCE_TOPIC_ID, "{}");
         createResource(
                 "/v1/projects/" + PROJECT_ID + "/subscriptions/" + 
SUBSCRIPTION_ID,
                 "{\"topic\":\"projects/"
@@ -101,6 +112,13 @@ public class GooglePubSubIT extends TestSuiteBase 
implements TestResource {
                         + "/topics/"
                         + TOPIC_ID
                         + "\",\"ackDeadlineSeconds\":10}");
+        createResource(
+                "/v1/projects/" + PROJECT_ID + "/subscriptions/" + 
SOURCE_SUBSCRIPTION_ID,
+                "{\"topic\":\"projects/"
+                        + PROJECT_ID
+                        + "/topics/"
+                        + SOURCE_TOPIC_ID
+                        + "\",\"ackDeadlineSeconds\":10}");
     }
 
     @AfterAll
@@ -148,6 +166,85 @@ public class GooglePubSubIT extends TestSuiteBase 
implements TestResource {
         acknowledge(message.path("ackId").asText());
     }
 
+    @TestTemplate
+    @DisabledOnContainer(
+            value = {},
+            type = {EngineType.FLINK, EngineType.SPARK},
+            disabledReason =
+                    "The source checkpoint assertion uses the Zeta REST job 
status and server logs")
+    public void 
testGooglePubSubSourceAcknowledgesAfterCheckpoint(TestContainer container)
+            throws Exception {
+        String jobId = String.valueOf(JobIdGenerator.newJobId());
+        CompletableFuture<Container.ExecResult> jobFuture =
+                CompletableFuture.supplyAsync(
+                        () -> {
+                            try {
+                                return container.executeJob(SOURCE_JOB_CONFIG, 
jobId);
+                            } catch (Exception e) {
+                                throw new RuntimeException(e);
+                            }
+                        });
+
+        try {
+            await().atMost(60, TimeUnit.SECONDS)
+                    .pollInterval(1, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                assertJobStillRunning(jobFuture);
+                                Assertions.assertEquals("RUNNING", 
container.getJobStatus(jobId));
+                            });
+
+            long checkpointCount = 
container.getCompletedCheckpointCount(jobId);
+            publish(SOURCE_TOPIC_ID, SOURCE_MESSAGE);
+
+            await().atMost(60, TimeUnit.SECONDS)
+                    .pollInterval(1, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                assertJobStillRunning(jobFuture);
+                                Assertions.assertTrue(
+                                        container
+                                                .getServerLogs()
+                                                
.contains(EXPECTED_SOURCE_EVENT_ID),
+                                        "Published Pub/Sub message was not 
emitted by the source");
+                            });
+            await().atMost(60, TimeUnit.SECONDS)
+                    .pollInterval(1, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                assertJobStillRunning(jobFuture);
+                                Assertions.assertTrue(
+                                        
container.getCompletedCheckpointCount(jobId)
+                                                > checkpointCount,
+                                        "No checkpoint completed after the 
Pub/Sub message was emitted");
+                            });
+        } finally {
+            if (!jobFuture.isDone()) {
+                Container.ExecResult cancelResult = container.cancelJob(jobId);
+                Assertions.assertEquals(0, cancelResult.getExitCode(), 
cancelResult.getStderr());
+            }
+        }
+
+        Container.ExecResult jobResult = jobFuture.get(120, TimeUnit.SECONDS);
+        Assertions.assertEquals(0, jobResult.getExitCode(), 
jobResult.getStderr());
+    }
+
+    private void publish(String topicId, String payload) throws IOException {
+        String data = 
Base64.getEncoder().encodeToString(payload.getBytes(StandardCharsets.UTF_8));
+        request(
+                "POST",
+                "/v1/projects/" + PROJECT_ID + "/topics/" + topicId + 
":publish",
+                "{\"messages\":[{\"data\":\"" + data + "\"}]}");
+    }
+
+    private void assertJobStillRunning(CompletableFuture<Container.ExecResult> 
jobFuture)
+            throws Exception {
+        if (jobFuture.isDone()) {
+            Container.ExecResult result = jobFuture.get();
+            Assertions.fail("Streaming source job terminated early:\n" + 
result.getStderr());
+        }
+    }
+
     private void acknowledge(String ackId) throws IOException {
         request(
                 "POST",
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/resources/pubsub/google_pubsub_to_console.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/resources/pubsub/google_pubsub_to_console.conf
new file mode 100644
index 0000000000..e5d65feecb
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-google-pubsub-e2e/src/test/resources/pubsub/google_pubsub_to_console.conf
@@ -0,0 +1,41 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+env {
+  parallelism = 1
+  job.mode = "STREAMING"
+  checkpoint.interval = 1000
+}
+
+source {
+  GooglePubSub {
+    project_id = "seatunnel-test"
+    subscription = "source-events-test"
+    emulator_host = "pubsub-emulator:8085"
+    format = json
+    schema = {
+      fields {
+        event_id = string
+        event_type = string
+      }
+    }
+  }
+}
+
+sink {
+  Console {}
+}

Reply via email to