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 7e468d60f6 [Feature][Connector-V2] Add Google Ads Source Connector
(#12230)
7e468d60f6 is described below
commit 7e468d60f67832e3e5fa00bf3035fbf193e74238
Author: Claire <[email protected]>
AuthorDate: Thu Sep 10 14:40:51 2026 +0000
[Feature][Connector-V2] Add Google Ads Source Connector (#12230)
---
.github/workflows/labeler/label-scope-conf.yml | 5 +
config/plugin_config | 1 +
docs/en/connector-v2/source/GoogleAds.md | 185 +++++++++
docs/zh/connector-v2/source/GoogleAds.md | 179 +++++++++
plugin-mapping.properties | 1 +
.../connector-google-ads/pom.xml | 55 +++
.../google/ads/client/GoogleAdsClient.java | 402 ++++++++++++++++++++
.../google/ads/config/GoogleAdsParameters.java | 59 +++
.../google/ads/config/GoogleAdsSourceOptions.java | 160 ++++++++
.../google/ads/config/GoogleAdsTableConfig.java | 53 +++
.../ads/exception/GoogleAdsConnectorErrorCode.java | 47 +++
.../ads/exception/GoogleAdsConnectorException.java | 36 ++
.../google/ads/source/GoogleAdsSource.java | 251 +++++++++++++
.../google/ads/source/GoogleAdsSourceFactory.java | 82 ++++
.../google/ads/source/GoogleAdsSourceReader.java | 98 +++++
.../google/ads/GoogleAdsSourceFactoryTest.java | 209 +++++++++++
.../google/ads/client/GoogleAdsClientTest.java | 418 +++++++++++++++++++++
.../google/ads/source/GoogleAdsSourceTest.java | 296 +++++++++++++++
seatunnel-connectors-v2/pom.xml | 1 +
seatunnel-dist/pom.xml | 6 +
20 files changed, 2544 insertions(+)
diff --git a/.github/workflows/labeler/label-scope-conf.yml
b/.github/workflows/labeler/label-scope-conf.yml
index 141c3f8469..be8f0f0261 100644
--- a/.github/workflows/labeler/label-scope-conf.yml
+++ b/.github/workflows/labeler/label-scope-conf.yml
@@ -169,6 +169,11 @@ fluss:
- changed-files:
- any-glob-to-any-file: seatunnel-connectors-v2/connector-fluss/**
- all-globs-to-all-files:
'!seatunnel-connectors-v2/connector-!(fluss)/**'
+google-ads:
+ - all:
+ - changed-files:
+ - any-glob-to-any-file:
seatunnel-connectors-v2/connector-google-ads/**
+ - all-globs-to-all-files:
'!seatunnel-connectors-v2/connector-!(google-ads)/**'
google-firestore:
- all:
- changed-files:
diff --git a/config/plugin_config b/config/plugin_config
index af6ce45b7e..50decdbe6f 100644
--- a/config/plugin_config
+++ b/config/plugin_config
@@ -58,6 +58,7 @@ connector-google-firestore
connector-google-pubsub
connector-azure-queue-storage
connector-google-bigtable
+connector-google-ads
connector-graphql
connector-hive
connector-http-base
diff --git a/docs/en/connector-v2/source/GoogleAds.md
b/docs/en/connector-v2/source/GoogleAds.md
new file mode 100644
index 0000000000..ed524f928f
--- /dev/null
+++ b/docs/en/connector-v2/source/GoogleAds.md
@@ -0,0 +1,185 @@
+# GoogleAds
+
+> Google Ads source connector
+
+## Support Those Engines
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## Key Features
+
+- [x] [batch](../../introduction/concepts/connector-v2-features.md)
+- [ ] [stream](../../introduction/concepts/connector-v2-features.md)
+- [ ] [exactly-once](../../introduction/concepts/connector-v2-features.md)
+- [x] [column projection](../../introduction/concepts/connector-v2-features.md)
+- [ ] [parallelism](../../introduction/concepts/connector-v2-features.md)
+- [x] [support multiple table
read](../../introduction/concepts/connector-v2-features.md)
+
+## Description
+
+Reads data from Google Ads resources (campaign, ad_group, keyword_view, ...)
+using the Google Ads REST API (`googleAds:search`) with GAQL queries. Supports
+single-resource, full-GAQL-query, and multi-table (`tables_configs`) batch
+ingestion. Results are streamed page by page via `nextPageToken`, so only one
+page is held in memory at a time.
+
+Authentication uses the OAuth 2.0 refresh-token grant plus a Google Ads
+developer token. The access token is refreshed automatically when it expires.
+
+The output schema is derived automatically: the connector queries the
+`googleAdsFields:search` metadata service for the data type of every selected
+field and builds the schema **in the SELECT field order**, so column `i`
+always corresponds to selected field `i`.
+
+## Supported DataSource Info
+
+| Datasource | Supported Versions |
+|------------|----------------------|
+| Google Ads | REST API v21 (default, configurable via `api_version`) |
+
+## Prerequisites
+
+1. A Google Ads **developer token** (from the API Center of a manager/MCC
account).
+2. An OAuth2 **client ID / client secret** (Google Cloud Console, OAuth consent
+ configured for the `https://www.googleapis.com/auth/adwords` scope).
+3. A **refresh token** authorized for that scope (e.g. via the OAuth 2.0
+ Playground or the Google Ads API oauth helper scripts).
+4. The **customer ID** of the account to query (digits only, no dashes). If it
+ is a client account under an MCC, also set `login_customer_id` to the MCC
+ customer ID.
+
+## Source Options
+
+| Name | Type | Required | Default | Description
|
+|--------------------|---------|----------|---------|------------------------------------------------------------------------------------------------------|
+| developer_token | String | Yes | - | Google Ads API developer
token. |
+| client_id | String | Yes | - | OAuth2 client ID.
|
+| client_secret | String | Yes | - | OAuth2 client secret.
|
+| refresh_token | String | Yes | - | OAuth2 refresh token for
the `adwords` scope. |
+| customer_id | String | Yes | - | Customer ID to query,
digits only, e.g. `1234567890`.
|
+| login_customer_id | String | No | - | Manager (MCC) customer
ID, digits only. Required when `customer_id` is managed by an MCC. |
+| api_version | String | No | v21 | Google Ads REST API
version.
|
+| resource | String | No* | - | Resource for
single-table mode, e.g. `campaign`. Requires `fields`. Exclusive with `query`
and `tables_configs`. |
+| fields | List | No | - | Ordered list of GAQL
field paths, e.g. `[campaign.id, metrics.clicks]`. The output schema follows
this order. |
+| filter | String | No | - | GAQL WHERE clause
appended to the auto-built `SELECT <fields> FROM <resource>` query.
|
+| query | String | No* | - | Full GAQL query.
Exclusive with `resource`/`fields`/`filter` and `tables_configs`.
|
+| tables_configs | List | No* | - | Multi-table
configuration list. Each entry requires `table_path`. Exclusive with `resource`
and `query`. |
+| request_timeout_ms | Integer | No | 60000 | HTTP request timeout in
milliseconds for a single call. |
+| max_retries | Integer | No | 3 | Maximum retries for
transient HTTP failures (429/5xx/network errors) of a single request.
|
+| retry_backoff_ms | Long | No | 1000 | Base backoff in
milliseconds between retries; doubled per attempt.
|
+| page_size | Integer | No | - | Page size for search
requests. When unset the server default is used. Note: recent API versions
ignore or reject an explicit page size. |
+
+\* Exactly one of `resource`, `query` or `tables_configs` must be provided.
+
+GAQL has no `SELECT *`, so the field list is always explicit — either via
+`fields` or inside `query`. Invalid field names are rejected before any data is
+read, with the offending name in the error message. Invalid GAQL (HTTP 400) is
+never retried and the API error message is surfaced as-is.
+
+### tables_configs entry options
+
+| Name | Type | Required | Description
|
+|-------------|--------|----------|--------------------------------------------------------------------------------|
+| table_path | String | Yes | Format: `database.resource`, e.g.
`google_ads.campaign`. |
+| fields | List | No* | Ordered GAQL field paths for this table.
|
+| query | String | No* | Full GAQL query for this table. Its `FROM`
resource must match `table_path`. |
+| filter | String | No | GAQL WHERE clause (only with `fields`).
|
+| customer_id | String | No | Per-table customer ID override; falls back
to the global `customer_id`. |
+
+\* Exactly one of `fields` or `query` per entry.
+
+## Data Type Mapping
+
+| Google Ads Data Type | SeaTunnel Data Type | Notes
|
+|----------------------------|---------------------|--------------------------------------------------------------------------|
+| INT64, UINT64 | BIGINT | The REST API serializes
int64 as a JSON string; parsed to long. |
+| INT32 | INT |
|
+| DOUBLE, FLOAT | DOUBLE |
|
+| BOOLEAN | BOOLEAN |
|
+| DATE | STRING | Deliberate: date-like
fields have non-uniform formats (`2026-09-01`, `2026-09`, `2026-36`). |
+| STRING, ENUM, RESOURCE_NAME | STRING | Enums are symbolic
strings, e.g. `ENABLED`. |
+| MESSAGE | STRING | Nested messages are
emitted as JSON text. |
+| Other / unknown | STRING | Forward-compatible
fallback. |
+
+Fields absent from a result row (the API omits empty fields entirely) are
+emitted as `null`.
+
+## Example
+
+### Single resource
+
+```hocon
+source {
+ GoogleAds {
+ developer_token = "your_developer_token"
+ client_id = "your_client_id"
+ client_secret = "your_client_secret"
+ refresh_token = "your_refresh_token"
+ customer_id = "1234567890"
+
+ resource = "campaign"
+ fields = ["campaign.id", "campaign.name", "campaign.status",
"metrics.clicks", "metrics.impressions"]
+ filter = "segments.date DURING LAST_30_DAYS"
+ }
+}
+```
+
+### Full GAQL query
+
+```hocon
+source {
+ GoogleAds {
+ developer_token = "your_developer_token"
+ client_id = "your_client_id"
+ client_secret = "your_client_secret"
+ refresh_token = "your_refresh_token"
+ customer_id = "1234567890"
+
+ query = "SELECT ad_group.id, ad_group.name, metrics.clicks FROM ad_group
WHERE metrics.clicks > 0"
+ }
+}
+```
+
+### Multiple tables (with per-table customer ID)
+
+```hocon
+source {
+ GoogleAds {
+ developer_token = "your_developer_token"
+ client_id = "your_client_id"
+ client_secret = "your_client_secret"
+ refresh_token = "your_refresh_token"
+ customer_id = "1234567890"
+ login_customer_id = "9876543210"
+
+ tables_configs = [
+ {
+ table_path = "google_ads.campaign"
+ fields = ["campaign.id", "campaign.name", "metrics.cost_micros"]
+ filter = "segments.date DURING LAST_7_DAYS"
+ },
+ {
+ table_path = "google_ads.ad_group"
+ query = "SELECT ad_group.id, ad_group.name FROM ad_group"
+ customer_id = "2345678901"
+ }
+ ]
+ }
+}
+```
+
+## Limitations
+
+- Batch only; no incremental/CDC reads (use `filter` with date segments for
+ windowed extraction).
+- No parallel reading; each job reads with a single split.
+- No exactly-once semantics; re-running a job re-reads the data.
+- Nested MESSAGE fields are emitted as JSON strings, not nested rows.
+
+## Changelog
+
+### next version
+
+- Add Google Ads source connector with GAQL, automatic schema derivation and
multi-table support
diff --git a/docs/zh/connector-v2/source/GoogleAds.md
b/docs/zh/connector-v2/source/GoogleAds.md
new file mode 100644
index 0000000000..c7916d408c
--- /dev/null
+++ b/docs/zh/connector-v2/source/GoogleAds.md
@@ -0,0 +1,179 @@
+# GoogleAds
+
+> Google Ads 源连接器
+
+## 支持的引擎
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## 主要特性
+
+- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
+- [ ] [流处理](../../introduction/concepts/connector-v2-features.md)
+- [ ] [精确一次](../../introduction/concepts/connector-v2-features.md)
+- [x] [列投影](../../introduction/concepts/connector-v2-features.md)
+- [ ] [并行读取](../../introduction/concepts/connector-v2-features.md)
+- [x] [支持多表读取](../../introduction/concepts/connector-v2-features.md)
+
+## 描述
+
+通过 Google Ads REST API(`googleAds:search`)和 GAQL 查询读取 Google Ads
+资源数据(campaign、ad_group、keyword_view 等)。支持单资源模式、完整 GAQL
+查询模式和多表(`tables_configs`)批量读取。结果通过 `nextPageToken` 逐页
+流式读取,内存中同一时间只保留一页数据。
+
+认证使用 OAuth 2.0 refresh-token 授权方式,并携带 Google Ads developer
+token。access token 过期时自动刷新。
+
+输出 schema 自动推导:连接器通过 `googleAdsFields:search` 元数据服务获取
+每个所选字段的数据类型,并**按 SELECT 字段顺序**构建 schema,第 i 列始终
+对应第 i 个所选字段。
+
+## 支持的数据源信息
+
+| 数据源 | 支持版本 |
+|------------|---------------------------------------------|
+| Google Ads | REST API v21(默认,可通过 `api_version` 配置) |
+
+## 前置条件
+
+1. Google Ads **developer token**(在管理者/MCC 账号的 API Center 获取)。
+2. OAuth2 **client ID / client secret**(Google Cloud Console,OAuth 同意屏幕
+ 需包含 `https://www.googleapis.com/auth/adwords` scope)。
+3. 已授权该 scope 的 **refresh token**(可通过 OAuth 2.0 Playground 或
+ Google Ads API 官方 oauth 脚本获取)。
+4. 要查询账号的 **customer ID**(纯数字,不含短横线)。如果该账号由 MCC
+ 管理,还需要将 `login_customer_id` 设置为 MCC 的 customer ID。
+
+## 源选项
+
+| 名称 | 类型 | 必填 | 默认值 | 描述
|
+|--------------------|---------|------|---------|--------------------------------------------------------------------------------------------|
+| developer_token | String | 是 | - | Google Ads API developer
token。 |
+| client_id | String | 是 | - | OAuth2 client ID。
|
+| client_secret | String | 是 | - | OAuth2 client secret。
|
+| refresh_token | String | 是 | - | 已授权 `adwords` scope 的 OAuth2
refresh token。 |
+| customer_id | String | 是 | - | 要查询的 customer ID,纯数字,例如
`1234567890`。 |
+| login_customer_id | String | 否 | - | 管理者(MCC)customer ID,纯数字。当
`customer_id` 由 MCC 管理时必须设置。 |
+| api_version | String | 否 | v21 | Google Ads REST API 版本。
|
+| resource | String | 否* | - | 单表模式的资源名,例如 `campaign`,需配合
`fields`。与 `query`、`tables_configs` 互斥。 |
+| fields | List | 否 | - | 有序的 GAQL 字段路径列表,例如
`[campaign.id, metrics.clicks]`。输出 schema 按此顺序排列。 |
+| filter | String | 否 | - | 追加到自动构建的 `SELECT <fields>
FROM <resource>` 查询后的 GAQL WHERE 子句。 |
+| query | String | 否* | - | 完整 GAQL 查询。与
`resource`/`fields`/`filter` 及 `tables_configs` 互斥。 |
+| tables_configs | List | 否* | - | 多表配置列表,每项必须包含 `table_path`。与
`resource`、`query` 互斥。 |
+| request_timeout_ms | Integer | 否 | 60000 | 单次 HTTP 请求超时时间(毫秒)。
|
+| max_retries | Integer | 否 | 3 |
单个请求瞬时失败(429/5xx/网络错误)的最大重试次数。 |
+| retry_backoff_ms | Long | 否 | 1000 | 重试基础退避时间(毫秒),每次翻倍。
|
+| page_size | Integer | 否 | - | 搜索请求分页大小,不设置时使用服务端默认值。注意:较新的
API 版本会忽略或拒绝显式分页大小。 |
+
+\* `resource`、`query`、`tables_configs` 三者必须且只能提供一个。
+
+GAQL 不支持 `SELECT *`,字段列表始终是显式的——通过 `fields` 或 `query`
+指定。非法字段名会在读取数据前被拦截,错误信息中包含出错的字段名。非法
+GAQL(HTTP 400)不会重试,API 错误信息原样透出。
+
+### tables_configs 条目选项
+
+| 名称 | 类型 | 必填 | 描述
|
+|-------------|--------|------|-----------------------------------------------------------------------|
+| table_path | String | 是 | 格式:`database.resource`,例如
`google_ads.campaign`。 |
+| fields | List | 否* | 该表的有序 GAQL 字段路径列表。
|
+| query | String | 否* | 该表的完整 GAQL 查询,其 `FROM` 资源必须与 `table_path` 一致。
|
+| filter | String | 否 | GAQL WHERE 子句(仅与 `fields` 搭配使用)。
|
+| customer_id | String | 否 | 该表的 customer ID 覆盖值,缺省时回退到全局 `customer_id`。
|
+
+\* 每个条目的 `fields` 与 `query` 必须且只能提供一个。
+
+## 数据类型映射
+
+| Google Ads 数据类型 | SeaTunnel 数据类型 | 说明
|
+|-----------------------------|--------------------|------------------------------------------------------------------------|
+| INT64、UINT64 | BIGINT | REST API 将 int64 序列化为 JSON
字符串,连接器解析为 long。 |
+| INT32 | INT |
|
+| DOUBLE、FLOAT | DOUBLE |
|
+| BOOLEAN | BOOLEAN |
|
+| DATE | STRING |
有意为之:日期类字段格式不统一(`2026-09-01`、`2026-09`、`2026-36`)。 |
+| STRING、ENUM、RESOURCE_NAME | STRING | 枚举为符号字符串,例如 `ENABLED`。
|
+| MESSAGE | STRING | 嵌套 message 以 JSON 文本输出。
|
+| 其他 / 未知 | STRING | 向前兼容兜底。
|
+
+结果行中缺失的字段(API 会整体省略空字段)输出为 `null`。
+
+## 示例
+
+### 单资源
+
+```hocon
+source {
+ GoogleAds {
+ developer_token = "your_developer_token"
+ client_id = "your_client_id"
+ client_secret = "your_client_secret"
+ refresh_token = "your_refresh_token"
+ customer_id = "1234567890"
+
+ resource = "campaign"
+ fields = ["campaign.id", "campaign.name", "campaign.status",
"metrics.clicks", "metrics.impressions"]
+ filter = "segments.date DURING LAST_30_DAYS"
+ }
+}
+```
+
+### 完整 GAQL 查询
+
+```hocon
+source {
+ GoogleAds {
+ developer_token = "your_developer_token"
+ client_id = "your_client_id"
+ client_secret = "your_client_secret"
+ refresh_token = "your_refresh_token"
+ customer_id = "1234567890"
+
+ query = "SELECT ad_group.id, ad_group.name, metrics.clicks FROM ad_group
WHERE metrics.clicks > 0"
+ }
+}
+```
+
+### 多表(含表级 customer ID)
+
+```hocon
+source {
+ GoogleAds {
+ developer_token = "your_developer_token"
+ client_id = "your_client_id"
+ client_secret = "your_client_secret"
+ refresh_token = "your_refresh_token"
+ customer_id = "1234567890"
+ login_customer_id = "9876543210"
+
+ tables_configs = [
+ {
+ table_path = "google_ads.campaign"
+ fields = ["campaign.id", "campaign.name", "metrics.cost_micros"]
+ filter = "segments.date DURING LAST_7_DAYS"
+ },
+ {
+ table_path = "google_ads.ad_group"
+ query = "SELECT ad_group.id, ad_group.name FROM ad_group"
+ customer_id = "2345678901"
+ }
+ ]
+ }
+}
+```
+
+## 限制
+
+- 仅支持批处理;不支持增量/CDC 读取(可用 `filter` 搭配日期 segment 做窗口抽取)。
+- 不支持并行读取;每个作业以单 split 读取。
+- 不支持精确一次语义;重跑作业会重新读取数据。
+- 嵌套 MESSAGE 字段以 JSON 字符串输出,不展开为嵌套行。
+
+## 变更日志
+
+### 下一版本
+
+- 新增 Google Ads 源连接器,支持 GAQL、schema 自动推导与多表读取
diff --git a/plugin-mapping.properties b/plugin-mapping.properties
index 1ac16a0fc5..4af70e7ed3 100644
--- a/plugin-mapping.properties
+++ b/plugin-mapping.properties
@@ -103,6 +103,7 @@ seatunnel.source.AzureQueueStorage =
connector-azure-queue-storage
seatunnel.sink.AzureQueueStorage = connector-azure-queue-storage
seatunnel.source.GoogleBigtable = connector-google-bigtable
seatunnel.sink.GoogleBigtable = connector-google-bigtable
+seatunnel.source.GoogleAds = connector-google-ads
seatunnel.sink.Tablestore = connector-tablestore
seatunnel.source.Tablestore = connector-tablestore
seatunnel.source.Lemlist = connector-http-lemlist
diff --git a/seatunnel-connectors-v2/connector-google-ads/pom.xml
b/seatunnel-connectors-v2/connector-google-ads/pom.xml
new file mode 100644
index 0000000000..f3bf79d3ad
--- /dev/null
+++ b/seatunnel-connectors-v2/connector-google-ads/pom.xml
@@ -0,0 +1,55 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+ ~ 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.
+ -->
+<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
+ <modelVersion>4.0.0</modelVersion>
+ <parent>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>seatunnel-connectors-v2</artifactId>
+ <version>${revision}</version>
+ </parent>
+
+ <artifactId>connector-google-ads</artifactId>
+ <name>SeaTunnel : Connectors V2 : Google Ads</name>
+
+ <properties>
+ <httpclient.version>4.5.13</httpclient.version>
+ </properties>
+
+ <dependencies>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-common</artifactId>
+ <version>${project.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.httpcomponents</groupId>
+ <artifactId>httpclient</artifactId>
+ <version>${httpclient.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>com.fasterxml.jackson.core</groupId>
+ <artifactId>jackson-databind</artifactId>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.commons</groupId>
+ <artifactId>commons-lang3</artifactId>
+ </dependency>
+ </dependencies>
+
+</project>
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/client/GoogleAdsClient.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/client/GoogleAdsClient.java
new file mode 100644
index 0000000000..5f8e5ce55e
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/client/GoogleAdsClient.java
@@ -0,0 +1,402 @@
+/*
+ * 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.ads.client;
+
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.JsonNode;
+import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.ObjectMapper;
+import
org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ObjectNode;
+
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsParameters;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorException;
+
+import org.apache.http.HttpHeaders;
+import org.apache.http.NameValuePair;
+import org.apache.http.client.config.RequestConfig;
+import org.apache.http.client.entity.UrlEncodedFormEntity;
+import org.apache.http.client.methods.CloseableHttpResponse;
+import org.apache.http.client.methods.HttpPost;
+import org.apache.http.entity.ContentType;
+import org.apache.http.entity.StringEntity;
+import org.apache.http.impl.client.CloseableHttpClient;
+import org.apache.http.impl.client.HttpClients;
+import org.apache.http.message.BasicNameValuePair;
+import org.apache.http.util.EntityUtils;
+
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.Closeable;
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.function.Consumer;
+
+@Slf4j
+public class GoogleAdsClient implements Closeable {
+
+ private static final String TOKEN_PATH = "/token";
+ private static final String FIELDS_SEARCH_PATH =
"/%s/googleAdsFields:search";
+ private static final String SEARCH_PATH =
"/%s/customers/%s/googleAds:search";
+ private static final String PLUGIN_NAME = "GoogleAds";
+ private static final long TOKEN_EXPIRY_SAFETY_MARGIN_MS = 60_000L;
+
+ private final GoogleAdsParameters params;
+ private final CloseableHttpClient httpClient;
+ private final ObjectMapper objectMapper = new ObjectMapper();
+ private String accessToken;
+ private long tokenExpiresAtMs;
+
+ public GoogleAdsClient(GoogleAdsParameters params) {
+ this.params = params;
+ RequestConfig requestConfig =
+ RequestConfig.custom()
+ .setConnectTimeout(params.getRequestTimeoutMs())
+ .setSocketTimeout(params.getRequestTimeoutMs())
+ .build();
+ this.httpClient =
HttpClients.custom().setDefaultRequestConfig(requestConfig).build();
+ }
+
+ public void authenticate() {
+ String tokenUrl = params.getOauthEndpoint() + TOKEN_PATH;
+ HttpPost post = new HttpPost(tokenUrl);
+
+ List<NameValuePair> form = new ArrayList<>();
+ form.add(new BasicNameValuePair("grant_type", "refresh_token"));
+ form.add(new BasicNameValuePair("client_id", params.getClientId()));
+ form.add(new BasicNameValuePair("client_secret",
params.getClientSecret()));
+ form.add(new BasicNameValuePair("refresh_token",
params.getRefreshToken()));
+
+ try {
+ post.setEntity(new UrlEncodedFormEntity(form,
StandardCharsets.UTF_8));
+ try (CloseableHttpResponse response = httpClient.execute(post)) {
+ int status = response.getStatusLine().getStatusCode();
+ String body = EntityUtils.toString(response.getEntity(),
StandardCharsets.UTF_8);
+ if (status != 200) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.AUTH_FAILED,
+ "HTTP " + status + ": " + body);
+ }
+ JsonNode json = objectMapper.readTree(body);
+ this.accessToken = json.get("access_token").asText();
+ long expiresInSec =
+ json.has("expires_in") ?
json.get("expires_in").asLong() : 3600L;
+ this.tokenExpiresAtMs =
+ System.currentTimeMillis()
+ + expiresInSec * 1000
+ - TOKEN_EXPIRY_SAFETY_MARGIN_MS;
+ log.info("Obtained Google OAuth2 access token");
+ }
+ } catch (GoogleAdsConnectorException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new
GoogleAdsConnectorException(GoogleAdsConnectorErrorCode.AUTH_FAILED, e);
+ }
+ }
+
+ /**
+ * Resolves the data type of every requested field via the
GoogleAdsFieldService and builds a
+ * CatalogTable whose column order follows fieldPaths (the GAQL SELECT
order), never the
+ * metadata response order, so schema position i always corresponds to
selected field i.
+ */
+ public CatalogTable describeFields(String database, String resource,
List<String> fieldPaths) {
+ StringBuilder inList = new StringBuilder();
+ for (int i = 0; i < fieldPaths.size(); i++) {
+ if (i > 0) {
+ inList.append(", ");
+ }
+ inList.append('\'').append(fieldPaths.get(i)).append('\'');
+ }
+ String metaQuery =
+ "SELECT name, data_type FROM google_ads_field WHERE name IN ("
+ inList + ")";
+
+ String url =
+ params.getApiEndpoint() + String.format(FIELDS_SEARCH_PATH,
params.getApiVersion());
+ ObjectNode requestBody = objectMapper.createObjectNode();
+ requestBody.put("query", metaQuery);
+
+ String body =
+ executeJsonPost(
+ url, requestBody,
GoogleAdsConnectorErrorCode.DESCRIBE_FIELDS_FAILED);
+ try {
+ JsonNode json = objectMapper.readTree(body);
+ Map<String, String> nameToType = new HashMap<>();
+ JsonNode results = json.get("results");
+ if (results != null) {
+ for (JsonNode field : results) {
+ nameToType.put(field.get("name").asText(),
field.get("dataType").asText());
+ }
+ }
+
+ TableSchema.Builder schemaBuilder = TableSchema.builder();
+ for (String fieldPath : fieldPaths) {
+ String dataType = nameToType.get(fieldPath);
+ if (dataType == null) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.DESCRIBE_FIELDS_FAILED,
+ "Unknown Google Ads field: "
+ + fieldPath
+ + ". Check the field name against the "
+ + resource
+ + " resource documentation.");
+ }
+ schemaBuilder.column(
+ PhysicalColumn.of(
+ fieldPath,
+ mapGoogleAdsType(dataType),
+ null,
+ null,
+ true,
+ null,
+ null));
+ }
+ return CatalogTable.of(
+ TableIdentifier.of(PLUGIN_NAME, database, resource),
+ schemaBuilder.build(),
+ Collections.emptyMap(),
+ Collections.emptyList(),
+ "");
+ } catch (GoogleAdsConnectorException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.DESCRIBE_FIELDS_FAILED, e);
+ }
+ }
+
+ private SeaTunnelDataType<?> mapGoogleAdsType(String dataType) {
+ switch (dataType) {
+ case "INT64":
+ case "UINT64":
+ return BasicType.LONG_TYPE;
+ case "INT32":
+ return BasicType.INT_TYPE;
+ case "DOUBLE":
+ case "FLOAT":
+ return BasicType.DOUBLE_TYPE;
+ case "BOOLEAN":
+ return BasicType.BOOLEAN_TYPE;
+ default:
+ // DATE (non-uniform formats like 2026-09 for segments.month),
STRING, ENUM,
+ // RESOURCE_NAME, MESSAGE (emitted as JSON text) and unknown
future types
+ return BasicType.STRING_TYPE;
+ }
+ }
+
+ /**
+ * Runs one GAQL query end-to-end, walking forward via nextPageToken. Each
result row is
+ * flattened to a typed Object[] aligned with fieldPaths and pushed to
rowConsumer immediately;
+ * only one page's JSON body is held in memory at a time.
+ */
+ public void search(
+ String customerId,
+ String gaql,
+ List<String> fieldPaths,
+ SeaTunnelRowType rowType,
+ Consumer<Object[]> rowConsumer) {
+ String url =
+ params.getApiEndpoint()
+ + String.format(SEARCH_PATH, params.getApiVersion(),
customerId);
+ String pageToken = null;
+
+ do {
+ ObjectNode requestBody = objectMapper.createObjectNode();
+ requestBody.put("query", gaql);
+ if (pageToken != null) {
+ requestBody.put("pageToken", pageToken);
+ }
+ if (params.getPageSize() != null) {
+ requestBody.put("pageSize", params.getPageSize());
+ }
+
+ String body =
+ executeJsonPost(url, requestBody,
GoogleAdsConnectorErrorCode.SEARCH_FAILED);
+ try {
+ JsonNode json = objectMapper.readTree(body);
+ JsonNode results = json.get("results");
+ if (results != null) {
+ for (JsonNode result : results) {
+ Object[] row = new Object[fieldPaths.size()];
+ for (int i = 0; i < fieldPaths.size(); i++) {
+ JsonNode node = resolvePath(result,
fieldPaths.get(i));
+ row[i] =
+ (node == null || node.isNull() ||
node.isMissingNode())
+ ? null
+ : extractValue(node,
rowType.getFieldType(i));
+ }
+ rowConsumer.accept(row);
+ }
+ }
+ JsonNode next = json.get("nextPageToken");
+ pageToken = (next != null && !next.asText().isEmpty()) ?
next.asText() : null;
+ } catch (GoogleAdsConnectorException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new
GoogleAdsConnectorException(GoogleAdsConnectorErrorCode.SEARCH_FAILED, e);
+ }
+ } while (pageToken != null);
+ }
+
+ /**
+ * GAQL field paths are snake_case (ad_group.cpc_bid_micros) but the REST
response uses protobuf
+ * JSON camelCase keys (adGroup.cpcBidMicros), so each segment is
converted before lookup.
+ * Absent fields (Google omits empty proto fields entirely) resolve to
null.
+ */
+ private JsonNode resolvePath(JsonNode result, String fieldPath) {
+ JsonNode node = result;
+ for (String segment : fieldPath.split("\\.")) {
+ if (node == null) {
+ return null;
+ }
+ node = node.get(snakeToCamel(segment));
+ }
+ return node;
+ }
+
+ private String snakeToCamel(String segment) {
+ if (segment.indexOf('_') < 0) {
+ return segment;
+ }
+ StringBuilder sb = new StringBuilder(segment.length());
+ boolean upperNext = false;
+ for (int i = 0; i < segment.length(); i++) {
+ char c = segment.charAt(i);
+ if (c == '_') {
+ upperNext = true;
+ } else {
+ sb.append(upperNext ? Character.toUpperCase(c) : c);
+ upperNext = false;
+ }
+ }
+ return sb.toString();
+ }
+
+ private Object extractValue(JsonNode node, SeaTunnelDataType<?>
targetType) {
+ switch (targetType.getSqlType()) {
+ case BIGINT:
+ // protobuf JSON serializes int64 as a string
+ return Long.parseLong(node.asText());
+ case INT:
+ return node.asInt();
+ case DOUBLE:
+ return node.asDouble();
+ case BOOLEAN:
+ return node.asBoolean();
+ default:
+ return node.isValueNode() ? node.asText() : node.toString();
+ }
+ }
+
+ /**
+ * POSTs a JSON body with Google Ads headers. On 401 the access token is
refreshed once and the
+ * request replayed (not counted against max_retries). Transient failures
(429/5xx/IO errors)
+ * are retried up to max_retries with exponential backoff; other non-200
statuses fail fast with
+ * the API error body surfaced.
+ */
+ private String executeJsonPost(
+ String url, ObjectNode requestBody, GoogleAdsConnectorErrorCode
errorCode) {
+ boolean tokenRefreshed = false;
+ int attempt = 0;
+ while (true) {
+ ensureToken();
+ try {
+ HttpPost post = new HttpPost(url);
+ post.setHeader(HttpHeaders.AUTHORIZATION, "Bearer " +
accessToken);
+ post.setHeader("developer-token", params.getDeveloperToken());
+ if (params.getLoginCustomerId() != null) {
+ post.setHeader("login-customer-id",
params.getLoginCustomerId());
+ }
+ post.setEntity(
+ new StringEntity(
+ objectMapper.writeValueAsString(requestBody),
+ ContentType.APPLICATION_JSON));
+
+ try (CloseableHttpResponse response =
httpClient.execute(post)) {
+ int status = response.getStatusLine().getStatusCode();
+ String body =
+ EntityUtils.toString(response.getEntity(),
StandardCharsets.UTF_8);
+ if (status == 200) {
+ return body;
+ }
+ if (status == 401 && !tokenRefreshed) {
+ log.info("Access token rejected (401); refreshing and
replaying request");
+ tokenRefreshed = true;
+ authenticate();
+ continue;
+ }
+ if (isTransient(status) && attempt <
params.getMaxRetries()) {
+ backoff(++attempt, "HTTP " + status);
+ continue;
+ }
+ throw new GoogleAdsConnectorException(
+ errorCode, "HTTP " + status + ": " + body);
+ }
+ } catch (GoogleAdsConnectorException e) {
+ throw e;
+ } catch (IOException e) {
+ if (attempt < params.getMaxRetries()) {
+ backoff(++attempt, e.getMessage());
+ continue;
+ }
+ throw new GoogleAdsConnectorException(errorCode, e);
+ }
+ }
+ }
+
+ private boolean isTransient(int status) {
+ return status == 429 || (status >= 500 && status <= 504);
+ }
+
+ private void backoff(int attempt, String reason) {
+ long sleepMs = params.getRetryBackoffMs() * (1L << (attempt - 1));
+ log.warn(
+ "Transient Google Ads API failure ({}); retry {}/{} after
{}ms",
+ reason,
+ attempt,
+ params.getMaxRetries(),
+ sleepMs);
+ try {
+ Thread.sleep(sleepMs);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.SEARCH_FAILED, "Interrupted
during retry backoff");
+ }
+ }
+
+ private void ensureToken() {
+ if (accessToken == null || System.currentTimeMillis() >=
tokenExpiresAtMs) {
+ authenticate();
+ }
+ }
+
+ @Override
+ public void close() throws IOException {
+ httpClient.close();
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsParameters.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsParameters.java
new file mode 100644
index 0000000000..bdc81fd24e
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsParameters.java
@@ -0,0 +1,59 @@
+/*
+ * 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.ads.config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+
+import lombok.Getter;
+
+import java.io.Serializable;
+
+@Getter
+public class GoogleAdsParameters implements Serializable {
+
+ private String developerToken;
+ private String clientId;
+ private String clientSecret;
+ private String refreshToken;
+ private String customerId;
+ private String loginCustomerId;
+ private String apiVersion;
+ private int requestTimeoutMs;
+ private int maxRetries;
+ private long retryBackoffMs;
+ private Integer pageSize;
+ private String apiEndpoint;
+ private String oauthEndpoint;
+
+ public void buildWithConfig(ReadonlyConfig config) {
+ this.developerToken =
config.get(GoogleAdsSourceOptions.DEVELOPER_TOKEN);
+ this.clientId = config.get(GoogleAdsSourceOptions.CLIENT_ID);
+ this.clientSecret = config.get(GoogleAdsSourceOptions.CLIENT_SECRET);
+ this.refreshToken = config.get(GoogleAdsSourceOptions.REFRESH_TOKEN);
+ this.customerId = config.get(GoogleAdsSourceOptions.CUSTOMER_ID);
+ this.loginCustomerId =
+
config.getOptional(GoogleAdsSourceOptions.LOGIN_CUSTOMER_ID).orElse(null);
+ this.apiVersion = config.get(GoogleAdsSourceOptions.API_VERSION);
+ this.requestTimeoutMs =
config.get(GoogleAdsSourceOptions.REQUEST_TIMEOUT_MS);
+ this.maxRetries = config.get(GoogleAdsSourceOptions.MAX_RETRIES);
+ this.retryBackoffMs =
config.get(GoogleAdsSourceOptions.RETRY_BACKOFF_MS);
+ this.pageSize =
config.getOptional(GoogleAdsSourceOptions.PAGE_SIZE).orElse(null);
+ this.apiEndpoint = config.get(GoogleAdsSourceOptions.API_ENDPOINT);
+ this.oauthEndpoint = config.get(GoogleAdsSourceOptions.OAUTH_ENDPOINT);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsSourceOptions.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsSourceOptions.java
new file mode 100644
index 0000000000..9eac2a4b5e
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsSourceOptions.java
@@ -0,0 +1,160 @@
+/*
+ * 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.ads.config;
+
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.Options;
+
+import java.util.List;
+
+public class GoogleAdsSourceOptions {
+
+ public static final Option<String> DEVELOPER_TOKEN =
+ Options.key("developer_token")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "Google Ads API developer token, obtained from a
Google Ads "
+ + "manager (MCC) account.");
+
+ public static final Option<String> CLIENT_ID =
+ Options.key("client_id")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("OAuth2 client ID of the Google Cloud
application.");
+
+ public static final Option<String> CLIENT_SECRET =
+ Options.key("client_secret")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("OAuth2 client secret of the Google Cloud
application.");
+
+ public static final Option<String> REFRESH_TOKEN =
+ Options.key("refresh_token")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "OAuth2 refresh token authorized for the Google
Ads API scope "
+ +
"(https://www.googleapis.com/auth/adwords).");
+
+ public static final Option<String> CUSTOMER_ID =
+ Options.key("customer_id")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "Google Ads customer ID to query, digits only
without dashes, "
+ + "e.g. 1234567890.");
+
+ public static final Option<String> LOGIN_CUSTOMER_ID =
+ Options.key("login_customer_id")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "Manager (MCC) account customer ID, digits only.
Required when "
+ + "customer_id is a client account managed
by an MCC.");
+
+ public static final Option<String> API_VERSION =
+ Options.key("api_version")
+ .stringType()
+ .defaultValue("v21")
+ .withDescription("Google Ads REST API version, e.g. v21.");
+
+ public static final Option<String> RESOURCE =
+ Options.key("resource")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "Google Ads resource to query in single-table
mode, e.g. campaign, "
+ + "ad_group, keyword_view. Mutually
exclusive with query and "
+ + "tables_configs.");
+
+ public static final Option<List<String>> FIELDS =
+ Options.key("fields")
+ .listType()
+ .noDefaultValue()
+ .withDescription(
+ "Ordered list of GAQL field paths to select, e.g. "
+ + "[campaign.id, campaign.name,
metrics.clicks]. The output "
+ + "schema follows this order.");
+
+ public static final Option<String> FILTER =
+ Options.key("filter")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "GAQL WHERE clause appended to the auto-built "
+ + "SELECT <fields> FROM <resource> query,
e.g. "
+ + "segments.date DURING LAST_30_DAYS.");
+
+ public static final Option<String> QUERY =
+ Options.key("query")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "Full GAQL query. Mutually exclusive with
resource/fields/filter. "
+ + "The output schema follows the SELECT
field order.");
+
+ public static final Option<Integer> REQUEST_TIMEOUT_MS =
+ Options.key("request_timeout_ms")
+ .intType()
+ .defaultValue(60000)
+ .withDescription("HTTP request timeout in milliseconds for
a single call.");
+
+ public static final Option<Integer> MAX_RETRIES =
+ Options.key("max_retries")
+ .intType()
+ .defaultValue(3)
+ .withDescription(
+ "Maximum retries for transient HTTP failures "
+ + "(429/5xx/network errors) of a single
request.");
+
+ public static final Option<Long> RETRY_BACKOFF_MS =
+ Options.key("retry_backoff_ms")
+ .longType()
+ .defaultValue(1000L)
+ .withDescription(
+ "Base backoff in milliseconds between retries;
doubled per attempt.");
+
+ public static final Option<Integer> PAGE_SIZE =
+ Options.key("page_size")
+ .intType()
+ .noDefaultValue()
+ .withDescription(
+ "Page size for search requests. When unset the
server default "
+ + "is used.");
+
+ public static final Option<String> TABLE_PATH =
+ Options.key("table_path")
+ .stringType()
+ .noDefaultValue()
+ .withDescription(
+ "Table path in 'database.resource' format, e.g.
google_ads.campaign. "
+ + "Used in tables_configs entries.");
+
+ public static final Option<String> API_ENDPOINT =
+ Options.key("api_endpoint")
+ .stringType()
+ .defaultValue("https://googleads.googleapis.com")
+ .withDescription("Google Ads API base endpoint. Intended
for testing.");
+
+ public static final Option<String> OAUTH_ENDPOINT =
+ Options.key("oauth_endpoint")
+ .stringType()
+ .defaultValue("https://oauth2.googleapis.com")
+ .withDescription("Google OAuth2 token endpoint base.
Intended for testing.");
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsTableConfig.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsTableConfig.java
new file mode 100644
index 0000000000..202a96cfad
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/config/GoogleAdsTableConfig.java
@@ -0,0 +1,53 @@
+/*
+ * 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.ads.config;
+
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+
+import lombok.Getter;
+
+import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.List;
+
+@Getter
+public class GoogleAdsTableConfig implements Serializable {
+
+ private final String gaql;
+ private final String resource;
+ private final String customerId;
+ private final ArrayList<String> fieldPaths;
+ private final CatalogTable catalogTable;
+
+ public GoogleAdsTableConfig(
+ String gaql,
+ String resource,
+ String customerId,
+ List<String> fieldPaths,
+ CatalogTable catalogTable) {
+ this.gaql = gaql;
+ this.resource = resource;
+ this.customerId = customerId;
+ this.fieldPaths = new ArrayList<>(fieldPaths);
+ this.catalogTable = catalogTable;
+ }
+
+ public String getTableId() {
+ return catalogTable.getTableId().toTablePath().toString();
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/exception/GoogleAdsConnectorErrorCode.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/exception/GoogleAdsConnectorErrorCode.java
new file mode 100644
index 0000000000..476464077a
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/exception/GoogleAdsConnectorErrorCode.java
@@ -0,0 +1,47 @@
+/*
+ * 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.ads.exception;
+
+import org.apache.seatunnel.common.exception.SeaTunnelErrorCode;
+
+public enum GoogleAdsConnectorErrorCode implements SeaTunnelErrorCode {
+ AUTH_FAILED("GOOGLE_ADS-01", "Failed to obtain Google OAuth2 access
token"),
+ DESCRIBE_FIELDS_FAILED("GOOGLE_ADS-02", "Failed to resolve Google Ads
field metadata"),
+ SEARCH_FAILED("GOOGLE_ADS-03", "Google Ads search request failed"),
+ INVALID_QUERY("GOOGLE_ADS-04", "Invalid GAQL query or fields
configuration"),
+ INVALID_TABLE_PATH("GOOGLE_ADS-05", "Invalid table_path; expected format:
database.resource"),
+ DUPLICATE_RESOURCE("GOOGLE_ADS-06", "Duplicate table found in
tables_configs");
+
+ private final String code;
+ private final String description;
+
+ GoogleAdsConnectorErrorCode(String code, String description) {
+ this.code = code;
+ this.description = description;
+ }
+
+ @Override
+ public String getCode() {
+ return code;
+ }
+
+ @Override
+ public String getDescription() {
+ return description;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/exception/GoogleAdsConnectorException.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/exception/GoogleAdsConnectorException.java
new file mode 100644
index 0000000000..69807cf597
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/exception/GoogleAdsConnectorException.java
@@ -0,0 +1,36 @@
+/*
+ * 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.ads.exception;
+
+import org.apache.seatunnel.common.exception.SeaTunnelErrorCode;
+import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
+
+public class GoogleAdsConnectorException extends SeaTunnelRuntimeException {
+ public GoogleAdsConnectorException(SeaTunnelErrorCode seaTunnelErrorCode,
String errorMessage) {
+ super(seaTunnelErrorCode, errorMessage);
+ }
+
+ public GoogleAdsConnectorException(
+ SeaTunnelErrorCode seaTunnelErrorCode, String errorMessage,
Throwable cause) {
+ super(seaTunnelErrorCode, errorMessage, cause);
+ }
+
+ public GoogleAdsConnectorException(SeaTunnelErrorCode seaTunnelErrorCode,
Throwable cause) {
+ super(seaTunnelErrorCode, cause);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSource.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSource.java
new file mode 100644
index 0000000000..c07cfb271d
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSource.java
@@ -0,0 +1,251 @@
+/*
+ * 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.ads.source;
+
+import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+import org.apache.seatunnel.api.source.Boundedness;
+import org.apache.seatunnel.api.source.SupportColumnProjection;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitSource;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.client.GoogleAdsClient;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsParameters;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsSourceOptions;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsTableConfig;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorException;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.regex.Matcher;
+import java.util.regex.Pattern;
+import java.util.stream.Collectors;
+
+public class GoogleAdsSource extends AbstractSingleSplitSource<SeaTunnelRow>
+ implements SupportColumnProjection {
+
+ private static final String PLUGIN_NAME = "GoogleAds";
+ private static final String DEFAULT_DATABASE = "google_ads";
+ private static final Pattern FIELD_PATH_PATTERN =
+ Pattern.compile("^[a-z0-9_]+(\\.[a-z0-9_]+)+$");
+ private static final Pattern GAQL_PATTERN =
+ Pattern.compile(
+ "^\\s*SELECT\\s+(.+?)\\s+FROM\\s+([a-z0-9_]+)\\b.*$",
+ Pattern.CASE_INSENSITIVE | Pattern.DOTALL);
+
+ private final GoogleAdsParameters params;
+ private final List<GoogleAdsTableConfig> tableConfigs;
+
+ public GoogleAdsSource(GoogleAdsParameters params, ReadonlyConfig config) {
+ this.params = params;
+ this.tableConfigs = buildTableConfigs(params, config);
+ }
+
+ /**
+ * Resolves resource/fields/query/tables_configs into one or more
GoogleAdsTableConfig
+ * instances, deriving each CatalogTable from the GoogleAdsFieldService.
Runs once during
+ * factory createSource with a one-shot client scoped to this call.
+ */
+ private List<GoogleAdsTableConfig> buildTableConfigs(
+ GoogleAdsParameters params, ReadonlyConfig config) {
+ try (GoogleAdsClient client = new GoogleAdsClient(params)) {
+ client.authenticate();
+
+ if
(config.getOptional(ConnectorCommonOptions.TABLE_CONFIGS).isPresent()) {
+ List<Map<String, Object>> tableConfigMaps =
+ config.get(ConnectorCommonOptions.TABLE_CONFIGS);
+ List<GoogleAdsTableConfig> configs = new ArrayList<>();
+ for (Map<String, Object> map : tableConfigMaps) {
+ ReadonlyConfig tableConfig = ReadonlyConfig.fromMap(map);
+ GoogleAdsTableConfig built =
buildOneTableConfig(tableConfig, client);
+ String tableId = built.getTableId();
+ boolean duplicate =
+ configs.stream().anyMatch(c ->
c.getTableId().equals(tableId));
+ if (duplicate) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.DUPLICATE_RESOURCE,
+ "Duplicate table in tables_configs: " +
tableId);
+ }
+ configs.add(built);
+ }
+ return configs;
+ } else {
+ return Collections.singletonList(
+ buildSingleTableConfig(config, DEFAULT_DATABASE, null,
client));
+ }
+ } catch (GoogleAdsConnectorException e) {
+ throw e;
+ } catch (Exception e) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.DESCRIBE_FIELDS_FAILED,
+ "Failed to build Google Ads table configs",
+ e);
+ }
+ }
+
+ private GoogleAdsTableConfig buildOneTableConfig(
+ ReadonlyConfig tableConfig, GoogleAdsClient client) {
+ String tablePath =
+ tableConfig
+ .getOptional(GoogleAdsSourceOptions.TABLE_PATH)
+ .orElseThrow(
+ () ->
+ new GoogleAdsConnectorException(
+
GoogleAdsConnectorErrorCode.INVALID_TABLE_PATH,
+ "table_path is required in
tables_configs"));
+ String[] parts = tablePath.split("\\.", 2);
+ if (parts.length != 2 || StringUtils.isBlank(parts[0]) ||
StringUtils.isBlank(parts[1])) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.INVALID_TABLE_PATH,
+ "table_path must be 'database.resource', got: " +
tablePath);
+ }
+ return buildSingleTableConfig(tableConfig, parts[0], parts[1], client);
+ }
+
+ /**
+ * Builds one table from either fields mode (resource + fields [+ filter])
or query mode (full
+ * GAQL). The ordered field list parsed here is the single source of truth
shared by the schema
+ * (describeFields), the query, and row value extraction. expectedResource
is non-null only in
+ * tables_configs mode, where the table_path resource must agree with the
query.
+ */
+ private GoogleAdsTableConfig buildSingleTableConfig(
+ ReadonlyConfig config,
+ String database,
+ String expectedResource,
+ GoogleAdsClient client) {
+ String query =
config.getOptional(GoogleAdsSourceOptions.QUERY).orElse(null);
+ List<String> fields =
config.getOptional(GoogleAdsSourceOptions.FIELDS).orElse(null);
+ String filter =
config.getOptional(GoogleAdsSourceOptions.FILTER).orElse(null);
+
+ String gaql;
+ String resource;
+ List<String> fieldPaths;
+ if (query != null) {
+ if (fields != null || filter != null) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
+ "query is mutually exclusive with fields and filter");
+ }
+ Matcher matcher = GAQL_PATTERN.matcher(query);
+ if (!matcher.matches()) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
+ "Cannot parse GAQL query; expected 'SELECT <fields>
FROM <resource> ...', "
+ + "got: "
+ + query);
+ }
+ fieldPaths = new ArrayList<>();
+ for (String token : matcher.group(1).split(",")) {
+ fieldPaths.add(validateFieldPath(token.trim()));
+ }
+ resource = matcher.group(2);
+ if (expectedResource != null &&
!expectedResource.equals(resource)) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
+ "table_path resource '"
+ + expectedResource
+ + "' does not match query FROM resource '"
+ + resource
+ + "'");
+ }
+ gaql = query;
+ } else {
+ resource =
+ expectedResource != null
+ ? expectedResource
+ :
config.getOptional(GoogleAdsSourceOptions.RESOURCE)
+ .orElseThrow(
+ () ->
+ new
GoogleAdsConnectorException(
+
GoogleAdsConnectorErrorCode
+
.INVALID_QUERY,
+ "One of resource,
query or "
+ +
"tables_configs is "
+ +
"required"));
+ if (fields == null || fields.isEmpty()) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
+ "fields must be a non-empty list when using resource
mode");
+ }
+ fieldPaths = new ArrayList<>();
+ for (String field : fields) {
+ fieldPaths.add(validateFieldPath(field.trim()));
+ }
+ StringBuilder sb =
+ new StringBuilder("SELECT ")
+ .append(String.join(", ", fieldPaths))
+ .append(" FROM ")
+ .append(resource);
+ if (StringUtils.isNotBlank(filter)) {
+ sb.append(" WHERE ").append(filter);
+ }
+ gaql = sb.toString();
+ }
+
+ String customerId =
+ config.getOptional(GoogleAdsSourceOptions.CUSTOMER_ID)
+ .orElse(params.getCustomerId());
+ CatalogTable catalogTable = client.describeFields(database, resource,
fieldPaths);
+ return new GoogleAdsTableConfig(gaql, resource, customerId,
fieldPaths, catalogTable);
+ }
+
+ List<GoogleAdsTableConfig> getTableConfigs() {
+ return tableConfigs;
+ }
+
+ private String validateFieldPath(String fieldPath) {
+ if (!FIELD_PATH_PATTERN.matcher(fieldPath).matches()) {
+ throw new GoogleAdsConnectorException(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
+ "Invalid GAQL field path: '"
+ + fieldPath
+ + "'. Expected dotted snake_case like campaign.id
or metrics.clicks.");
+ }
+ return fieldPath;
+ }
+
+ @Override
+ public String getPluginName() {
+ return PLUGIN_NAME;
+ }
+
+ @Override
+ public Boundedness getBoundedness() {
+ return Boundedness.BOUNDED;
+ }
+
+ @Override
+ public List<CatalogTable> getProducedCatalogTables() {
+ return tableConfigs.stream()
+ .map(GoogleAdsTableConfig::getCatalogTable)
+ .collect(Collectors.toList());
+ }
+
+ @Override
+ public AbstractSingleSplitReader<SeaTunnelRow> createReader(
+ SingleSplitReaderContext readerContext) throws Exception {
+ return new GoogleAdsSourceReader(params, tableConfigs, readerContext);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceFactory.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceFactory.java
new file mode 100644
index 0000000000..086da2e469
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceFactory.java
@@ -0,0 +1,82 @@
+/*
+ * 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.ads.source;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+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.connectors.seatunnel.google.ads.config.GoogleAdsParameters;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsSourceOptions;
+
+import com.google.auto.service.AutoService;
+
+import java.io.Serializable;
+
+@AutoService(Factory.class)
+public class GoogleAdsSourceFactory implements TableSourceFactory {
+
+ @Override
+ public String factoryIdentifier() {
+ return "GoogleAds";
+ }
+
+ @Override
+ public OptionRule optionRule() {
+ return OptionRule.builder()
+ .required(
+ GoogleAdsSourceOptions.DEVELOPER_TOKEN,
+ GoogleAdsSourceOptions.CLIENT_ID,
+ GoogleAdsSourceOptions.CLIENT_SECRET,
+ GoogleAdsSourceOptions.REFRESH_TOKEN,
+ GoogleAdsSourceOptions.CUSTOMER_ID)
+ .exclusive(
+ GoogleAdsSourceOptions.RESOURCE,
+ GoogleAdsSourceOptions.QUERY,
+ ConnectorCommonOptions.TABLE_CONFIGS)
+ .optional(
+ GoogleAdsSourceOptions.LOGIN_CUSTOMER_ID,
+ GoogleAdsSourceOptions.API_VERSION,
+ GoogleAdsSourceOptions.FIELDS,
+ GoogleAdsSourceOptions.FILTER,
+ GoogleAdsSourceOptions.REQUEST_TIMEOUT_MS,
+ GoogleAdsSourceOptions.MAX_RETRIES,
+ GoogleAdsSourceOptions.RETRY_BACKOFF_MS,
+ GoogleAdsSourceOptions.PAGE_SIZE)
+ .build();
+ }
+
+ @Override
+ public <T, SplitT extends SourceSplit, StateT extends Serializable>
+ TableSource<T, SplitT, StateT>
createSource(TableSourceFactoryContext context) {
+ GoogleAdsParameters params = new GoogleAdsParameters();
+ params.buildWithConfig(context.getOptions());
+ return () ->
+ (SeaTunnelSource<T, SplitT, StateT>)
+ new GoogleAdsSource(params, context.getOptions());
+ }
+
+ @Override
+ public Class<? extends SeaTunnelSource> getSourceClass() {
+ return GoogleAdsSource.class;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceReader.java
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceReader.java
new file mode 100644
index 0000000000..801c11cc8a
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/main/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceReader.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.ads.source;
+
+import org.apache.seatunnel.api.source.Collector;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
+import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.client.GoogleAdsClient;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsParameters;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsTableConfig;
+
+import lombok.extern.slf4j.Slf4j;
+
+import java.io.IOException;
+import java.util.List;
+
+@Slf4j
+public class GoogleAdsSourceReader extends
AbstractSingleSplitReader<SeaTunnelRow> {
+
+ private final GoogleAdsParameters params;
+ private final List<GoogleAdsTableConfig> tableConfigs;
+ private final SingleSplitReaderContext readerContext;
+ private GoogleAdsClient client;
+
+ GoogleAdsSourceReader(
+ GoogleAdsParameters params,
+ List<GoogleAdsTableConfig> tableConfigs,
+ SingleSplitReaderContext readerContext) {
+ this.params = params;
+ this.tableConfigs = tableConfigs;
+ this.readerContext = readerContext;
+ }
+
+ @Override
+ public void open() throws Exception {
+ client = new GoogleAdsClient(params);
+ client.authenticate();
+ }
+
+ @Override
+ public void close() throws IOException {
+ if (client != null) {
+ client.close();
+ }
+ }
+
+ /**
+ * Single-pass bounded read for the assigned split. For each configured
table, runs its GAQL
+ * query and forwards each already-typed Object[] downstream as a
table-tagged SeaTunnelRow;
+ * rows stream through page by page without buffering the whole result
set. After every table
+ * has drained, signals no-more-elements so the framework can close the
split.
+ */
+ @Override
+ public void pollNext(Collector<SeaTunnelRow> output) throws Exception {
+ try {
+ for (GoogleAdsTableConfig tableConfig : tableConfigs) {
+ SeaTunnelRowType rowType =
tableConfig.getCatalogTable().getSeaTunnelRowType();
+ String tableId = tableConfig.getTableId();
+ log.info(
+ "Reading rows from Google Ads resource {} for customer
{}",
+ tableConfig.getResource(),
+ tableConfig.getCustomerId());
+ client.search(
+ tableConfig.getCustomerId(),
+ tableConfig.getGaql(),
+ tableConfig.getFieldPaths(),
+ rowType,
+ values -> {
+ SeaTunnelRow row = new SeaTunnelRow(values);
+ row.setTableId(tableId);
+ output.collect(row);
+ });
+ }
+ } finally {
+ readerContext.signalNoMoreElement();
+ }
+ }
+
+ @Override
+ public void notifyCheckpointComplete(long checkpointId) throws Exception {}
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/GoogleAdsSourceFactoryTest.java
b/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/GoogleAdsSourceFactoryTest.java
new file mode 100644
index 0000000000..a621e3781f
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/GoogleAdsSourceFactoryTest.java
@@ -0,0 +1,209 @@
+/*
+ * 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.ads;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.source.GoogleAdsSourceFactory;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+class GoogleAdsSourceFactoryTest {
+
+ private static final GoogleAdsSourceFactory FACTORY = new
GoogleAdsSourceFactory();
+
+ @Test
+ void testOptionRuleIsNotNull() {
+ Assertions.assertNotNull(FACTORY.optionRule());
+ }
+
+ @Test
+ void testFactoryIdentifier() {
+ Assertions.assertEquals("GoogleAds", FACTORY.factoryIdentifier());
+ }
+
+ @Test
+ void testResourceModeValid() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("resource", "campaign");
+ config.put("fields", Arrays.asList("campaign.id", "campaign.name",
"metrics.clicks"));
+
+ Assertions.assertDoesNotThrow(
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testResourceModeWithOptionalParamsValid() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("resource", "campaign");
+ config.put("fields", Arrays.asList("campaign.id"));
+ config.put("filter", "segments.date DURING LAST_30_DAYS");
+ config.put("login_customer_id", "9999999999");
+ config.put("api_version", "v21");
+ config.put("max_retries", 5);
+ config.put("request_timeout_ms", 30000);
+ config.put("retry_backoff_ms", 500L);
+ config.put("page_size", 1000);
+
+ Assertions.assertDoesNotThrow(
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testQueryModeValid() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("query", "SELECT campaign.id, metrics.clicks FROM
campaign");
+
+ Assertions.assertDoesNotThrow(
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testTablesConfigsModeValid() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("tables_configs", multiTableConfigs());
+
+ Assertions.assertDoesNotThrow(
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testResourceAndQueryTogetherThrows() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("resource", "campaign");
+ config.put("query", "SELECT campaign.id FROM campaign");
+
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testResourceAndTablesConfigsTogetherThrows() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("resource", "campaign");
+ config.put("tables_configs", multiTableConfigs());
+
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testQueryAndTablesConfigsTogetherThrows() {
+ Map<String, Object> config = requiredAuthConfig();
+ config.put("query", "SELECT campaign.id FROM campaign");
+ config.put("tables_configs", multiTableConfigs());
+
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testNoModeSelectedThrows() {
+ Map<String, Object> config = requiredAuthConfig();
+
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ @Test
+ void testMissingDeveloperTokenThrows() {
+ assertMissingRequiredThrows("developer_token");
+ }
+
+ @Test
+ void testMissingClientIdThrows() {
+ assertMissingRequiredThrows("client_id");
+ }
+
+ @Test
+ void testMissingClientSecretThrows() {
+ assertMissingRequiredThrows("client_secret");
+ }
+
+ @Test
+ void testMissingRefreshTokenThrows() {
+ assertMissingRequiredThrows("refresh_token");
+ }
+
+ @Test
+ void testMissingCustomerIdThrows() {
+ assertMissingRequiredThrows("customer_id");
+ }
+
+ private static void assertMissingRequiredThrows(String key) {
+ Map<String, Object> config = requiredAuthConfig();
+ config.remove(key);
+ config.put("query", "SELECT campaign.id FROM campaign");
+
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
+ ConfigValidator.of(ReadonlyConfig.fromMap(config))
+ .validate(FACTORY.optionRule()));
+ }
+
+ private static Map<String, Object> requiredAuthConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put("developer_token", "test_developer_token");
+ config.put("client_id", "test_client_id");
+ config.put("client_secret", "test_client_secret");
+ config.put("refresh_token", "test_refresh_token");
+ config.put("customer_id", "1234567890");
+ return config;
+ }
+
+ private static List<Map<String, Object>> multiTableConfigs() {
+ Map<String, Object> campaign = new HashMap<>();
+ campaign.put("table_path", "google_ads.campaign");
+ campaign.put("fields", Arrays.asList("campaign.id", "campaign.name"));
+
+ Map<String, Object> adGroup = new HashMap<>();
+ adGroup.put("table_path", "google_ads.ad_group");
+ adGroup.put("query", "SELECT ad_group.id, metrics.clicks FROM
ad_group");
+ adGroup.put("customer_id", "2345678901");
+
+ return Arrays.asList(campaign, adGroup);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/client/GoogleAdsClientTest.java
b/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/client/GoogleAdsClientTest.java
new file mode 100644
index 0000000000..e0c60648c4
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/client/GoogleAdsClientTest.java
@@ -0,0 +1,418 @@
+/*
+ * 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.ads.client;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.api.table.type.SqlType;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsParameters;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorException;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import com.sun.net.httpserver.HttpExchange;
+import com.sun.net.httpserver.HttpServer;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+class GoogleAdsClientTest {
+
+ private static final String CUSTOMER_ID = "1234567890";
+ private static final String FIELDS_PATH = "/v21/googleAdsFields:search";
+ private static final String SEARCH_PATH = "/v21/customers/" + CUSTOMER_ID
+ "/googleAds:search";
+
+ private HttpServer server;
+ private String baseUrl;
+
+ @BeforeEach
+ void setUp() throws IOException {
+ server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
+ server.start();
+ baseUrl = "http://127.0.0.1:" + server.getAddress().getPort();
+ }
+
+ @AfterEach
+ void tearDown() {
+ if (server != null) {
+ server.stop(0);
+ }
+ }
+
+ @Test
+ void authenticateFailureRaisesAuthException() {
+ server.createContext(
+ "/token", exchange -> respondJson(exchange, 400,
"{\"error\":\"invalid_grant\"}"));
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ GoogleAdsConnectorException ex =
+ Assertions.assertThrows(
+ GoogleAdsConnectorException.class,
client::authenticate);
+ Assertions.assertTrue(ex.getMessage().contains("400"));
+ } catch (IOException e) {
+ Assertions.fail(e);
+ }
+ }
+
+ @Test
+ void describeFieldsBuildsSchemaInRequestOrderDespiteShuffledResponse()
throws IOException {
+ registerToken();
+ server.createContext(
+ FIELDS_PATH,
+ exchange ->
+ respondJson(
+ exchange,
+ 200,
+ "{\"results\":["
+ + fieldMeta("metrics.clicks", "INT64")
+ + ","
+ + fieldMeta("campaign.name", "STRING")
+ + ","
+ + fieldMeta("segments.date", "DATE")
+ + ","
+ + fieldMeta("campaign.id", "INT64")
+ + ","
+ + fieldMeta("metrics.ctr", "DOUBLE")
+ + ","
+ + fieldMeta("campaign.status", "ENUM")
+ + "]}"));
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ CatalogTable table =
+ client.describeFields(
+ "google_ads",
+ "campaign",
+ Arrays.asList(
+ "campaign.id",
+ "campaign.name",
+ "campaign.status",
+ "metrics.clicks",
+ "metrics.ctr",
+ "segments.date"));
+ SeaTunnelRowType rowType = table.getSeaTunnelRowType();
+ Assertions.assertEquals(6, rowType.getTotalFields());
+ Assertions.assertEquals("campaign.id", rowType.getFieldName(0));
+ Assertions.assertEquals("campaign.name", rowType.getFieldName(1));
+ Assertions.assertEquals("campaign.status",
rowType.getFieldName(2));
+ Assertions.assertEquals("metrics.clicks", rowType.getFieldName(3));
+ Assertions.assertEquals("metrics.ctr", rowType.getFieldName(4));
+ Assertions.assertEquals("segments.date", rowType.getFieldName(5));
+ Assertions.assertEquals(SqlType.BIGINT,
rowType.getFieldType(0).getSqlType());
+ Assertions.assertEquals(SqlType.STRING,
rowType.getFieldType(1).getSqlType());
+ Assertions.assertEquals(SqlType.STRING,
rowType.getFieldType(2).getSqlType());
+ Assertions.assertEquals(SqlType.BIGINT,
rowType.getFieldType(3).getSqlType());
+ Assertions.assertEquals(SqlType.DOUBLE,
rowType.getFieldType(4).getSqlType());
+ // DATE deliberately maps to STRING (non-uniform formats like
2026-09)
+ Assertions.assertEquals(SqlType.STRING,
rowType.getFieldType(5).getSqlType());
+ }
+ }
+
+ @Test
+ void describeFieldsThrowsOnUnknownField() throws IOException {
+ registerToken();
+ server.createContext(
+ FIELDS_PATH,
+ exchange ->
+ respondJson(
+ exchange,
+ 200,
+ "{\"results\":[" + fieldMeta("campaign.id",
"INT64") + "]}"));
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ GoogleAdsConnectorException ex =
+ Assertions.assertThrows(
+ GoogleAdsConnectorException.class,
+ () ->
+ client.describeFields(
+ "google_ads",
+ "campaign",
+ Arrays.asList("campaign.id",
"campaign.namee")));
+ Assertions.assertTrue(ex.getMessage().contains("campaign.namee"));
+ }
+ }
+
+ @Test
+ void searchStreamsPagesAndConvertsTypesWithSnakeToCamelResolution() throws
IOException {
+ registerToken();
+ AtomicInteger searchHits = new AtomicInteger();
+ AtomicReference<String> secondRequestBody = new AtomicReference<>();
+ server.createContext(
+ SEARCH_PATH,
+ exchange -> {
+ String requestBody = readBody(exchange);
+ if (searchHits.incrementAndGet() == 1) {
+ respondJson(
+ exchange,
+ 200,
+ "{\"results\":["
+ +
"{\"adGroup\":{\"id\":\"111\",\"cpcBidMicros\":\"2500000\","
+ + "\"status\":\"ENABLED\"},"
+ +
"\"metrics\":{\"ctr\":0.052,\"clicks\":\"42\"}}"
+ + "],\"nextPageToken\":\"page-2\"}");
+ } else {
+ secondRequestBody.set(requestBody);
+ respondJson(
+ exchange,
+ 200,
+ "{\"results\":["
+ +
"{\"adGroup\":{\"id\":\"222\",\"status\":\"PAUSED\"},"
+ +
"\"metrics\":{\"ctr\":0.01,\"clicks\":\"7\"}}"
+ + "]}");
+ }
+ });
+
+ List<String> fieldPaths =
+ Arrays.asList(
+ "ad_group.id",
+ "ad_group.cpc_bid_micros",
+ "ad_group.status",
+ "metrics.ctr",
+ "metrics.clicks");
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ fieldPaths.toArray(new String[0]),
+ new SeaTunnelDataType[] {
+ BasicType.LONG_TYPE,
+ BasicType.LONG_TYPE,
+ BasicType.STRING_TYPE,
+ BasicType.DOUBLE_TYPE,
+ BasicType.LONG_TYPE
+ });
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ List<Object[]> rows = new ArrayList<>();
+ client.search(
+ CUSTOMER_ID,
+ "SELECT ad_group.id FROM ad_group",
+ fieldPaths,
+ rowType,
+ rows::add);
+
+ Assertions.assertEquals(2, rows.size());
+ Assertions.assertEquals(2, searchHits.get());
+ // string-encoded int64 parsed to Long; snake_case path resolved
via camelCase keys
+ Assertions.assertArrayEquals(
+ new Object[] {111L, 2500000L, "ENABLED", 0.052, 42L},
rows.get(0));
+ // absent cpc_bid_micros on page 2 -> null
+ Assertions.assertArrayEquals(
+ new Object[] {222L, null, "PAUSED", 0.01, 7L},
rows.get(1));
+
Assertions.assertTrue(secondRequestBody.get().contains("\"pageToken\":\"page-2\""));
+ }
+ }
+
+ @Test
+ void searchSendsRequiredGoogleAdsHeaders() throws IOException {
+ registerToken();
+ AtomicReference<String> authHeader = new AtomicReference<>();
+ AtomicReference<String> devTokenHeader = new AtomicReference<>();
+ AtomicReference<String> loginCidHeader = new AtomicReference<>();
+ server.createContext(
+ SEARCH_PATH,
+ exchange -> {
+
authHeader.set(exchange.getRequestHeaders().getFirst("Authorization"));
+
devTokenHeader.set(exchange.getRequestHeaders().getFirst("developer-token"));
+
loginCidHeader.set(exchange.getRequestHeaders().getFirst("login-customer-id"));
+ respondJson(exchange, 200, "{\"results\":[]}");
+ });
+
+ try (GoogleAdsClient client = new
GoogleAdsClient(buildParams("9999999999"))) {
+ client.search(
+ CUSTOMER_ID,
+ "SELECT campaign.id FROM campaign",
+ Arrays.asList("campaign.id"),
+ new SeaTunnelRowType(
+ new String[] {"campaign.id"},
+ new SeaTunnelDataType[] {BasicType.LONG_TYPE}),
+ row -> {});
+ Assertions.assertEquals("Bearer tok-1", authHeader.get());
+ Assertions.assertEquals("dev-token", devTokenHeader.get());
+ Assertions.assertEquals("9999999999", loginCidHeader.get());
+ }
+ }
+
+ @Test
+ void searchRefreshesTokenOnceAndReplaysOn401() throws IOException {
+ AtomicInteger tokenHits = new AtomicInteger();
+ server.createContext(
+ "/token",
+ exchange ->
+ respondJson(
+ exchange,
+ 200,
+ "{\"access_token\":\"tok-"
+ + tokenHits.incrementAndGet()
+ + "\",\"expires_in\":3600}"));
+ AtomicInteger searchHits = new AtomicInteger();
+ server.createContext(
+ SEARCH_PATH,
+ exchange -> {
+ if (searchHits.incrementAndGet() == 1) {
+ respondJson(exchange, 401,
"{\"error\":{\"status\":\"UNAUTHENTICATED\"}}");
+ } else {
+ Assertions.assertEquals(
+ "Bearer tok-2",
+
exchange.getRequestHeaders().getFirst("Authorization"));
+ respondJson(exchange, 200,
"{\"results\":[{\"campaign\":{\"id\":\"1\"}}]}");
+ }
+ });
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ List<Object[]> rows = new ArrayList<>();
+ client.search(
+ CUSTOMER_ID,
+ "SELECT campaign.id FROM campaign",
+ Arrays.asList("campaign.id"),
+ new SeaTunnelRowType(
+ new String[] {"campaign.id"},
+ new SeaTunnelDataType[] {BasicType.LONG_TYPE}),
+ rows::add);
+ Assertions.assertEquals(1, rows.size());
+ Assertions.assertEquals(2, searchHits.get());
+ Assertions.assertEquals(2, tokenHits.get());
+ }
+ }
+
+ @Test
+ void searchRetriesTransient503UpToMaxRetries() throws IOException {
+ registerToken();
+ AtomicInteger searchHits = new AtomicInteger();
+ server.createContext(
+ SEARCH_PATH,
+ exchange -> {
+ searchHits.incrementAndGet();
+ respondJson(exchange, 503, "{\"error\":\"unavailable\"}");
+ });
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ Assertions.assertThrows(
+ GoogleAdsConnectorException.class,
+ () ->
+ client.search(
+ CUSTOMER_ID,
+ "SELECT campaign.id FROM campaign",
+ Arrays.asList("campaign.id"),
+ new SeaTunnelRowType(
+ new String[] {"campaign.id"},
+ new SeaTunnelDataType[]
{BasicType.LONG_TYPE}),
+ row -> {}));
+ // initial attempt + max_retries(2) retries
+ Assertions.assertEquals(3, searchHits.get());
+ }
+ }
+
+ @Test
+ void searchDoesNotRetryInvalidQuery400AndSurfacesApiError() throws
IOException {
+ registerToken();
+ AtomicInteger searchHits = new AtomicInteger();
+ server.createContext(
+ SEARCH_PATH,
+ exchange -> {
+ searchHits.incrementAndGet();
+ respondJson(
+ exchange,
+ 400,
+ "{\"error\":{\"status\":\"INVALID_ARGUMENT\","
+ + "\"message\":\"Unrecognized field in the
query\"}}");
+ });
+
+ try (GoogleAdsClient client = new GoogleAdsClient(buildParams(null))) {
+ GoogleAdsConnectorException ex =
+ Assertions.assertThrows(
+ GoogleAdsConnectorException.class,
+ () ->
+ client.search(
+ CUSTOMER_ID,
+ "SELECT bad.field FROM campaign",
+ Arrays.asList("bad.field"),
+ new SeaTunnelRowType(
+ new String[] {"bad.field"},
+ new SeaTunnelDataType[] {
+ BasicType.STRING_TYPE
+ }),
+ row -> {}));
+ Assertions.assertEquals(1, searchHits.get());
+ Assertions.assertTrue(ex.getMessage().contains("Unrecognized
field"));
+ }
+ }
+
+ private void registerToken() {
+ server.createContext(
+ "/token",
+ exchange ->
+ respondJson(
+ exchange, 200,
"{\"access_token\":\"tok-1\",\"expires_in\":3600}"));
+ }
+
+ private static String fieldMeta(String name, String dataType) {
+ return "{\"name\":\"" + name + "\",\"dataType\":\"" + dataType + "\"}";
+ }
+
+ private static String readBody(HttpExchange exchange) throws IOException {
+ byte[] buf = new byte[8192];
+ StringBuilder sb = new StringBuilder();
+ int n;
+ while ((n = exchange.getRequestBody().read(buf)) > 0) {
+ sb.append(new String(buf, 0, n, StandardCharsets.UTF_8));
+ }
+ return sb.toString();
+ }
+
+ private static void respondJson(HttpExchange exchange, int status, String
body)
+ throws IOException {
+ byte[] bytes = body.getBytes(StandardCharsets.UTF_8);
+ exchange.getResponseHeaders().add("Content-Type", "application/json");
+ exchange.sendResponseHeaders(status, bytes.length);
+ try (OutputStream out = exchange.getResponseBody()) {
+ out.write(bytes);
+ }
+ }
+
+ private GoogleAdsParameters buildParams(String loginCustomerId) {
+ Map<String, Object> map = new HashMap<>();
+ map.put("developer_token", "dev-token");
+ map.put("client_id", "cid");
+ map.put("client_secret", "csec");
+ map.put("refresh_token", "rtok");
+ map.put("customer_id", CUSTOMER_ID);
+ map.put("api_endpoint", baseUrl);
+ map.put("oauth_endpoint", baseUrl);
+ map.put("max_retries", 2);
+ map.put("retry_backoff_ms", 10L);
+ if (loginCustomerId != null) {
+ map.put("login_customer_id", loginCustomerId);
+ }
+ GoogleAdsParameters params = new GoogleAdsParameters();
+ params.buildWithConfig(ReadonlyConfig.fromMap(map));
+ return params;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceTest.java
b/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceTest.java
new file mode 100644
index 0000000000..22df8d4ae7
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-google-ads/src/test/java/org/apache/seatunnel/connectors/seatunnel/google/ads/source/GoogleAdsSourceTest.java
@@ -0,0 +1,296 @@
+/*
+ * 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.ads.source;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import org.apache.seatunnel.api.table.type.SqlType;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsParameters;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.config.GoogleAdsTableConfig;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.google.ads.exception.GoogleAdsConnectorException;
+
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import com.sun.net.httpserver.HttpServer;
+
+import java.io.IOException;
+import java.io.OutputStream;
+import java.net.InetSocketAddress;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Covers the config-resolution logic of {@link GoogleAdsSource}: three-mode
dispatch, GAQL
+ * building/parsing, table_path handling and field-path validation. The
metadata service is stubbed
+ * with a local HTTP server (same pattern as GoogleAdsClientTest).
+ */
+class GoogleAdsSourceTest {
+
+ private static final String CUSTOMER_ID = "1234567890";
+
+ private HttpServer server;
+ private String baseUrl;
+
+ @BeforeEach
+ void setUp() throws IOException {
+ server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
+ server.createContext(
+ "/token",
+ exchange ->
+ respondJson(
+ exchange, 200,
"{\"access_token\":\"tok-1\",\"expires_in\":3600}"));
+ server.createContext(
+ "/v21/googleAdsFields:search",
+ exchange ->
+ respondJson(
+ exchange,
+ 200,
+ "{\"results\":["
+ + fieldMeta("campaign.id", "INT64")
+ + ","
+ + fieldMeta("campaign.name", "STRING")
+ + ","
+ + fieldMeta("metrics.clicks", "INT64")
+ + ","
+ + fieldMeta("ad_group.id", "INT64")
+ + ","
+ + fieldMeta("ad_group.name", "STRING")
+ + "]}"));
+ server.start();
+ baseUrl = "http://127.0.0.1:" + server.getAddress().getPort();
+ }
+
+ @AfterEach
+ void tearDown() {
+ if (server != null) {
+ server.stop(0);
+ }
+ }
+
+ @Test
+ void resourceModeBuildsGaqlWithFilterAndSchemaInFieldOrder() {
+ Map<String, Object> map = baseConfig();
+ map.put("resource", "campaign");
+ map.put("fields", Arrays.asList("campaign.id", "campaign.name",
"metrics.clicks"));
+ map.put("filter", "segments.date DURING LAST_7_DAYS");
+
+ GoogleAdsSource source = buildSource(map);
+
+ List<GoogleAdsTableConfig> configs = source.getTableConfigs();
+ Assertions.assertEquals(1, configs.size());
+ GoogleAdsTableConfig table = configs.get(0);
+ Assertions.assertEquals(
+ "SELECT campaign.id, campaign.name, metrics.clicks FROM
campaign "
+ + "WHERE segments.date DURING LAST_7_DAYS",
+ table.getGaql());
+ Assertions.assertEquals("campaign", table.getResource());
+ Assertions.assertEquals("google_ads.campaign", table.getTableId());
+ SeaTunnelRowType rowType =
source.getProducedCatalogTables().get(0).getSeaTunnelRowType();
+ Assertions.assertEquals("campaign.id", rowType.getFieldName(0));
+ Assertions.assertEquals("campaign.name", rowType.getFieldName(1));
+ Assertions.assertEquals("metrics.clicks", rowType.getFieldName(2));
+ Assertions.assertEquals(SqlType.BIGINT,
rowType.getFieldType(0).getSqlType());
+ Assertions.assertEquals(SqlType.STRING,
rowType.getFieldType(1).getSqlType());
+ }
+
+ @Test
+ void queryModeParsesFieldsAndResourceFromGaql() {
+ Map<String, Object> map = baseConfig();
+ map.put("query", "select campaign.id , campaign.name from campaign
WHERE campaign.id > 0");
+
+ GoogleAdsTableConfig table = buildSource(map).getTableConfigs().get(0);
+ Assertions.assertEquals(
+ Arrays.asList("campaign.id", "campaign.name"),
table.getFieldPaths());
+ Assertions.assertEquals("campaign", table.getResource());
+ Assertions.assertEquals(
+ "select campaign.id , campaign.name from campaign WHERE
campaign.id > 0",
+ table.getGaql());
+ }
+
+ @Test
+ void queryModeRejectsQueryCombinedWithFields() {
+ Map<String, Object> map = baseConfig();
+ map.put("query", "SELECT campaign.id FROM campaign");
+ map.put("fields", Arrays.asList("campaign.id"));
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
ex.getSeaTunnelErrorCode());
+ Assertions.assertTrue(ex.getMessage().contains("mutually exclusive"));
+ }
+
+ @Test
+ void queryModeRejectsUnparsableGaql() {
+ Map<String, Object> map = baseConfig();
+ map.put("query", "DELETE FROM campaign");
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
ex.getSeaTunnelErrorCode());
+ }
+
+ @Test
+ void resourceModeRequiresNonEmptyFields() {
+ Map<String, Object> map = baseConfig();
+ map.put("resource", "campaign");
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
ex.getSeaTunnelErrorCode());
+ Assertions.assertTrue(ex.getMessage().contains("fields"));
+ }
+
+ @Test
+ void missingAllModesIsRejected() {
+ GoogleAdsConnectorException ex = assertBuildFails(baseConfig());
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
ex.getSeaTunnelErrorCode());
+ }
+
+ @Test
+ void invalidFieldPathsAreRejectedWithOffendingName() {
+ for (String bad : Arrays.asList("Campaign.Id", "campaign",
"campaign.id;drop")) {
+ Map<String, Object> map = baseConfig();
+ map.put("resource", "campaign");
+ map.put("fields", Arrays.asList(bad));
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
ex.getSeaTunnelErrorCode());
+ Assertions.assertTrue(ex.getMessage().contains(bad), "message
should name: " + bad);
+ }
+ }
+
+ @Test
+ void tablesConfigsBuildsMultipleTablesWithPerTableCustomerIdFallback() {
+ Map<String, Object> entryWithOverride = new HashMap<>();
+ entryWithOverride.put("table_path", "google_ads.campaign");
+ entryWithOverride.put("fields", Arrays.asList("campaign.id",
"campaign.name"));
+ entryWithOverride.put("customer_id", "2222222222");
+ Map<String, Object> entryWithFallback = new HashMap<>();
+ entryWithFallback.put("table_path", "google_ads.ad_group");
+ entryWithFallback.put("query", "SELECT ad_group.id, ad_group.name FROM
ad_group");
+
+ Map<String, Object> map = baseConfig();
+ map.put("tables_configs", Arrays.asList(entryWithOverride,
entryWithFallback));
+
+ List<GoogleAdsTableConfig> configs =
buildSource(map).getTableConfigs();
+ Assertions.assertEquals(2, configs.size());
+ Assertions.assertEquals("google_ads.campaign",
configs.get(0).getTableId());
+ Assertions.assertEquals("2222222222", configs.get(0).getCustomerId());
+ Assertions.assertEquals("google_ads.ad_group",
configs.get(1).getTableId());
+ Assertions.assertEquals(CUSTOMER_ID, configs.get(1).getCustomerId());
+ }
+
+ @Test
+ void tablesConfigsRejectsResourceMismatchBetweenTablePathAndQuery() {
+ Map<String, Object> entry = new HashMap<>();
+ entry.put("table_path", "google_ads.campaign");
+ entry.put("query", "SELECT ad_group.id FROM ad_group");
+
+ Map<String, Object> map = baseConfig();
+ map.put("tables_configs", Arrays.asList(entry));
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_QUERY,
ex.getSeaTunnelErrorCode());
+ Assertions.assertTrue(ex.getMessage().contains("does not match"));
+ }
+
+ @Test
+ void tablesConfigsRejectsDuplicateTablePath() {
+ Map<String, Object> entry1 = new HashMap<>();
+ entry1.put("table_path", "google_ads.campaign");
+ entry1.put("fields", Arrays.asList("campaign.id"));
+ Map<String, Object> entry2 = new HashMap<>();
+ entry2.put("table_path", "google_ads.campaign");
+ entry2.put("fields", Arrays.asList("campaign.name"));
+
+ Map<String, Object> map = baseConfig();
+ map.put("tables_configs", Arrays.asList(entry1, entry2));
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.DUPLICATE_RESOURCE,
ex.getSeaTunnelErrorCode());
+ }
+
+ @Test
+ void tablesConfigsRejectsMalformedTablePath() {
+ for (String badPath : Arrays.asList("campaign", ".campaign",
"google_ads.")) {
+ Map<String, Object> entry = new HashMap<>();
+ entry.put("table_path", badPath);
+ entry.put("fields", Arrays.asList("campaign.id"));
+
+ Map<String, Object> map = baseConfig();
+ map.put("tables_configs", Arrays.asList(entry));
+
+ GoogleAdsConnectorException ex = assertBuildFails(map);
+ Assertions.assertEquals(
+ GoogleAdsConnectorErrorCode.INVALID_TABLE_PATH,
+ ex.getSeaTunnelErrorCode(),
+ "table_path should be rejected: " + badPath);
+ }
+ }
+
+ private GoogleAdsConnectorException assertBuildFails(Map<String, Object>
map) {
+ return Assertions.assertThrows(GoogleAdsConnectorException.class, ()
-> buildSource(map));
+ }
+
+ private GoogleAdsSource buildSource(Map<String, Object> map) {
+ ReadonlyConfig config = ReadonlyConfig.fromMap(map);
+ GoogleAdsParameters params = new GoogleAdsParameters();
+ params.buildWithConfig(config);
+ return new GoogleAdsSource(params, config);
+ }
+
+ private Map<String, Object> baseConfig() {
+ Map<String, Object> map = new HashMap<>();
+ map.put("developer_token", "dev-token");
+ map.put("client_id", "cid");
+ map.put("client_secret", "csec");
+ map.put("refresh_token", "rtok");
+ map.put("customer_id", CUSTOMER_ID);
+ map.put("api_endpoint", baseUrl);
+ map.put("oauth_endpoint", baseUrl);
+ map.put("max_retries", 1);
+ map.put("retry_backoff_ms", 10L);
+ return map;
+ }
+
+ private static String fieldMeta(String name, String dataType) {
+ return "{\"name\":\"" + name + "\",\"dataType\":\"" + dataType + "\"}";
+ }
+
+ private static void respondJson(
+ com.sun.net.httpserver.HttpExchange exchange, int status, String
body)
+ throws IOException {
+ byte[] bytes = body.getBytes(StandardCharsets.UTF_8);
+ exchange.getResponseHeaders().add("Content-Type", "application/json");
+ exchange.sendResponseHeaders(status, bytes.length);
+ try (OutputStream out = exchange.getResponseBody()) {
+ out.write(bytes);
+ }
+ }
+}
diff --git a/seatunnel-connectors-v2/pom.xml b/seatunnel-connectors-v2/pom.xml
index 43d7ebfcba..4b508ef7e8 100644
--- a/seatunnel-connectors-v2/pom.xml
+++ b/seatunnel-connectors-v2/pom.xml
@@ -73,6 +73,7 @@
<module>connector-google-sheets</module>
<module>connector-google-firestore</module>
<module>connector-google-pubsub</module>
+ <module>connector-google-ads</module>
<module>connector-azure-queue-storage</module>
<module>connector-slack</module>
<module>connector-rabbitmq</module>
diff --git a/seatunnel-dist/pom.xml b/seatunnel-dist/pom.xml
index 008f658cff..2e3f39d51e 100644
--- a/seatunnel-dist/pom.xml
+++ b/seatunnel-dist/pom.xml
@@ -548,6 +548,12 @@
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-google-ads</artifactId>
+ <version>${project.version}</version>
+ <scope>provided</scope>
+ </dependency>
<dependency>
<groupId>org.apache.seatunnel</groupId>
<artifactId>connector-datahub</artifactId>